S1MdpSyncTransformJob.cs 2.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354
  1. using Admin.NET.Plugin.AiDOP.DataPlatform.S1Refresh;
  2. using Admin.NET.Plugin.AiDOP.Order;
  3. using Furion.Schedule;
  4. using Microsoft.Extensions.DependencyInjection;
  5. using Microsoft.Extensions.Logging;
  6. using System.Text.Json;
  7. namespace Admin.NET.Plugin.AiDOP.Job;
  8. /// <summary>
  9. /// S1 首批 MDP 同步与标准化转换定时任务:按租户/工厂作用域分别入队,由 worker 受控并行。
  10. /// </summary>
  11. [JobDetail("job_s1_mdp_sync_transform", Description = "S1 MDP同步与标准化转换(含订单交付域原子聚合层)", GroupName = "default", Concurrent = false)]
  12. [Period(3600000, TriggerId = "trigger_s1_mdp_sync_transform", Description = "每60分钟执行")]
  13. public class S1MdpSyncTransformJob : IJob
  14. {
  15. private readonly IServiceScopeFactory _scopeFactory;
  16. private readonly ILogger _logger;
  17. public S1MdpSyncTransformJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  18. {
  19. _scopeFactory = scopeFactory;
  20. _logger = loggerFactory.CreateLogger(nameof(S1MdpSyncTransformJob));
  21. }
  22. public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
  23. {
  24. using var scope = _scopeFactory.CreateScope();
  25. var catalog = scope.ServiceProvider.GetRequiredService<S1MdpScopeCatalog>();
  26. var rebuild = scope.ServiceProvider.GetRequiredService<S1DashboardRebuildService>();
  27. var scopes = await catalog.ListEnabledScopesAsync(stoppingToken);
  28. var accepted = 0;
  29. var skipped = 0;
  30. foreach (var runScope in scopes)
  31. {
  32. stoppingToken.ThrowIfCancellationRequested();
  33. try
  34. {
  35. var (status, body) = await rebuild.EnqueueAsync(runScope.TenantId, runScope.FactoryId, null, "AUTO", stoppingToken);
  36. if (status == 202) accepted++;
  37. else skipped++;
  38. _logger.LogInformation(
  39. "S1MdpSyncTransformJob 作用域 tenant={TenantId} factory={FactoryId} status={Status} jobId={JobId}",
  40. runScope.TenantId, runScope.FactoryId, status, body.JobId);
  41. }
  42. catch (Exception ex)
  43. {
  44. _logger.LogError(ex, "S1MdpSyncTransformJob 作用域失败 tenant={TenantId} factory={FactoryId}", runScope.TenantId, runScope.FactoryId);
  45. }
  46. }
  47. _logger.LogInformation("S1MdpSyncTransformJob 完成 accepted={Accepted} skipped={Skipped} payload={Payload}",
  48. accepted, skipped, JsonSerializer.Serialize(new { accepted, skipped, scopes = scopes.Count }));
  49. }
  50. }