feat(aihitMainSyncer): refactoring+изменение запуска Worker.

This commit is contained in:
Mikhail Kuznetsov
2025-11-19 11:52:41 +10:00
parent 629a487875
commit d486c72416
5 changed files with 117 additions and 81 deletions

View File

@@ -54,9 +54,7 @@ namespace PARR.AIHITMainLoader
logger.LogInformation("Запуск загрузки данных из АИХ ИТ."); logger.LogInformation("Запуск загрузки данных из АИХ ИТ.");
using (var scope = serviceProvider.CreateScope()) using (var scope = serviceProvider.CreateScope())
{ {
var aihitService = scope.ServiceProvider.GetService<IAihitService>(); var aihitService = scope.ServiceProvider.GetRequiredService<IAihitService>();
if (aihitService == null)
throw new Exception($"Не найден сервис: {nameof(IAihitService)}");
var services = GetServices(aihitService); var services = GetServices(aihitService);
@@ -89,14 +87,21 @@ namespace PARR.AIHITMainLoader
int count = preparedDataList.Count; int count = preparedDataList.Count;
int packageSize = loaderSettings.PackageSize; int packageSize = loaderSettings.PackageSize;
if (packageSize <= 0)
{
logger.LogError("Некорректный размер пакета: {PackageSize}. Ожидалось > 0.", packageSize);
return;
}
var parts = count == 0 ? 0 : (count + packageSize - 1) / packageSize; var parts = count == 0 ? 0 : (count + packageSize - 1) / packageSize;
for (var i = 0; i < (count + packageSize - 1) / packageSize; i++) for (var i = 0; i < parts; i++)
{ {
int start = i * packageSize; int start = i * packageSize;
int length = Math.Min(packageSize, count - start); int length = Math.Min(packageSize, count - start);
var batch = preparedDataList.GetRange(start, length).ToArray(); var batch = preparedDataList.GetRange(start, length).ToArray();
var sendResult = await mqService.SendAsync(mqSettings, batch.ToArray()); var sendResult = await mqService.SendAsync(mqSettings, batch);
if (sendResult.IsSuccess) if (sendResult.IsSuccess)
{ {
logger.LogInformation($"Данные переданы в RabbitMQ: {batch.Count()}"); logger.LogInformation($"Данные переданы в RabbitMQ: {batch.Count()}");
@@ -106,7 +111,8 @@ namespace PARR.AIHITMainLoader
logger.LogError($"Ошибка при передаче данных в Rabbit, не переданные сообщения: {string.Join(", ", sendResult.NotSendMessages!)}"); logger.LogError($"Ошибка при передаче данных в Rabbit, не переданные сообщения: {string.Join(", ", sendResult.NotSendMessages!)}");
} }
} }
} catch (Exception ex) }
catch (Exception ex)
{ {
logger.LogError(ex, "Ошибка при выполнении загрузки данных из АИХ ИТ."); logger.LogError(ex, "Ошибка при выполнении загрузки данных из АИХ ИТ.");
} }
@@ -115,13 +121,12 @@ namespace PARR.AIHITMainLoader
private List<Func<string, IEnumerable<IMainData>?>> GetServices(IAihitService service) private List<Func<string, IEnumerable<IMainData>?>> GetServices(IAihitService service)
{ {
var serviceList = new List<Func<string, IEnumerable<IMainData>?>>(); return new()
{
serviceList.Add(service.GetCvkData); service.GetCvkData,
serviceList.Add(service.GetPtkData); service.GetPtkData,
serviceList.Add(service.GetRegionalEKData); service.GetRegionalEKData
};
return serviceList;
} }
} }
} }

View File

