ModuleRebuildStore.cs 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339
  1. using System.Threading.Channels;
  2. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  3. using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
  4. using Microsoft.Extensions.Logging;
  5. namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  6. public sealed class ModuleRebuildQueue : ISingleton
  7. {
  8. private readonly Channel<byte> _channel = Channel.CreateBounded<byte>(new BoundedChannelOptions(8)
  9. {
  10. FullMode = BoundedChannelFullMode.DropOldest,
  11. SingleReader = true,
  12. SingleWriter = false
  13. });
  14. public void Pulse() => _channel.Writer.TryWrite(0);
  15. public ChannelReader<byte> Reader => _channel.Reader;
  16. }
  17. public sealed class ModuleRebuildJobAccepted
  18. {
  19. public bool Ok { get; set; }
  20. public string ModuleCode { get; set; } = string.Empty;
  21. public long? JobId { get; set; }
  22. public string Status { get; set; }
  23. public string Message { get; set; }
  24. }
  25. public sealed class ModuleRebuildJobDto
  26. {
  27. public bool Ok { get; set; } = true;
  28. public string ModuleCode { get; set; } = string.Empty;
  29. public long JobId { get; set; }
  30. public string Status { get; set; }
  31. public string CurrentStage { get; set; }
  32. public int StageIndex { get; set; }
  33. public int StageTotal { get; set; }
  34. public int ProgressPercent { get; set; }
  35. public string ProgressMessage { get; set; }
  36. public DateTime? LastProgressAt { get; set; }
  37. public DateTime? HeartbeatAt { get; set; }
  38. public string FailedStage { get; set; }
  39. public DateTime SubmittedAt { get; set; }
  40. public DateTime? StartedAt { get; set; }
  41. public DateTime? FinishedAt { get; set; }
  42. public int? DurationMs { get; set; }
  43. public string BatchId { get; set; }
  44. public int StageRows { get; set; }
  45. public int StandardRows { get; set; }
  46. public int DwdRows { get; set; }
  47. public int KpiRows { get; set; }
  48. public int AtomicRows { get; set; }
  49. public string DetailJson { get; set; }
  50. public string ErrorMessage { get; set; }
  51. public string Message { get; set; }
  52. }
  53. public interface IModuleRebuildJobStore
  54. {
  55. Task<AdoModuleDashboardRebuildJob> InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default);
  56. Task<AdoModuleDashboardRebuildJob> FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default);
  57. /// <summary>
  58. /// 把一条仍在 QUEUED 的任务收口为 CANCELLED,给人工重算腾位。
  59. ///
  60. /// <para>条件更新(<c>WHERE status='QUEUED'</c>):若 Worker 已把它抢成 RUNNING 则返回 <c>false</c>,
  61. /// 调用方必须据此退回 409,不得把正在跑的任务掀掉。</para>
  62. /// </summary>
  63. Task<bool> TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default);
  64. /// <summary>该 scope 最近一条 SUCCESS,用于 AUTO 冷却判定。没有则返回 null。</summary>
  65. Task<AdoModuleDashboardRebuildJob> FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default);
  66. Task<AdoModuleDashboardRebuildJob> GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default);
  67. Task<AdoModuleDashboardRebuildJob> GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default);
  68. /// <param name="enabledModules">本实例启用的模块码集合,空集合直接返回 null。</param>
  69. /// <param name="globalMaxParallelScopes">
  70. /// 全库并发上限:统计**整个库**处于 RUNNING 的 job 条数,达到上限即不再抢占。
  71. /// 与单进程的 <c>MaxParallelScopes</c> 不同,它跨实例生效。
  72. /// </param>
  73. /// <param name="runner">
  74. /// 本实例是否为 ETL 执行机(<c>AidopJobGate.IsRunner</c>)。
  75. /// <c>false</c> 时只领取页面手工任务,不领 <c>AUTO</c> / <c>BOOTSTRAP</c>。
  76. /// <para><b>不要改成「非执行机干脆不启动 Worker」</b>:队列同时承载页面「数据重算」按钮,
  77. /// 停掉 Worker 会让开发机上的手工重算永远停在 QUEUED。</para>
  78. /// </param>
  79. /// <param name="ct">取消令牌。</param>
  80. Task<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(IReadOnlyCollection<string> enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default);
  81. Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default);
  82. Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default);
  83. Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default);
  84. Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default);
  85. Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default);
  86. /// <summary>置协作式取消标志。只对仍在 QUEUED / RUNNING 的行生效,返回是否改到行。</summary>
  87. Task<bool> RequestCancelAsync(long jobId, CancellationToken ct = default);
  88. /// <summary>阶段边界查询取消标志。查不到行视为未取消。</summary>
  89. Task<bool> IsCancelRequestedAsync(long jobId, CancellationToken ct = default);
  90. /// <summary>
  91. /// 找该模块、该租户在本次重算开始后仍为 RUNNING 的 transform run log。
  92. /// job_code 约定为 <c>{module}_MDP_SYNC_TRANSFORM</c>。找不到返回 0。
  93. /// </summary>
  94. Task<long> FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter);
  95. }
  96. public sealed class ModuleRebuildJobStore : IModuleRebuildJobStore, ITransient
  97. {
  98. private readonly ISqlSugarClient _db;
  99. private readonly ILogger _logger;
  100. public ModuleRebuildJobStore(ISqlSugarClient db, ILoggerFactory loggerFactory)
  101. {
  102. _db = db;
  103. _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildJobStore));
  104. }
  105. public async Task<AdoModuleDashboardRebuildJob> InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default)
  106. {
  107. var id = await _db.Insertable(row).ExecuteReturnIdentityAsync();
  108. row.Id = id;
  109. return row;
  110. }
  111. public Task<AdoModuleDashboardRebuildJob> FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) =>
  112. _db.Queryable<AdoModuleDashboardRebuildJob>()
  113. .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
  114. && (x.Status == ModuleRebuildStatus.Queued || x.Status == ModuleRebuildStatus.Running))
  115. .OrderBy(x => x.Id)
  116. .FirstAsync(ct);
  117. public async Task<bool> TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default)
  118. {
  119. var now = DateTime.Now;
  120. var n = await _db.Updateable<AdoModuleDashboardRebuildJob>()
  121. .SetColumns(x => new AdoModuleDashboardRebuildJob
  122. {
  123. Status = ModuleRebuildStatus.Cancelled,
  124. CurrentStage = ModuleRebuildStages.Cancelled,
  125. FinishedAt = now,
  126. ErrorMessage = reason,
  127. UpdateTime = now
  128. })
  129. .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Queued)
  130. .ExecuteCommandAsync(ct);
  131. return n > 0;
  132. }
  133. public Task<AdoModuleDashboardRebuildJob> FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) =>
  134. _db.Queryable<AdoModuleDashboardRebuildJob>()
  135. .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
  136. && x.Status == ModuleRebuildStatus.Success)
  137. .OrderBy(x => x.FinishedAt, OrderByType.Desc)
  138. .FirstAsync(ct);
  139. public Task<AdoModuleDashboardRebuildJob> GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default) =>
  140. _db.Queryable<AdoModuleDashboardRebuildJob>()
  141. .FirstAsync(x => x.Id == id && x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId, ct);
  142. public Task<AdoModuleDashboardRebuildJob> GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) =>
  143. _db.Queryable<AdoModuleDashboardRebuildJob>()
  144. .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId)
  145. .OrderBy(x => x.Id, OrderByType.Desc)
  146. .FirstAsync(ct);
  147. public async Task<bool> RequestCancelAsync(long jobId, CancellationToken ct = default)
  148. {
  149. var affected = await _db.Updateable<AdoModuleDashboardRebuildJob>()
  150. .SetColumns(x => new AdoModuleDashboardRebuildJob { CancelRequestedFlag = true, UpdateTime = DateTime.Now })
  151. .Where(x => x.Id == jobId
  152. && (x.Status == ModuleRebuildStatus.Queued || x.Status == ModuleRebuildStatus.Running))
  153. .ExecuteCommandAsync(ct);
  154. return affected > 0;
  155. }
  156. public async Task<long> FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter)
  157. {
  158. var jobCode = $"{moduleCode}_MDP_SYNC_TRANSFORM";
  159. var id = await _db.Ado.GetLongAsync(
  160. """
  161. SELECT id FROM mdp_transform_run_log
  162. WHERE job_code=@JobCode AND tenant_id=@TenantId AND status='RUNNING' AND start_time>=@StartedAfter
  163. ORDER BY id DESC LIMIT 1
  164. """,
  165. new List<SugarParameter>
  166. {
  167. new("@JobCode", jobCode),
  168. new("@TenantId", tenantId),
  169. new("@StartedAfter", startedAfter.AddMinutes(-1))
  170. });
  171. return id;
  172. }
  173. public async Task<bool> IsCancelRequestedAsync(long jobId, CancellationToken ct = default)
  174. {
  175. var flag = await _db.Queryable<AdoModuleDashboardRebuildJob>()
  176. .Where(x => x.Id == jobId)
  177. .Select(x => x.CancelRequestedFlag)
  178. .FirstAsync(ct);
  179. return flag;
  180. }
  181. /// <summary>跨实例抢占互斥锁键。只保护「数 RUNNING + 置 RUNNING」这一小段,不覆盖重算本体。</summary>
  182. public const string ClaimLockKey = "aidop:mdp:rebuild:claim";
  183. public async Task<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(
  184. IReadOnlyCollection<string> enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default)
  185. {
  186. if (enabledModules == null || enabledModules.Count == 0)
  187. return null;
  188. // 计数与抢占必须在同一把跨实例锁内:否则 N 个实例会在同一瞬间都读到
  189. // 「RUNNING 还没到上限」,各自抢走一个 job,上限形同虚设。
  190. // 非阻塞取锁——抢不到说明别的实例正在抢占,直接让出,5 秒后下一拍再来。
  191. await using var claimGate = await InventoryInboundLockGuard.TryAcquireAsync(
  192. _db, _logger, ClaimLockKey, ct);
  193. if (!claimGate.Acquired)
  194. return null;
  195. var running = await _db.Queryable<AdoModuleDashboardRebuildJob>()
  196. .Where(x => x.Status == ModuleRebuildStatus.Running)
  197. .CountAsync(ct);
  198. if (running >= globalMaxParallelScopes)
  199. return null;
  200. var enabled = enabledModules.ToList();
  201. var row = await _db.Queryable<AdoModuleDashboardRebuildJob>()
  202. .Where(x => x.Status == ModuleRebuildStatus.Queued && enabled.Contains(x.ModuleCode))
  203. // 非执行机只领页面手工任务。用**白名单**而非「排除 AUTO/BOOTSTRAP」:
  204. // 后续新增的自动触发类型(如夜间兜底的 AUTO_NIGHTLY)天然落在白名单外,
  205. // 不会因为有人忘了往黑名单里补一项就被开发机领去跑全量。
  206. .WhereIF(!runner, x => x.TriggerType == ModuleRebuildTriggerType.Manual || x.RequestedBy != null)
  207. .Where(x => SqlFunc.Subqueryable<AdoModuleDashboardRebuildJob>()
  208. .Where(y => y.ModuleCode == x.ModuleCode && y.TenantId == x.TenantId && y.FactoryId == x.FactoryId
  209. && y.Status == ModuleRebuildStatus.Running)
  210. .NotAny())
  211. .OrderBy(x => x.Id)
  212. .FirstAsync(ct);
  213. if (row == null)
  214. return null;
  215. var now = DateTime.Now;
  216. var n = await _db.Updateable<AdoModuleDashboardRebuildJob>()
  217. .SetColumns(x => new AdoModuleDashboardRebuildJob
  218. {
  219. Status = ModuleRebuildStatus.Running,
  220. CurrentStage = ModuleRebuildStages.AcquiringLock,
  221. StageIndex = 0,
  222. ProgressPercent = 2,
  223. ProgressMessage = $"等待现有 {row.ModuleCode} 全量任务完成",
  224. LastProgressAt = now,
  225. StartedAt = now,
  226. HeartbeatAt = now,
  227. UpdateTime = now
  228. })
  229. .Where(x => x.Id == row.Id && x.Status == ModuleRebuildStatus.Queued)
  230. .ExecuteCommandAsync(ct);
  231. if (n <= 0)
  232. return null;
  233. row.Status = ModuleRebuildStatus.Running;
  234. row.CurrentStage = ModuleRebuildStages.AcquiringLock;
  235. row.ProgressPercent = 2;
  236. row.StartedAt = now;
  237. row.HeartbeatAt = now;
  238. row.UpdateTime = now;
  239. return row;
  240. }
  241. public Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) =>
  242. _db.Updateable(row).ExecuteCommandAsync(ct);
  243. public Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default) =>
  244. _db.Updateable<AdoModuleDashboardRebuildJob>()
  245. .SetColumns(x => new AdoModuleDashboardRebuildJob
  246. {
  247. CurrentStage = currentStage,
  248. StageIndex = stageIndex,
  249. ProgressPercent = progressPercent,
  250. ProgressMessage = message,
  251. LastProgressAt = now,
  252. HeartbeatAt = now,
  253. UpdateTime = now
  254. })
  255. .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Running)
  256. .ExecuteCommandAsync(ct);
  257. public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default)
  258. {
  259. var column = completedStage switch
  260. {
  261. ModuleRebuildStages.Staging or ModuleRebuildStages.T8Inbound => "stage_rows",
  262. ModuleRebuildStages.Standard or ModuleRebuildStages.KpiPreparing => "standard_rows",
  263. ModuleRebuildStages.Dwd => "dwd_rows",
  264. ModuleRebuildStages.Kpi or ModuleRebuildStages.KpiCalculating => "kpi_rows",
  265. ModuleRebuildStages.Atomic => "atomic_rows",
  266. _ => null
  267. };
  268. if (column == null)
  269. return UpdateProgressAsync(jobId, completedStage, 0, progressPercent, nextMessage, now, ct);
  270. return _db.Ado.ExecuteCommandAsync(
  271. $"""
  272. UPDATE ado_module_dashboard_rebuild_job
  273. SET `{column}`=@Rows,
  274. progress_percent=@Pct,
  275. progress_message=@Msg,
  276. detail_json=IFNULL(@Detail, detail_json),
  277. last_progress_at=@Now,
  278. heartbeat_at=@Now,
  279. update_time=@Now
  280. WHERE id=@Id AND status='RUNNING'
  281. """,
  282. new SugarParameter("@Rows", rows),
  283. new SugarParameter("@Pct", progressPercent),
  284. new SugarParameter("@Msg", nextMessage),
  285. new SugarParameter("@Detail", (object?)detailJson ?? DBNull.Value),
  286. new SugarParameter("@Now", now),
  287. new SugarParameter("@Id", jobId));
  288. }
  289. public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default) =>
  290. _db.Updateable<AdoModuleDashboardRebuildJob>()
  291. .SetColumns(x => new AdoModuleDashboardRebuildJob { HeartbeatAt = now, UpdateTime = now })
  292. .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Running)
  293. .ExecuteCommandAsync(ct);
  294. public Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default)
  295. {
  296. var cutoff = DateTime.Now - staleAfter;
  297. var now = DateTime.Now;
  298. return _db.Updateable<AdoModuleDashboardRebuildJob>()
  299. .SetColumns(x => new AdoModuleDashboardRebuildJob
  300. {
  301. Status = ModuleRebuildStatus.Failed,
  302. CurrentStage = ModuleRebuildStages.Failed,
  303. FinishedAt = now,
  304. LastProgressAt = now,
  305. UpdateTime = now,
  306. ErrorMessage = "服务中断,任务未正常结束"
  307. })
  308. .Where(x => x.Status == ModuleRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff))
  309. .ExecuteCommandAsync(ct);
  310. }
  311. }