Files
parr_api/PARR.AIHITMainLoader/AihitMainLoader.cs

145 lines
5.5 KiB
C#
Raw Permalink Normal View History

using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging;
using PARR.AIHITMainLoader.Models;
using PARR.AIHITMainLoader.Services;
using PARR.AIHITMainLoader.Settings;
using PARR.Core.Common.Interfaces;
using PARR.Core.Common.Interfaces.RabbitServices;
using System.Text.Encodings.Web;
using System.Text.Json;
namespace PARR.AIHITMainLoader
{
internal class AihitMainLoader : IAihitMainLoader
{
private readonly ILogger<AihitMainLoader> logger;
private readonly IIntervalService intervalService;
private readonly WorkerSettings workerSettings;
private readonly LoaderSettings loaderSettings;
private readonly MqSettings mqSettings;
private readonly IRabbitService mqService;
private readonly IServiceProvider serviceProvider;
private static readonly JsonSerializerOptions jsonOptions = new JsonSerializerOptions
{
Encoder = JavaScriptEncoder.UnsafeRelaxedJsonEscaping,
};
public AihitMainLoader(
ILogger<AihitMainLoader> logger,
IIntervalService intervalService,
WorkerSettings workerSettings,
LoaderSettings loaderSettings,
MqSettings mqSettings,
IRabbitService mqService,
IServiceProvider serviceProvider
)
{
this.logger = logger;
this.intervalService = intervalService;
this.workerSettings = workerSettings;
this.loaderSettings = loaderSettings;
this.mqSettings = mqSettings;
this.mqService = mqService;
this.serviceProvider = serviceProvider;
}
public async Task StartAsync()
{
logger.LogInformation("Запуск сервиса загрузки данных из АИХ ИТ.");
await intervalService.IntervalInitAsync(LoadDataAsync, workerSettings.RepeatEvery);
}
private async Task LoadDataAsync()
{
try
{
logger.LogInformation("Запуск загрузки данных из АИХ ИТ.");
using (var scope = serviceProvider.CreateScope())
{
var aihitService = scope.ServiceProvider.GetRequiredService<IAihitService>();
var services = GetServices(aihitService);
var currentBatch = new List<string>();
int packageSize = loaderSettings.PackageSize;
if (packageSize <= 0)
{
logger.LogError("Некорректный размер пакета: {PackageSize}. Ожидалось > 0.", packageSize);
return;
}
foreach (var service in services)
{
foreach (var responseArea in loaderSettings.ResponseAreas)
{
logger.LogDebug($"responseArea = {responseArea}");
var EKs = service.Invoke(responseArea);
if (EKs?.Any() != true) continue;
foreach (var item in EKs)
{
if (item == null) continue;
var mainData = item.ToMainData();
if (mainData == null) continue;
var json = JsonSerializer.Serialize(mainData, jsonOptions);
currentBatch.Add(json);
// Отправка, если набрали полный пакет
if (currentBatch.Count >= packageSize)
{
await SendBatchAsync(currentBatch);
currentBatch.Clear();
}
}
}
}
// Отправка остатка
if (currentBatch.Count > 0)
{
await SendBatchAsync(currentBatch);
}
}
}
catch (Exception ex)
{
logger.LogError(ex, "Ошибка при выполнении загрузки данных из АИХ ИТ.");
}
}
private async Task SendBatchAsync(List<string> batch)
{
var sendResult = await mqService.SendAsync(mqSettings, batch.ToArray());
if (sendResult.IsSuccess)
{
logger.LogInformation($"Данные переданы в RabbitMQ: {batch.Count}");
logger.LogDebug($"Отправлено сообщений: {batch.Count}, пример первого: {batch.FirstOrDefault()?.Substring(0, Math.Min(200, batch.FirstOrDefault()?.Length ?? 0))}...");
}
else
{
logger.LogError($"Ошибка при передаче данных в RabbitMQ. Не переданные сообщения: {string.Join(", ", sendResult.NotSendMessages!)}");
}
}
private List<Func<string, IEnumerable<IMainData>?>> GetServices(IAihitService service)
{
return new()
{
service.GetCvkData,
service.GetPtkData,
service.GetRegionalEKData,
service.GetStoData,
service.GetOrgData
};
}
}
}