@@ -33,10 +33,7 @@ namespace PARR.AIHITMainSyncer
{ {
using (var scope = serviceProvider.CreateScope()) using (var scope = serviceProvider.CreateScope())
{ {
var syncerService = scope.ServiceProvider.GetService<ISyncerService>(); var syncerService = scope.ServiceProvider.GetRequiredService<ISyncerService>();
if (syncerService == null)
throw new Exception("Не смог получить серивс ISyncerService, scope.ServiceProvider.GetService<ISyncerService>()");
await syncerService.SyncAsync(msg); await syncerService.SyncAsync(msg);
} }

View File

@@ -3,7 +3,6 @@ using Microsoft.Extensions.Logging;
using Newtonsoft.Json.Linq; using Newtonsoft.Json.Linq;
using PARR.BLL.Domain.Mq; using PARR.BLL.Domain.Mq;
using PARR.BLL.Services.Interfaces; using PARR.BLL.Services.Interfaces;
using PARR.DAL.Extensions;
using PARR.DAL.Models.Unit; using PARR.DAL.Models.Unit;
using PARR.DAL.Services.Interfaces.Unit; using PARR.DAL.Services.Interfaces.Unit;
@@ -52,6 +51,9 @@ namespace PARR.AIHITMainSyncer.Services
if (!objFromQuery.Properties.Any()) if (!objFromQuery.Properties.Any())
return; return;
foreach (var kvp in objFromQuery.Properties.ToList())
objFromQuery.Properties[kvp.Key.Trim()] = kvp.Value?.Trim();
var listAihitData = new List<KeyValuePair<string, string?>>(); var listAihitData = new List<KeyValuePair<string, string?>>();
var propertisWithSimpleValues = objFromQuery.Properties var propertisWithSimpleValues = objFromQuery.Properties
@@ -69,45 +71,54 @@ namespace PARR.AIHITMainSyncer.Services
var unit = await GetUnit(objFromQuery.Name); var unit = await GetUnit(objFromQuery.Name);
var unitValuesToRemove = unit.UnitValues.Where(t => !listAihitData.Any(a => var incomingPairs = listAihitData
IsStringEqual(t.Field!.AihitName, a.Key) .Select(kv => new { kv.Key, kv.Value })
) || .ToList();
!listAihitData.Any(a =>
IsStringEqual(t.Field!.AihitName, a.Key) && IsStringEqual(t.Value!.Value, a.Value) var unitValuesToRemove = unit.UnitValues
)).ToList(); .Where(uv => !incomingPairs.Any(ip =>
IsStringEqual(uv.Field!.AihitName, ip.Key) &&
IsStringEqual(uv.Value!.Value, ip.Value)))
.ToList();
unitValuesToRemove.ForEach(uv => unit.UnitValues.Remove(uv)); unitValuesToRemove.ForEach(uv => unit.UnitValues.Remove(uv));
var unitValueExisting = unit.UnitValues.Where(t => listAihitData.Any(a => var unitValueMissing = listAihitData.Where(t =>
IsStringEqual(t.Field!.AihitName, a.Key) &&
IsStringEqual(t.Value!.Value, a.Value)
)).ToList();
var unitValueMising = listAihitData.Where(t =>
!unit.UnitValues.Any(a => !unit.UnitValues.Any(a =>
IsStringEqual(t.Key, a.Field!.AihitName) && IsStringEqual(t.Key, a.Field!.AihitName) &&
IsStringEqual(t.Value, a.Value!.Value) IsStringEqual(t.Value, a.Value!.Value)
)).ToList(); )).ToList();
if (unitValueMising.Any()) if (unitValueMissing.Any())
{ {
var fieldsInChanges = unitValueMising.Select(s => s.Key).Distinct().ToList(); var fieldsInChanges = unitValueMissing.Select(s => s.Key).Distinct().ToList();
var valuesInChanges = unitValueMising.Select(s => s.Value).Distinct().ToList(); var valuesInChanges = unitValueMissing.Select(s => s.Value).Distinct().ToList();
await SyncFieldsAsync(fieldsInChanges!, fieldsFromDB); await SyncFieldsAsync(fieldsInChanges!, fieldsFromDB);
await SyncValuesAsync(valuesInChanges); await SyncValuesAsync(valuesInChanges);
//fieldsInChanges = fieldsInChanges.Select(s => s.ToLower()).ToList()!; var normalizedFieldNames = fieldsInChanges
//valuesInChanges = valuesInChanges.Select(s => (s == null) ? null : Sanitize(s.ToLower())).ToList()!; .Select(n => n?.Trim())
.Where(n => n != null)
.ToHashSet(StringComparer.OrdinalIgnoreCase);
var fieldsNameWithId = await unitFieldService.Get().AsNoTracking().Where(t => fieldsInChanges.Any(c => var fieldsNameWithId = await unitFieldService.Get().AsNoTracking()
c == t.AihitName//.AihitName.ToLower() .Where(f => normalizedFieldNames.Contains(f.AihitName))
)).ToListAsync(); .ToListAsync();
var valuesWithId = await unitFieldValueService.Get().AsNoTracking().Where(t => valuesInChanges.Any(c =>
c == t.Value//((t.Value == null) ? null : t.Value.ToLower())
)).ToListAsync();
foreach (var change in unitValueMising) var normalizedValues = valuesInChanges
.Where(v => v != null)
.Select(v => v!.Trim())
.ToHashSet(StringComparer.OrdinalIgnoreCase);
var valuesWithId = await unitFieldValueService.Get().AsNoTracking()
.Where(v =>
(v.Value == null && valuesInChanges.Contains(null)) ||
(v.Value != null && normalizedValues.Contains(v.Value)))
.ToListAsync();
foreach (var change in unitValueMissing)
ApplyChange(unit, fieldsNameWithId, valuesWithId, change); ApplyChange(unit, fieldsNameWithId, valuesWithId, change);
if (!await unitService.CommitAsync()) if (!await unitService.CommitAsync())
@@ -161,50 +172,72 @@ namespace PARR.AIHITMainSyncer.Services
private async Task SyncValuesAsync(List<string?> values) private async Task SyncValuesAsync(List<string?> values)
{ {
bool hasNull = values.Any(v => v is null);
// все поля из АИХИТ // все поля из АИХИТ
var valuesFromAihit = values.Distinct(); var normalizedNoneNullValues = values
.Where(v => v != null)
.Select(v => v!.Trim())
.Distinct(StringComparer.OrdinalIgnoreCase)
.ToList();
// список values для поиска в БД без null // Получаем существующие значения из БД
var valuesFromAihitWithoutNull = valuesFromAihit.Where(x => x != null).ToList(); var existingValuesInDb = await unitFieldValueService.Get()
.Where(v => normalizedNoneNullValues.Contains(v.Value!) || v.Value == null)
var valuesInDb = await unitFieldValueService.Get().Where(t => .Select(v => v.Value)
(t.Value != null &&
valuesFromAihitWithoutNull.Any(x => x == t.Value)
)
|| (t.Value == null)
).Select(t => t.Value)
.ToListAsync(); .ToListAsync();
if (!valuesFromAihit.Any(x => x == null)) if (!hasNull)
valuesInDb.Remove(null); existingValuesInDb.RemoveAll(v => v is null);
var newValues = valuesFromAihit.Where(t => !valuesInDb.Any(x => IsStringEqual(x, t))).ToList(); var newNonNullValues = normalizedNoneNullValues
.Where(v => !existingValuesInDb.Any(ev =>
string.Equals(ev, v, StringComparison.OrdinalIgnoreCase)))
.ToList();
if (newValues.Any()) var newValuesToInsert = new List<string?>();
{ newValuesToInsert.AddRange(newNonNullValues);
foreach (var item in newValues)
{ if (hasNull && !existingValuesInDb.Contains(null))
var value = new UnitFieldValue newValuesToInsert.Add(null);
if (newValuesToInsert.Count == 0)
return;
// Создаём сущности — сохраняем оригинальный регистр из nonNullValues!
// (для null — просто null)
var newFieldValues = newValuesToInsert.Select(v =>
new UnitFieldValue
{ {
Id = Guid.NewGuid(), Id = Guid.NewGuid(),
DateCreated = DateTimeOffset.UtcNow, Value = v // ← v — уже trim()-нутый non-null, или null
Value = item })
}; .ToList();
if (!await unitFieldValueService.CreateAsync(value)) if (!await unitFieldValueService.AddRangeAsync(newFieldValues) || !await unitFieldValueService.CommitAsync())
logger.LogError($"Не удалось создать запись в таблице FieldValues: {item}, {value.ToJson()}"); {
logger.LogError(
"Не удалось сохранить {Count} новых значений FieldValues",
newValuesToInsert.Count);
}
else else
logger.LogInformation($"Создана запись а таблице FieldValues: {item}, {value.ToJson()}"); {
if (newValuesToInsert.Count <= 10)
{
var preview = string.Join(", ", newValuesToInsert.Select(v => v ?? "<null>"));
logger.LogInformation("Добавлено {Count} значений FieldValues: [{Values}]", newValuesToInsert.Count, preview);
}
else
{
logger.LogInformation("Добавлено {Count} значений FieldValues (первые 5: {Preview})",
newValuesToInsert.Count,
string.Join(", ", newValuesToInsert.Take(5).Select(v => v ?? "<null>")));
} }
if (!await unitFieldValueService.CommitAsync())
logger.LogError($"Не удалось применить изменения по добавлению новых Fields в базе данных");
} }
} }
//private async Task SyncFieldsAsync(List<string> fieldsFromAihit, List<FieldDto> fieldsFromDB)
private async Task SyncFieldsAsync(List<string> fieldsFromAihit, List<UnitField> fieldsFromDB) private async Task SyncFieldsAsync(List<string> fieldsFromAihit, List<UnitField> fieldsFromDB)
{ {
//var newFields = fieldsFromAihit.Where(t => !fieldsFromDB.Any(f => f.Name == t));
var newFields = fieldsFromAihit.Where(t => !fieldsFromDB.Any(f => f.AihitName == t)); var newFields = fieldsFromAihit.Where(t => !fieldsFromDB.Any(f => f.AihitName == t));
if (newFields.Any()) if (newFields.Any())
@@ -231,11 +264,10 @@ namespace PARR.AIHITMainSyncer.Services
private bool IsStringEqual(string? value1, string? value2) private bool IsStringEqual(string? value1, string? value2)
{ {
//var _value1 = value1?.ToLower().Trim(); return string.Equals(
//var _value2 = value2?.ToLower().Trim(); value1?.Trim(),
var _value1 = value1; value2?.Trim(),
var _value2 = value2; StringComparison.OrdinalIgnoreCase);
return _value1 == _value2;
} }
@@ -252,7 +284,7 @@ namespace PARR.AIHITMainSyncer.Services
private async Task<Unit> CreateUnitAsync(string name) private async Task<Unit> CreateUnitAsync(string name)
{ {
var unit = new Unit { Name = name.Trim() }; var unit = new Unit { Name = name.Trim().ToUpperInvariant() };
if (!await unitService.CreateAsync(unit) || !await unitService.CommitAsync()) if (!await unitService.CreateAsync(unit) || !await unitService.CommitAsync())
logger.LogError($"Не удалось создать Unit {unit.Name}"); logger.LogError($"Не удалось создать Unit {unit.Name}");
@@ -265,8 +297,9 @@ namespace PARR.AIHITMainSyncer.Services
public async Task<Unit?> GetUnitByName(string name) public async Task<Unit?> GetUnitByName(string name)
{ {
var normalized = name.Trim().ToUpperInvariant();
return await GetUnitWithFieldsAndValues() return await GetUnitWithFieldsAndValues()
.FirstOrDefaultAsync(u => u.Name.ToLower() == name.ToLower()); .FirstOrDefaultAsync(u => u.Name == normalized);
} }

View File

@@ -17,12 +17,13 @@ namespace PARR.AIHITSyncerWorker
protected override async Task ExecuteAsync(CancellationToken stoppingToken) protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{ {
await aihitSyncer.StartAsync(); await aihitSyncer.StartAsync();
await Task.Delay(Timeout.Infinite, stoppingToken);
} }
public override Task StopAsync(CancellationToken cancellationToken) public override Task StopAsync(CancellationToken cancellationToken)
{ {
aihitSyncer.StopAsync().Wait(); aihitSyncer.StopAsync();
return base.StopAsync(cancellationToken); return base.StopAsync(cancellationToken);
} }

View File

@@ -1,7 +1,7 @@
{ {
"Serilog": { "Serilog": {
"MinimumLevel": { "MinimumLevel": {
"Default": "Information", "Default": "Debug",
"Override": { "Override": {
//"Microsoft": "Information", //"Microsoft": "Information",
"Microsoft.Hosting.Lifetime": "Information" "Microsoft.Hosting.Lifetime": "Information"