ModuleRebuildService.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419
  1. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  2. using Microsoft.Extensions.Logging;
  3. using Microsoft.Extensions.Logging.Abstractions;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  5. public sealed class ModuleRebuildService : ITransient
  6. {
  7. private readonly IModuleRebuildJobStore _jobs;
  8. private readonly ModuleRebuildQueue _queue;
  9. private readonly IModuleRebuildCapability _capability;
  10. private readonly ILogger _logger;
  11. public ModuleRebuildService(IModuleRebuildJobStore jobs, ModuleRebuildQueue queue, IModuleRebuildCapability capability)
  12. : this(jobs, queue, capability, NullLoggerFactory.Instance)
  13. {
  14. }
  15. private readonly MdpNeutralSourceCleanup? _cleanup;
  16. private readonly TransformRunLogFinalizer? _runLogFinalizer;
  17. public ModuleRebuildService(
  18. IModuleRebuildJobStore jobs,
  19. ModuleRebuildQueue queue,
  20. IModuleRebuildCapability capability,
  21. ILoggerFactory loggerFactory)
  22. : this(jobs, queue, capability, loggerFactory, null)
  23. {
  24. }
  25. public ModuleRebuildService(
  26. IModuleRebuildJobStore jobs,
  27. ModuleRebuildQueue queue,
  28. IModuleRebuildCapability capability,
  29. ILoggerFactory loggerFactory,
  30. MdpNeutralSourceCleanup? cleanup,
  31. TransformRunLogFinalizer? runLogFinalizer = null)
  32. {
  33. _jobs = jobs;
  34. _queue = queue;
  35. _capability = capability;
  36. _cleanup = cleanup;
  37. _runLogFinalizer = runLogFinalizer;
  38. _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildService));
  39. }
  40. public async Task<(int StatusCode, ModuleRebuildJobAccepted Body)> EnqueueAsync(
  41. string moduleCode,
  42. long tenantId,
  43. long factoryId,
  44. long? requestedBy,
  45. string triggerType = "MANUAL",
  46. CancellationToken ct = default)
  47. {
  48. var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
  49. if (!_capability.IsEnabled(scope.ModuleCode))
  50. {
  51. return (404, new ModuleRebuildJobAccepted
  52. {
  53. Ok = false,
  54. ModuleCode = scope.ModuleCode,
  55. Status = "DISABLED",
  56. Message = $"{scope.ModuleCode} 数据重算尚未启用"
  57. });
  58. }
  59. var normalizedTrigger = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant();
  60. var manualRequest = ModuleRebuildTriggerType.IsManualTrigger(normalizedTrigger, requestedBy);
  61. var active = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  62. if (active != null)
  63. {
  64. // 人工重算可以取代同 scope 里「还没开跑的自动任务」:二者都是同 scope 全量,取代不丢工作量。
  65. // 不让位则会死锁——ClaimNextQueuedAsync 只让执行机领自动任务,非执行机上这条 QUEUED 永远
  66. // 不会消失,页面「数据重算」按钮就被它永久挡住(执行机指派还会随进程重启失效)。
  67. var supersedable = manualRequest
  68. && active.Status == ModuleRebuildStatus.Queued
  69. && !ModuleRebuildTriggerType.IsManualTrigger(active.TriggerType, active.RequestedBy);
  70. if (!supersedable
  71. || !await _jobs.TrySupersedeQueuedAsync(active.Id, $"已被人工重算取代(原 {active.TriggerType})", ct))
  72. {
  73. return (409, Conflict(scope.ModuleCode, active));
  74. }
  75. _logger.LogInformation(
  76. "[ModuleRebuild] {Module} tenant={Tenant} factory={Factory} 人工重算取代待执行任务 job={JobId} trigger={Trigger}",
  77. scope.ModuleCode, scope.TenantId, scope.FactoryId, active.Id, active.TriggerType);
  78. }
  79. // AUTO / BOOTSTRAP 才冷却。MANUAL 是页面「数据重算」按钮,AUTO_NIGHTLY 是全量兜底,两者都必须放行。
  80. var cooldownHours = _capability.AutoMinIntervalHours;
  81. if (cooldownHours > 0
  82. && !manualRequest
  83. && normalizedTrigger != ModuleRebuildTriggerType.AutoNightly)
  84. {
  85. var lastSuccess = await _jobs.FindLastSuccessAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  86. if (lastSuccess?.FinishedAt is { } finishedAt
  87. && DateTime.Now - finishedAt < TimeSpan.FromHours(cooldownHours))
  88. {
  89. return (409, new ModuleRebuildJobAccepted
  90. {
  91. Ok = false,
  92. ModuleCode = scope.ModuleCode,
  93. JobId = lastSuccess.Id,
  94. Status = "COOLDOWN",
  95. Message = $"{scope.ModuleCode} 距上次成功重算不足 {cooldownHours} 小时,本次跳过"
  96. });
  97. }
  98. }
  99. var now = DateTime.Now;
  100. var created = await _jobs.InsertQueuedAsync(new AdoModuleDashboardRebuildJob
  101. {
  102. ModuleCode = scope.ModuleCode,
  103. TenantId = scope.TenantId,
  104. FactoryId = scope.FactoryId,
  105. Status = ModuleRebuildStatus.Queued,
  106. CurrentStage = ModuleRebuildStages.Queued,
  107. StageIndex = 0,
  108. StageTotal = ModuleRebuildStages.StageTotal(scope.ModuleCode),
  109. ProgressPercent = 0,
  110. ProgressMessage = "已入队",
  111. LastProgressAt = now,
  112. TriggerType = normalizedTrigger,
  113. RequestedBy = requestedBy,
  114. SubmittedAt = now,
  115. CreateTime = now,
  116. UpdateTime = now
  117. }, ct);
  118. var oldest = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  119. if (oldest != null && oldest.Id != created.Id)
  120. {
  121. created.Status = ModuleRebuildStatus.Failed;
  122. created.CurrentStage = ModuleRebuildStages.Failed;
  123. created.FinishedAt = DateTime.Now;
  124. created.ErrorMessage = $"{scope.ModuleCode} 数据重算正在执行,请勿重复提交";
  125. created.UpdateTime = DateTime.Now;
  126. await _jobs.UpdateAsync(created, ct);
  127. return (409, Conflict(scope.ModuleCode, oldest));
  128. }
  129. _queue.Pulse();
  130. return (202, new ModuleRebuildJobAccepted
  131. {
  132. Ok = true,
  133. ModuleCode = scope.ModuleCode,
  134. JobId = created.Id,
  135. Status = ModuleRebuildStatus.Queued,
  136. Message = $"{scope.ModuleCode} 数据重算已排队"
  137. });
  138. }
  139. public async Task<ModuleRebuildJobDto> GetAsync(string moduleCode, long jobId, long tenantId, long factoryId, CancellationToken ct = default)
  140. {
  141. var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
  142. var row = await _jobs.GetByIdAsync(scope.ModuleCode, jobId, scope.TenantId, scope.FactoryId, ct);
  143. if (row == null)
  144. throw Oops.Oh("任务不存在");
  145. return ToDto(row);
  146. }
  147. public async Task<ModuleRebuildJobDto> LatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
  148. {
  149. var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
  150. var row = await _jobs.GetLatestAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  151. return row == null ? null : ToDto(row);
  152. }
  153. /// <summary>
  154. /// 陈旧任务收口:RUNNING 无心跳、以及 QUEUED 长期无人消费,两类都要捞。
  155. ///
  156. /// <para>QUEUED 那一支是 2026-09-28 补的:非执行机只领手工任务,所以执行机指派一丢,
  157. /// 自动任务在队列里没有任何消费者,却会一直挡着同 scope 的后续入队。
  158. /// 只扫 RUNNING 的旧实现看不见这种滞留。</para>
  159. /// </summary>
  160. public async Task FailStaleAsync(CancellationToken ct = default)
  161. {
  162. await _jobs.FailStaleRunningAsync(ModuleRebuildLock.StaleAfter, ct);
  163. var reaped = await _jobs.FailStaleQueuedAsync(
  164. ModuleRebuildLock.QueuedStaleAfter,
  165. $"排队超过 {ModuleRebuildLock.QueuedStaleAfter.TotalHours:0} 小时仍无人领取,已收口;"
  166. + "常见成因是 ETL 执行机未指派或指派已失效",
  167. ct);
  168. if (reaped > 0)
  169. _logger.LogWarning(
  170. "[ModuleRebuild] 收口长期滞留的 QUEUED 任务 {Count} 条(阈值 {Hours} 小时)。"
  171. + "请检查 MDP 运行监控页是否存在「零台存活执行机」",
  172. reaped, ModuleRebuildLock.QueuedStaleAfter.TotalHours);
  173. }
  174. public async Task TouchHeartbeatAsync(long jobId, CancellationToken ct = default)
  175. {
  176. try
  177. {
  178. await _jobs.TouchHeartbeatAsync(jobId, DateTime.Now, ct);
  179. }
  180. catch (Exception ex)
  181. {
  182. _logger.LogWarning(ex, "[MdpRebuild] heartbeat failed jobId={JobId}", jobId);
  183. }
  184. }
  185. public async Task ApplyProgressAsync(long jobId, string moduleCode, ModuleProgressUpdate update, CancellationToken ct = default)
  186. {
  187. try
  188. {
  189. var now = DateTime.Now;
  190. if (update.CompletedStage != null && update.Rows.HasValue)
  191. {
  192. await _jobs.UpdateStageResultAsync(
  193. jobId, update.CompletedStage, update.Rows.Value, update.ProgressPercent, update.Message, update.DetailJson, now, ct);
  194. }
  195. await _jobs.UpdateProgressAsync(
  196. jobId, update.Stage, update.StageIndex, update.ProgressPercent, update.Message, now, ct);
  197. }
  198. catch (Exception ex)
  199. {
  200. _logger.LogWarning(ex, "[MdpRebuild] progress write failed jobId={JobId} module={Module} stage={Stage}", jobId, moduleCode, update.Stage);
  201. }
  202. }
  203. public async Task RunClaimedAsync(
  204. AdoModuleDashboardRebuildJob job,
  205. IModuleRebuildHandler handler,
  206. IModuleRebuildLock runLock,
  207. CancellationToken stoppingToken)
  208. {
  209. var failedStage = job.CurrentStage;
  210. var scope = MdpRebuildScope.Create(job.ModuleCode, job.TenantId, job.FactoryId);
  211. try
  212. {
  213. await ApplyProgressAsync(job.Id, scope.ModuleCode, new ModuleProgressUpdate(
  214. ModuleRebuildStages.AcquiringLock, 0, 2, $"等待现有 {scope.ModuleCode} 全量任务完成"), CancellationToken.None);
  215. var holderId = $"{job.TriggerType}:{scope.ScopeKey}:{Environment.MachineName}:{Guid.NewGuid():N}";
  216. var lease = await runLock.TryAcquireAsync(scope, holderId, job.Id, stoppingToken);
  217. if (lease == null)
  218. throw new ModuleRebuildAlreadyRunningException(scope.ModuleCode);
  219. await using (lease)
  220. {
  221. using var lockHeartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
  222. var lockHeartbeat = KeepLeaseAliveAsync(lease, lockHeartbeatCts.Token);
  223. try
  224. {
  225. if (_cleanup != null)
  226. await _cleanup.PurgeModuleAsync(scope.TenantId, scope.ModuleCode, stoppingToken);
  227. var result = await handler.RunAsync(
  228. scope,
  229. job.TriggerType,
  230. job.Id,
  231. async update =>
  232. {
  233. failedStage = update.Stage;
  234. await ApplyProgressAsync(job.Id, scope.ModuleCode, update, CancellationToken.None);
  235. // 阶段边界:每个 handler 每跨一个阶段都会走到这里,是唯一的通用取消点。
  236. if (await _jobs.IsCancelRequestedAsync(job.Id, CancellationToken.None))
  237. throw new ModuleRebuildCancelledException();
  238. },
  239. stoppingToken);
  240. var now = DateTime.Now;
  241. job.Status = ModuleRebuildStatus.Success;
  242. job.CurrentStage = ModuleRebuildStages.Success;
  243. job.StageIndex = ModuleRebuildStages.StageTotal(scope.ModuleCode);
  244. job.ProgressPercent = 100;
  245. job.ProgressMessage = $"{scope.ModuleCode} 数据重算已完成";
  246. job.FailedStage = null;
  247. job.LastProgressAt = now;
  248. job.BatchId = result.BatchId;
  249. job.TransformRunLogId = result.RunLogId;
  250. job.StageRows = result.StageRows;
  251. job.StandardRows = result.StandardRows;
  252. job.DwdRows = result.DwdRows;
  253. job.KpiRows = result.KpiRows;
  254. job.AtomicRows = result.AtomicRows;
  255. job.DetailJson = result.DetailJson;
  256. job.FinishedAt = now;
  257. job.HeartbeatAt = now;
  258. job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null;
  259. job.UpdateTime = now;
  260. job.ErrorMessage = null;
  261. await _jobs.UpdateAsync(job, CancellationToken.None);
  262. }
  263. finally
  264. {
  265. lockHeartbeatCts.Cancel();
  266. try { await lockHeartbeat; } catch (OperationCanceledException) { }
  267. }
  268. }
  269. }
  270. catch (ModuleRebuildAlreadyRunningException)
  271. {
  272. var now = DateTime.Now;
  273. job.Status = ModuleRebuildStatus.Queued;
  274. job.CurrentStage = ModuleRebuildStages.AcquiringLock;
  275. job.StageIndex = 0;
  276. job.ProgressPercent = 2;
  277. job.ProgressMessage = $"等待现有 {scope.ModuleCode} 全量任务完成";
  278. job.StartedAt = null;
  279. job.HeartbeatAt = now;
  280. job.LastProgressAt = now;
  281. job.UpdateTime = now;
  282. job.ErrorMessage = null;
  283. await _jobs.UpdateAsync(job, CancellationToken.None);
  284. _queue.Pulse();
  285. }
  286. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  287. {
  288. await FailJobAsync(job, "服务停止,任务已取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage);
  289. }
  290. catch (ModuleRebuildCancelledException)
  291. {
  292. await FailJobAsync(job, "已由管理员取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage);
  293. if (_runLogFinalizer != null)
  294. {
  295. var runLogId = await _jobs.FindOpenRunLogIdAsync(scope.ModuleCode, scope.TenantId, job.StartedAt ?? job.SubmittedAt);
  296. if (runLogId > 0)
  297. await _runLogFinalizer.FinalizeAsAdminCancelledAsync(runLogId, job.StartedAt ?? DateTime.Now);
  298. }
  299. }
  300. catch (Exception ex)
  301. {
  302. await FailJobAsync(
  303. job,
  304. Truncate(ex.Message, 2000),
  305. ModuleRebuildStatus.Failed,
  306. ModuleRebuildStages.Failed,
  307. failedStage);
  308. }
  309. }
  310. private async Task FailJobAsync(AdoModuleDashboardRebuildJob job, string message, string status, string currentStage, string failedStage)
  311. {
  312. var now = DateTime.Now;
  313. job.Status = status;
  314. job.CurrentStage = currentStage;
  315. job.FailedStage = status == ModuleRebuildStatus.Failed ? failedStage : job.FailedStage;
  316. job.ProgressMessage = status == ModuleRebuildStatus.Failed
  317. ? $"在【{ModuleRebuildStages.ToChinese(failedStage)}】失败"
  318. : message;
  319. job.FinishedAt = now;
  320. job.HeartbeatAt = now;
  321. job.LastProgressAt = now;
  322. job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null;
  323. job.ErrorMessage = message;
  324. job.UpdateTime = now;
  325. await _jobs.UpdateAsync(job, CancellationToken.None);
  326. }
  327. public static ModuleRebuildJobDto ToDto(AdoModuleDashboardRebuildJob row) => new()
  328. {
  329. Ok = true,
  330. ModuleCode = row.ModuleCode,
  331. JobId = row.Id,
  332. Status = row.Status,
  333. CurrentStage = row.CurrentStage,
  334. StageIndex = row.StageIndex,
  335. StageTotal = row.StageTotal <= 0 ? ModuleRebuildStages.StageTotal(row.ModuleCode) : row.StageTotal,
  336. ProgressPercent = row.ProgressPercent,
  337. ProgressMessage = row.ProgressMessage,
  338. LastProgressAt = row.LastProgressAt,
  339. HeartbeatAt = row.HeartbeatAt,
  340. FailedStage = row.FailedStage,
  341. SubmittedAt = row.SubmittedAt,
  342. StartedAt = row.StartedAt,
  343. FinishedAt = row.FinishedAt,
  344. DurationMs = row.DurationMs,
  345. BatchId = row.BatchId,
  346. StageRows = row.StageRows,
  347. StandardRows = row.StandardRows,
  348. DwdRows = row.DwdRows,
  349. KpiRows = row.KpiRows,
  350. AtomicRows = row.AtomicRows,
  351. DetailJson = row.DetailJson,
  352. ErrorMessage = row.ErrorMessage
  353. };
  354. private static ModuleRebuildJobAccepted Conflict(string moduleCode, AdoModuleDashboardRebuildJob active) => new()
  355. {
  356. Ok = false,
  357. ModuleCode = moduleCode,
  358. JobId = active.Id,
  359. Status = active.Status,
  360. Message = $"{moduleCode} 数据重算正在执行,请勿重复提交"
  361. };
  362. private static async Task KeepLeaseAliveAsync(IModuleRebuildLease lease, CancellationToken ct)
  363. {
  364. while (!ct.IsCancellationRequested)
  365. {
  366. try
  367. {
  368. await Task.Delay(TimeSpan.FromSeconds(30), ct);
  369. await lease.HeartbeatAsync(CancellationToken.None);
  370. }
  371. catch (OperationCanceledException)
  372. {
  373. return;
  374. }
  375. catch
  376. {
  377. // 单例 SqlSugar 连接被业务阶段占用时跳过本次锁心跳,避免打断 MDP。
  378. }
  379. }
  380. }
  381. private static int ToDurationMs(TimeSpan elapsed)
  382. {
  383. var ms = elapsed.TotalMilliseconds;
  384. if (double.IsNaN(ms) || ms <= 0) return 0;
  385. return ms >= int.MaxValue ? int.MaxValue : (int)ms;
  386. }
  387. private static string Truncate(string value, int max) =>
  388. string.IsNullOrEmpty(value) || value.Length <= max ? value : value[..max];
  389. }