ModuleRebuildService.cs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330
  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. public ModuleRebuildService(
  16. IModuleRebuildJobStore jobs,
  17. ModuleRebuildQueue queue,
  18. IModuleRebuildCapability capability,
  19. ILoggerFactory loggerFactory)
  20. {
  21. _jobs = jobs;
  22. _queue = queue;
  23. _capability = capability;
  24. _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildService));
  25. }
  26. public async Task<(int StatusCode, ModuleRebuildJobAccepted Body)> EnqueueAsync(
  27. string moduleCode,
  28. long tenantId,
  29. long factoryId,
  30. long? requestedBy,
  31. string triggerType = "MANUAL",
  32. CancellationToken ct = default)
  33. {
  34. var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
  35. if (!_capability.IsEnabled(scope.ModuleCode))
  36. {
  37. return (404, new ModuleRebuildJobAccepted
  38. {
  39. Ok = false,
  40. ModuleCode = scope.ModuleCode,
  41. Status = "DISABLED",
  42. Message = $"{scope.ModuleCode} 数据重算尚未启用"
  43. });
  44. }
  45. var active = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  46. if (active != null)
  47. {
  48. return (409, Conflict(scope.ModuleCode, active));
  49. }
  50. var now = DateTime.Now;
  51. var created = await _jobs.InsertQueuedAsync(new AdoModuleDashboardRebuildJob
  52. {
  53. ModuleCode = scope.ModuleCode,
  54. TenantId = scope.TenantId,
  55. FactoryId = scope.FactoryId,
  56. Status = ModuleRebuildStatus.Queued,
  57. CurrentStage = ModuleRebuildStages.Queued,
  58. StageIndex = 0,
  59. StageTotal = ModuleRebuildStages.StageTotal(scope.ModuleCode),
  60. ProgressPercent = 0,
  61. ProgressMessage = "已入队",
  62. LastProgressAt = now,
  63. TriggerType = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant(),
  64. RequestedBy = requestedBy,
  65. SubmittedAt = now,
  66. CreateTime = now,
  67. UpdateTime = now
  68. }, ct);
  69. var oldest = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  70. if (oldest != null && oldest.Id != created.Id)
  71. {
  72. created.Status = ModuleRebuildStatus.Failed;
  73. created.CurrentStage = ModuleRebuildStages.Failed;
  74. created.FinishedAt = DateTime.Now;
  75. created.ErrorMessage = $"{scope.ModuleCode} 数据重算正在执行,请勿重复提交";
  76. created.UpdateTime = DateTime.Now;
  77. await _jobs.UpdateAsync(created, ct);
  78. return (409, Conflict(scope.ModuleCode, oldest));
  79. }
  80. _queue.Pulse();
  81. return (202, new ModuleRebuildJobAccepted
  82. {
  83. Ok = true,
  84. ModuleCode = scope.ModuleCode,
  85. JobId = created.Id,
  86. Status = ModuleRebuildStatus.Queued,
  87. Message = $"{scope.ModuleCode} 数据重算已排队"
  88. });
  89. }
  90. public async Task<ModuleRebuildJobDto> GetAsync(string moduleCode, long jobId, long tenantId, long factoryId, CancellationToken ct = default)
  91. {
  92. var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
  93. var row = await _jobs.GetByIdAsync(scope.ModuleCode, jobId, scope.TenantId, scope.FactoryId, ct);
  94. if (row == null)
  95. throw Oops.Oh("任务不存在");
  96. return ToDto(row);
  97. }
  98. public async Task<ModuleRebuildJobDto> LatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
  99. {
  100. var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
  101. var row = await _jobs.GetLatestAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
  102. return row == null ? null : ToDto(row);
  103. }
  104. public async Task FailStaleAsync(CancellationToken ct = default) =>
  105. await _jobs.FailStaleRunningAsync(ModuleRebuildLock.StaleAfter, ct);
  106. public async Task TouchHeartbeatAsync(long jobId, CancellationToken ct = default)
  107. {
  108. try
  109. {
  110. await _jobs.TouchHeartbeatAsync(jobId, DateTime.Now, ct);
  111. }
  112. catch (Exception ex)
  113. {
  114. _logger.LogWarning(ex, "[MdpRebuild] heartbeat failed jobId={JobId}", jobId);
  115. }
  116. }
  117. public async Task ApplyProgressAsync(long jobId, string moduleCode, ModuleProgressUpdate update, CancellationToken ct = default)
  118. {
  119. try
  120. {
  121. var now = DateTime.Now;
  122. if (update.CompletedStage != null && update.Rows.HasValue)
  123. {
  124. await _jobs.UpdateStageResultAsync(
  125. jobId, update.CompletedStage, update.Rows.Value, update.ProgressPercent, update.Message, update.DetailJson, now, ct);
  126. }
  127. await _jobs.UpdateProgressAsync(
  128. jobId, update.Stage, update.StageIndex, update.ProgressPercent, update.Message, now, ct);
  129. }
  130. catch (Exception ex)
  131. {
  132. _logger.LogWarning(ex, "[MdpRebuild] progress write failed jobId={JobId} module={Module} stage={Stage}", jobId, moduleCode, update.Stage);
  133. }
  134. }
  135. public async Task RunClaimedAsync(
  136. AdoModuleDashboardRebuildJob job,
  137. IModuleRebuildHandler handler,
  138. IModuleRebuildLock runLock,
  139. CancellationToken stoppingToken)
  140. {
  141. var failedStage = job.CurrentStage;
  142. var scope = MdpRebuildScope.Create(job.ModuleCode, job.TenantId, job.FactoryId);
  143. try
  144. {
  145. await ApplyProgressAsync(job.Id, scope.ModuleCode, new ModuleProgressUpdate(
  146. ModuleRebuildStages.AcquiringLock, 0, 2, $"等待现有 {scope.ModuleCode} 全量任务完成"), CancellationToken.None);
  147. var holderId = $"{job.TriggerType}:{scope.ScopeKey}:{Environment.MachineName}:{Guid.NewGuid():N}";
  148. var lease = await runLock.TryAcquireAsync(scope, holderId, job.Id, stoppingToken);
  149. if (lease == null)
  150. throw new ModuleRebuildAlreadyRunningException(scope.ModuleCode);
  151. await using (lease)
  152. {
  153. using var lockHeartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
  154. var lockHeartbeat = KeepLeaseAliveAsync(lease, lockHeartbeatCts.Token);
  155. try
  156. {
  157. var result = await handler.RunAsync(
  158. scope,
  159. job.TriggerType,
  160. job.Id,
  161. async update =>
  162. {
  163. failedStage = update.Stage;
  164. await ApplyProgressAsync(job.Id, scope.ModuleCode, update, CancellationToken.None);
  165. },
  166. stoppingToken);
  167. var now = DateTime.Now;
  168. job.Status = ModuleRebuildStatus.Success;
  169. job.CurrentStage = ModuleRebuildStages.Success;
  170. job.StageIndex = ModuleRebuildStages.StageTotal(scope.ModuleCode);
  171. job.ProgressPercent = 100;
  172. job.ProgressMessage = $"{scope.ModuleCode} 数据重算已完成";
  173. job.FailedStage = null;
  174. job.LastProgressAt = now;
  175. job.BatchId = result.BatchId;
  176. job.TransformRunLogId = result.RunLogId;
  177. job.StageRows = result.StageRows;
  178. job.StandardRows = result.StandardRows;
  179. job.DwdRows = result.DwdRows;
  180. job.KpiRows = result.KpiRows;
  181. job.AtomicRows = result.AtomicRows;
  182. job.DetailJson = result.DetailJson;
  183. job.FinishedAt = now;
  184. job.HeartbeatAt = now;
  185. job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null;
  186. job.UpdateTime = now;
  187. job.ErrorMessage = null;
  188. await _jobs.UpdateAsync(job, CancellationToken.None);
  189. }
  190. finally
  191. {
  192. lockHeartbeatCts.Cancel();
  193. try { await lockHeartbeat; } catch (OperationCanceledException) { }
  194. }
  195. }
  196. }
  197. catch (ModuleRebuildAlreadyRunningException)
  198. {
  199. var now = DateTime.Now;
  200. job.Status = ModuleRebuildStatus.Queued;
  201. job.CurrentStage = ModuleRebuildStages.AcquiringLock;
  202. job.StageIndex = 0;
  203. job.ProgressPercent = 2;
  204. job.ProgressMessage = $"等待现有 {scope.ModuleCode} 全量任务完成";
  205. job.StartedAt = null;
  206. job.HeartbeatAt = now;
  207. job.LastProgressAt = now;
  208. job.UpdateTime = now;
  209. job.ErrorMessage = null;
  210. await _jobs.UpdateAsync(job, CancellationToken.None);
  211. _queue.Pulse();
  212. }
  213. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  214. {
  215. await FailJobAsync(job, "服务停止,任务已取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage);
  216. }
  217. catch (Exception ex)
  218. {
  219. await FailJobAsync(
  220. job,
  221. Truncate(ex.Message, 2000),
  222. ModuleRebuildStatus.Failed,
  223. ModuleRebuildStages.Failed,
  224. failedStage);
  225. }
  226. }
  227. private async Task FailJobAsync(AdoModuleDashboardRebuildJob job, string message, string status, string currentStage, string failedStage)
  228. {
  229. var now = DateTime.Now;
  230. job.Status = status;
  231. job.CurrentStage = currentStage;
  232. job.FailedStage = status == ModuleRebuildStatus.Failed ? failedStage : job.FailedStage;
  233. job.ProgressMessage = status == ModuleRebuildStatus.Failed
  234. ? $"在【{ModuleRebuildStages.ToChinese(failedStage)}】失败"
  235. : message;
  236. job.FinishedAt = now;
  237. job.HeartbeatAt = now;
  238. job.LastProgressAt = now;
  239. job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null;
  240. job.ErrorMessage = message;
  241. job.UpdateTime = now;
  242. await _jobs.UpdateAsync(job, CancellationToken.None);
  243. }
  244. public static ModuleRebuildJobDto ToDto(AdoModuleDashboardRebuildJob row) => new()
  245. {
  246. Ok = true,
  247. ModuleCode = row.ModuleCode,
  248. JobId = row.Id,
  249. Status = row.Status,
  250. CurrentStage = row.CurrentStage,
  251. StageIndex = row.StageIndex,
  252. StageTotal = row.StageTotal <= 0 ? ModuleRebuildStages.StageTotal(row.ModuleCode) : row.StageTotal,
  253. ProgressPercent = row.ProgressPercent,
  254. ProgressMessage = row.ProgressMessage,
  255. LastProgressAt = row.LastProgressAt,
  256. HeartbeatAt = row.HeartbeatAt,
  257. FailedStage = row.FailedStage,
  258. SubmittedAt = row.SubmittedAt,
  259. StartedAt = row.StartedAt,
  260. FinishedAt = row.FinishedAt,
  261. DurationMs = row.DurationMs,
  262. BatchId = row.BatchId,
  263. StageRows = row.StageRows,
  264. StandardRows = row.StandardRows,
  265. DwdRows = row.DwdRows,
  266. KpiRows = row.KpiRows,
  267. AtomicRows = row.AtomicRows,
  268. DetailJson = row.DetailJson,
  269. ErrorMessage = row.ErrorMessage
  270. };
  271. private static ModuleRebuildJobAccepted Conflict(string moduleCode, AdoModuleDashboardRebuildJob active) => new()
  272. {
  273. Ok = false,
  274. ModuleCode = moduleCode,
  275. JobId = active.Id,
  276. Status = active.Status,
  277. Message = $"{moduleCode} 数据重算正在执行,请勿重复提交"
  278. };
  279. private static async Task KeepLeaseAliveAsync(IModuleRebuildLease lease, CancellationToken ct)
  280. {
  281. while (!ct.IsCancellationRequested)
  282. {
  283. try
  284. {
  285. await Task.Delay(TimeSpan.FromSeconds(30), ct);
  286. await lease.HeartbeatAsync(CancellationToken.None);
  287. }
  288. catch (OperationCanceledException)
  289. {
  290. return;
  291. }
  292. catch
  293. {
  294. // 单例 SqlSugar 连接被业务阶段占用时跳过本次锁心跳,避免打断 MDP。
  295. }
  296. }
  297. }
  298. private static int ToDurationMs(TimeSpan elapsed)
  299. {
  300. var ms = elapsed.TotalMilliseconds;
  301. if (double.IsNaN(ms) || ms <= 0) return 0;
  302. return ms >= int.MaxValue ? int.MaxValue : (int)ms;
  303. }
  304. private static string Truncate(string value, int max) =>
  305. string.IsNullOrEmpty(value) || value.Length <= max ? value : value[..max];
  306. }