145 lines
5.5 KiB
C#
145 lines
5.5 KiB
C#
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
|
||
};
|
||
}
|
||
}
|
||
}
|