S5MdpRefreshJob.cs 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125
  1. using Admin.NET.Core.Service;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform;
  3. using Admin.NET.Plugin.AiDOP.FinishedWarehouse;
  4. using Admin.NET.Plugin.AiDOP.Infrastructure;
  5. using Admin.NET.Plugin.AiDOP.Manufacturing;
  6. using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  7. using Furion.Schedule;
  8. using Microsoft.Extensions.DependencyInjection;
  9. using Microsoft.Extensions.Logging;
  10. using System.Text.Json;
  11. namespace Admin.NET.Plugin.AiDOP.Job;
  12. /// <summary>
  13. /// T8 基础数据 + S5/S6/S7 KPI 统一刷新编排任务(方案④ owner)。
  14. /// 一个 tick 内:先执行 1 次 T8 基表双模式入站(源→mdp_stg_t8_*→mdp_std_t8_*),
  15. /// 再顺序计算 S5/S6/S7 三个模块的 T8 KPI(均读同一份本轮刷新后的 std,保证强一致)。
  16. /// JobId 保留 job_s5_t8_kpi_refresh(不迁移 Job 身份/历史);原 S6/S7 独立 Job 已合并至本 Job。
  17. /// 每日 02:00 / 07:00 / 12:00 / 17:00 / 22:00 共 5 次。
  18. ///
  19. /// 失败语义(关键,勿机械照抄 Bootstrap):
  20. /// ① T8 inbound 是整轮硬门——失败则记错误日志并立即结束本轮,
  21. /// 禁止继续算 S5/S6/S7(不允许退化为用上一周期 std 静默计算的 eventual consistency);
  22. /// ② inbound 成功后,S5→S6→S7 采用步级失败隔离(某模块失败记错误、继续下一模块);
  23. /// ③ OperationCanceledException 保持取消语义直接上抛,不吞。
  24. /// 手工补偿入口 /t8-base-mdp/inbound 仍独立保留。
  25. ///
  26. /// 并发:与手动刷新(AidopKanbanController → AidopT8KpiManualRefreshService)共享同一把硬 TTL 刷新锁
  27. /// (RefreshLockKey)。占用时本轮定时跳过(5 次/天,下一 tick 自愈),避免共享入站/std 与手动刷新并发写;
  28. /// 本轮结束在 finally 主动释放,即使异常/挂死也由 TTL 到点自动过期,不永久阻塞后续刷新。
  29. /// </summary>
  30. [JobDetail("job_s5_t8_kpi_refresh",
  31. Description = "T8 基础数据 + S5/S6/S7 KPI 统一刷新编排(inbound 硬门 → S5→S6→S7;5 次/天:02/07/12/17/22)",
  32. GroupName = "default",
  33. Concurrent = false)]
  34. [Cron("0 2,7,12,17,22 * * *",
  35. TriggerId = "trigger_s5_t8_kpi_refresh",
  36. Description = "每日 02:00/07:00/12:00/17:00/22:00 触发(5 字段:分 时 日 月 周,默认 CronStringFormat.Default)")]
  37. public class S5MdpRefreshJob : IJob
  38. {
  39. private readonly IServiceScopeFactory _scopeFactory;
  40. private readonly ILogger _logger;
  41. public S5MdpRefreshJob(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  42. {
  43. _scopeFactory = scopeFactory;
  44. _logger = loggerFactory.CreateLogger(nameof(S5MdpRefreshJob));
  45. }
  46. public async Task ExecuteAsync(JobExecutingContext context, CancellationToken stoppingToken)
  47. {
  48. // 与手动刷新共享同一把硬 TTL 刷新锁:占用时本轮跳过;本轮结束 finally 释放,异常/挂死由 TTL 自动过期兜底。
  49. using var lockScope = _scopeFactory.CreateScope();
  50. var cache = lockScope.ServiceProvider.GetRequiredService<SysCacheService>();
  51. if (cache.ExistKey(AidopT8KpiManualRefreshService.RefreshLockKey))
  52. {
  53. _logger.LogWarning("S5MdpRefreshJob 检测到刷新锁被占用(手动刷新或上轮未释放),本轮跳过");
  54. return;
  55. }
  56. cache.Set(AidopT8KpiManualRefreshService.RefreshLockKey, $"AUTO:{DateTime.Now:yyyyMMddHHmmss}",
  57. TimeSpan.FromSeconds(AidopT8KpiManualRefreshService.RefreshLockTtlSeconds));
  58. try
  59. {
  60. // ① 硬门:T8 基表入站(源→stg→std)必须先成功;失败则本轮结束,不继续下游、不降级用旧 std。
  61. try
  62. {
  63. using var inboundScope = _scopeFactory.CreateScope();
  64. var inbound = inboundScope.ServiceProvider.GetRequiredService<T8BaseInboundMdpSyncService>();
  65. // T8 源表无 tenant_id;定时任务显式经账套映射解析真实租户。
  66. var t8TenantId = AidopSourceTenantMap.ResolveTenantId("pbxfxp", 0);
  67. var inboundResult = await inbound.RunInboundAsync(
  68. tenantId: t8TenantId, fullRefresh: true, entityCode: null, cancellationToken: stoppingToken);
  69. _logger.LogInformation("S5MdpRefreshJob T8 inbound 完成 {Payload}", JsonSerializer.Serialize(inboundResult));
  70. }
  71. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  72. {
  73. _logger.LogInformation("S5MdpRefreshJob T8 inbound 收到停止信号,本轮结束");
  74. throw;
  75. }
  76. catch (Exception ex)
  77. {
  78. _logger.LogError(ex, "S5MdpRefreshJob T8 inbound 失败,本轮跳过 S5/S6/S7 KPI 刷新(不使用上一周期 std 继续计算)");
  79. return;
  80. }
  81. // ② 下游:inbound 成功后,S5→S6→S7 顺序执行,步级失败隔离(各自独立 DI scope)。
  82. await RunStepAsync("S5", stoppingToken, async sp =>
  83. (object)await sp.GetRequiredService<S5MdpSyncTransformService>().RunFullAsync(stoppingToken, "AUTO"));
  84. await RunStepAsync("S6", stoppingToken, async sp =>
  85. (object)await sp.GetRequiredService<S6MdpSyncTransformService>().RunFullAsync(stoppingToken, "AUTO"));
  86. await RunStepAsync("S7", stoppingToken, async sp =>
  87. (object)await sp.GetRequiredService<S7MdpSyncTransformService>().RunFullAsync(stoppingToken, "AUTO"));
  88. }
  89. finally
  90. {
  91. cache.Remove(AidopT8KpiManualRefreshService.RefreshLockKey);
  92. }
  93. }
  94. /// <summary>
  95. /// 下游模块步级执行:独立 DI scope + 失败隔离(仿 SmartOpsKpiMdpBootstrapJob.RunStepAsync)。
  96. /// 仅用于 inbound 成功之后的 S5/S6/S7;不得用于包裹 inbound(inbound 失败必须中断本轮)。
  97. /// </summary>
  98. private async Task RunStepAsync(string step, CancellationToken stoppingToken, Func<IServiceProvider, Task<object>> action)
  99. {
  100. using var scope = _scopeFactory.CreateScope();
  101. try
  102. {
  103. stoppingToken.ThrowIfCancellationRequested();
  104. var payload = await action(scope.ServiceProvider);
  105. _logger.LogInformation("S5MdpRefreshJob {Step} 完成 {Payload}", step, JsonSerializer.Serialize(payload));
  106. }
  107. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  108. {
  109. _logger.LogInformation("S5MdpRefreshJob {Step} 收到停止信号", step);
  110. throw;
  111. }
  112. catch (Exception ex)
  113. {
  114. // 失败通知由各 service 内部 MarkTransformRunFailedAsync 触发;此处记录并继续下一步
  115. _logger.LogError(ex, "S5MdpRefreshJob {Step} 执行失败,继续后续步骤", step);
  116. }
  117. }
  118. }