ModuleRebuildWorker.cs 7.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171
  1. using Microsoft.Extensions.DependencyInjection;
  2. using Microsoft.Extensions.Hosting;
  3. using Microsoft.Extensions.Logging;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
  5. /// <summary>
  6. /// 模块看板完整重算的队列消费者。
  7. ///
  8. /// <para><b>不受 <see cref="AidopJobGate"/> 开关管辖,必须常开</b>:它消费的队列既来自定时作业,
  9. /// 也来自页面上的「数据重算」按钮;整体关掉会让手工触发永远停在 QUEUED。
  10. /// 跨实例重复消费由两道闸收口:<c>ClaimNextQueuedAsync</c> 的全局并发上限(<c>GET_LOCK</c> +
  11. /// <c>GlobalMaxParallelScopes</c>),以及按执行机身份过滤领取范围——
  12. /// 非执行机只领手工任务,不领 <c>AUTO</c> / <c>BOOTSTRAP</c>。</para>
  13. /// </summary>
  14. public sealed class ModuleRebuildWorker : BackgroundService
  15. {
  16. private static readonly TimeSpan IdleDelay = TimeSpan.FromSeconds(5);
  17. private static readonly TimeSpan StaleSweepInterval = TimeSpan.FromMinutes(1);
  18. private readonly IServiceScopeFactory _scopeFactory;
  19. private readonly ModuleRebuildQueue _queue;
  20. private readonly ILogger _logger;
  21. public ModuleRebuildWorker(
  22. IServiceScopeFactory scopeFactory,
  23. ModuleRebuildQueue queue,
  24. ILoggerFactory loggerFactory)
  25. {
  26. _scopeFactory = scopeFactory;
  27. _queue = queue;
  28. _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildWorker));
  29. }
  30. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  31. {
  32. try { await Task.Delay(TimeSpan.FromSeconds(8), stoppingToken); }
  33. catch (OperationCanceledException) { return; }
  34. await FailStaleAsync(stoppingToken);
  35. var running = new List<Task>();
  36. var nextStaleSweepAt = DateTimeOffset.UtcNow + StaleSweepInterval;
  37. while (!stoppingToken.IsCancellationRequested)
  38. {
  39. try
  40. {
  41. if (DateTimeOffset.UtcNow >= nextStaleSweepAt)
  42. {
  43. await FailStaleAsync(stoppingToken);
  44. nextStaleSweepAt = DateTimeOffset.UtcNow + StaleSweepInterval;
  45. }
  46. await FillSlotsAsync(running, stoppingToken);
  47. }
  48. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  49. {
  50. break;
  51. }
  52. catch (Exception ex)
  53. {
  54. _logger.LogWarning(ex, "[ModuleRebuildWorker] run failed");
  55. }
  56. running.RemoveAll(t => t.IsCompleted);
  57. if (running.Count == 0)
  58. {
  59. try
  60. {
  61. using var linked = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
  62. linked.CancelAfter(IdleDelay);
  63. try { await _queue.Reader.ReadAsync(linked.Token); }
  64. catch (OperationCanceledException) when (!stoppingToken.IsCancellationRequested) { }
  65. }
  66. catch (OperationCanceledException) { break; }
  67. }
  68. else
  69. {
  70. var pulse = WaitPulseAsync(stoppingToken);
  71. var completed = await Task.WhenAny(running.Append(pulse));
  72. if (completed != pulse)
  73. running.Remove(completed);
  74. }
  75. }
  76. try { await Task.WhenAll(running); }
  77. catch (Exception ex) { _logger.LogWarning(ex, "[ModuleRebuildWorker] drain failed"); }
  78. }
  79. private async Task FillSlotsAsync(List<Task> running, CancellationToken stoppingToken)
  80. {
  81. int max = 2;
  82. try
  83. {
  84. using var probe = _scopeFactory.CreateScope();
  85. max = probe.ServiceProvider.GetRequiredService<IModuleRebuildCapability>().MaxParallelScopes;
  86. }
  87. catch
  88. {
  89. max = 2;
  90. }
  91. while (running.Count < max && !stoppingToken.IsCancellationRequested)
  92. {
  93. var started = await TryStartOneAsync(stoppingToken);
  94. if (started == null)
  95. break;
  96. running.Add(started);
  97. }
  98. }
  99. private async Task<Task> TryStartOneAsync(CancellationToken stoppingToken)
  100. {
  101. using var probe = _scopeFactory.CreateScope();
  102. var capability = probe.ServiceProvider.GetRequiredService<IModuleRebuildCapability>();
  103. var enabled = MdpRebuildScope.RebuildModules.Where(capability.IsEnabled).ToArray();
  104. if (enabled.Length == 0)
  105. return null;
  106. // 逐拍求值而非启动时一次:执行机指派可在运行中被超管切换,不需重启进程。
  107. // IsRunner 是纯内存读(指派位由 EtlInstanceRegistrar 心跳时发布),
  108. // 本方法每 5 秒进一次,不会因此多产生任何查询。
  109. var runner = AidopJobGate.IsRunner;
  110. var job = await probe.ServiceProvider.GetRequiredService<IModuleRebuildJobStore>()
  111. .ClaimNextQueuedAsync(enabled, capability.GlobalMaxParallelScopes, runner, stoppingToken);
  112. if (job == null)
  113. return null;
  114. return ExecuteClaimedAsync(job, stoppingToken);
  115. }
  116. private async Task ExecuteClaimedAsync(Admin.NET.Plugin.AiDOP.Entity.SmartOps.AdoModuleDashboardRebuildJob job, CancellationToken stoppingToken)
  117. {
  118. using var scope = _scopeFactory.CreateScope();
  119. var svc = scope.ServiceProvider.GetRequiredService<ModuleRebuildService>();
  120. var runLock = scope.ServiceProvider.GetRequiredService<IModuleRebuildLock>();
  121. var handler = scope.ServiceProvider.GetServices<IModuleRebuildHandler>()
  122. .FirstOrDefault(x => string.Equals(x.ModuleCode, job.ModuleCode, StringComparison.OrdinalIgnoreCase));
  123. if (handler == null)
  124. {
  125. job.Status = ModuleRebuildStatus.Failed;
  126. job.CurrentStage = ModuleRebuildStages.Failed;
  127. job.ErrorMessage = $"未注册 {job.ModuleCode} 重算处理器";
  128. job.FinishedAt = DateTime.Now;
  129. job.UpdateTime = DateTime.Now;
  130. await scope.ServiceProvider.GetRequiredService<IModuleRebuildJobStore>().UpdateAsync(job, CancellationToken.None);
  131. return;
  132. }
  133. // 进度写入已带 heartbeat;后台另开连接会与单例 ISqlSugarClient 抢同一条 MySQL 连接。
  134. await svc.RunClaimedAsync(job, handler, runLock, stoppingToken);
  135. }
  136. private async Task WaitPulseAsync(CancellationToken stoppingToken)
  137. {
  138. using var linked = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
  139. linked.CancelAfter(IdleDelay);
  140. try { await _queue.Reader.ReadAsync(linked.Token); }
  141. catch (OperationCanceledException) { }
  142. }
  143. private async Task FailStaleAsync(CancellationToken stoppingToken)
  144. {
  145. try
  146. {
  147. using var scope = _scopeFactory.CreateScope();
  148. await scope.ServiceProvider.GetRequiredService<ModuleRebuildService>().FailStaleAsync(stoppingToken);
  149. }
  150. catch (Exception ex)
  151. {
  152. _logger.LogWarning(ex, "[ModuleRebuildWorker] fail-stale failed");
  153. }
  154. }
  155. }