EtlInstanceRegistrar.cs 9.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207
  1. using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
  2. using Microsoft.Extensions.DependencyInjection;
  3. using Microsoft.Extensions.Hosting;
  4. using Microsoft.Extensions.Logging;
  5. namespace Admin.NET.Plugin.AiDOP.Infrastructure;
  6. /// <summary>
  7. /// ETL 执行机实例注册器:把本进程登记进 <c>ado_etl_instance</c> 并定期心跳,
  8. /// 同时把本实例的指派位发布到 <see cref="AidopRunnerState"/> 供 <see cref="AidopJobGate"/> 读取。
  9. ///
  10. /// <para><b>只读指派、从不自行指派</b>:本注册器插入新行时 <c>is_runner</c> 恒为 0。
  11. /// 指派是人工动作(超管在 MDP 运行监控页操作,或直接改库)。
  12. /// 禁止在这里实现「抢锁自荐」或「心跳超时自动接管」——
  13. /// 那会让开发机在某次重启后意外变成执行机,正是本次治理要消除的情形。</para>
  14. ///
  15. /// <para><b>本服务不受 <see cref="AidopJobGate"/> 管辖</b>,必须常开。
  16. /// 它是闸门的数据来源,被闸门关掉就成了死锁:没人心跳 → 没人发布指派 → 永远不是执行机。</para>
  17. /// </summary>
  18. public sealed class EtlInstanceRegistrar : BackgroundService
  19. {
  20. /// <summary>心跳周期。<see cref="AidopRunnerState.FreshnessWindow"/> 按其 3 倍取值。</summary>
  21. public static readonly TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(20);
  22. /// <summary>启动延时:等宿主配置与建表脚本就绪,避免首拍必然失败刷一条 warn。</summary>
  23. private static readonly TimeSpan StartupDelay = TimeSpan.FromSeconds(5);
  24. /// <summary>下线实例行的保留期。超期且未被指派的行会被清理,防止表随重启无限增长。</summary>
  25. private static readonly TimeSpan DeadInstanceRetention = TimeSpan.FromDays(7);
  26. private readonly IServiceScopeFactory _scopeFactory;
  27. private readonly ILogger _logger;
  28. private int _consecutiveFailures;
  29. private int _guardrailsLogged;
  30. public EtlInstanceRegistrar(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  31. {
  32. _scopeFactory = scopeFactory;
  33. _logger = loggerFactory.CreateLogger(nameof(EtlInstanceRegistrar));
  34. }
  35. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  36. {
  37. try { await Task.Delay(StartupDelay, stoppingToken); }
  38. catch (OperationCanceledException) { return; }
  39. _logger.LogInformation(
  40. "[EtlInstanceRegistrar] 本实例 instanceId={InstanceId} machine={Machine} pid={Pid} version={Version}",
  41. AidopInstanceIdentity.InstanceId, AidopInstanceIdentity.MachineName,
  42. AidopInstanceIdentity.ProcessId, AidopInstanceIdentity.AppVersion);
  43. var nextCleanupAt = DateTimeOffset.UtcNow;
  44. while (!stoppingToken.IsCancellationRequested)
  45. {
  46. try
  47. {
  48. var assigned = await HeartbeatAsync(stoppingToken);
  49. AidopRunnerState.Publish(assigned);
  50. if (Interlocked.Exchange(ref _guardrailsLogged, 1) == 0)
  51. await LogMySqlGuardrailsAsync(assigned, stoppingToken);
  52. if (_consecutiveFailures > 0)
  53. {
  54. _logger.LogInformation(
  55. "[EtlInstanceRegistrar] 心跳恢复,连续失败 {Count} 次后已重新发布指派 isRunner={IsRunner}",
  56. _consecutiveFailures, assigned);
  57. _consecutiveFailures = 0;
  58. }
  59. if (DateTimeOffset.UtcNow >= nextCleanupAt)
  60. {
  61. await CleanupDeadAsync(stoppingToken);
  62. nextCleanupAt = DateTimeOffset.UtcNow + TimeSpan.FromHours(6);
  63. }
  64. }
  65. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  66. {
  67. break;
  68. }
  69. catch (Exception ex)
  70. {
  71. // 查库失败必须打回未知,让闸门回退到配置值:
  72. // 指派可能已被运维改掉,继续沿用旧值会让两台实例同时自认执行机。
  73. AidopRunnerState.Invalidate();
  74. // 库不可用时本循环每 20 秒一次,不做节流会刷屏淹没真正的错误。
  75. _consecutiveFailures++;
  76. if (_consecutiveFailures == 1 || _consecutiveFailures % 30 == 0)
  77. _logger.LogWarning(ex,
  78. "[EtlInstanceRegistrar] 心跳失败(连续第 {Count} 次),指派状态已置未知,闸门回退到配置值",
  79. _consecutiveFailures);
  80. }
  81. try { await Task.Delay(HeartbeatInterval, stoppingToken); }
  82. catch (OperationCanceledException) { break; }
  83. }
  84. }
  85. /// <summary>
  86. /// 心跳一拍:更新本实例行(不存在则插入),回读指派位。
  87. /// 返回本实例是否被指派为执行机。
  88. /// </summary>
  89. private async Task<bool> HeartbeatAsync(CancellationToken ct)
  90. {
  91. using var scope = _scopeFactory.CreateScope();
  92. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  93. var instanceId = AidopInstanceIdentity.InstanceId;
  94. var now = DateTime.Now;
  95. // 刻意不更新 is_runner:指派归人工,注册器只报到。
  96. var affected = await db.Updateable<AdoEtlInstance>()
  97. .SetColumns(x => new AdoEtlInstance
  98. {
  99. LastHeartbeatAt = now,
  100. AppVersion = AidopInstanceIdentity.AppVersion,
  101. UpdateTime = now
  102. })
  103. .Where(x => x.InstanceId == instanceId)
  104. .ExecuteCommandAsync(ct);
  105. if (affected <= 0)
  106. {
  107. await db.Insertable(new AdoEtlInstance
  108. {
  109. InstanceId = instanceId,
  110. MachineName = AidopInstanceIdentity.MachineName,
  111. ProcessId = AidopInstanceIdentity.ProcessId,
  112. AppVersion = AidopInstanceIdentity.AppVersion,
  113. IsRunner = false,
  114. StartedAt = AidopInstanceIdentity.StartedAt,
  115. LastHeartbeatAt = now,
  116. CreateTime = now,
  117. UpdateTime = now
  118. }).ExecuteCommandAsync(ct);
  119. }
  120. return await db.Queryable<AdoEtlInstance>()
  121. .Where(x => x.InstanceId == instanceId)
  122. .Select(x => x.IsRunner)
  123. .FirstAsync(ct);
  124. }
  125. /// <summary>
  126. /// 首次心跳成功后打一条现场参数,全进程只打一次。
  127. ///
  128. /// <para><b>为什么要打</b>:2026-09-24 的 1040 事故里,165 的
  129. /// <c>innodb_buffer_pool_size</c> 一直是默认 128MB、<c>max_connections</c> 是默认 151,
  130. /// 而没有任何人能从日志看出这一点——只能事后 SSH 上去查。
  131. /// 把它连同执行机身份打进启动日志,下次排障第一眼就能看到。</para>
  132. ///
  133. /// <para>不打连接串、不打任何口令。失败只 warn,不影响心跳。</para>
  134. /// </summary>
  135. private async Task LogMySqlGuardrailsAsync(bool assigned, CancellationToken ct)
  136. {
  137. try
  138. {
  139. using var scope = _scopeFactory.CreateScope();
  140. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  141. var row = await db.Ado.SqlQuerySingleAsync<MySqlGuardrails>(
  142. "SELECT @@max_connections AS MaxConnections, "
  143. + "@@innodb_buffer_pool_size AS BufferPoolBytes, "
  144. + "DATABASE() AS DbName");
  145. _logger.LogInformation(
  146. "[EtlInstanceRegistrar] 现场参数 db={Db} maxConnections={MaxConn} bufferPoolMB={PoolMb} "
  147. + "isRunner={IsRunner} instanceId={InstanceId}",
  148. row?.DbName, row?.MaxConnections, row?.BufferPoolBytes / 1024 / 1024,
  149. assigned, AidopInstanceIdentity.InstanceId);
  150. }
  151. catch (Exception ex)
  152. {
  153. _logger.LogWarning(ex, "[EtlInstanceRegistrar] 读取现场 MySQL 参数失败,跳过");
  154. }
  155. }
  156. private sealed class MySqlGuardrails
  157. {
  158. public long MaxConnections { get; set; }
  159. public long BufferPoolBytes { get; set; }
  160. public string DbName { get; set; }
  161. }
  162. /// <summary>
  163. /// 清理早已下线的实例行。被指派的那行即使心跳超期也保留——
  164. /// 运维需要看到「指派的机器挂了」,而不是记录凭空消失。
  165. /// </summary>
  166. private async Task CleanupDeadAsync(CancellationToken ct)
  167. {
  168. try
  169. {
  170. using var scope = _scopeFactory.CreateScope();
  171. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  172. var cutoff = DateTime.Now - DeadInstanceRetention;
  173. var removed = await db.Deleteable<AdoEtlInstance>()
  174. .Where(x => !x.IsRunner && x.LastHeartbeatAt < cutoff)
  175. .ExecuteCommandAsync(ct);
  176. if (removed > 0)
  177. _logger.LogInformation("[EtlInstanceRegistrar] 已清理下线实例 {Count} 行", removed);
  178. }
  179. catch (Exception ex)
  180. {
  181. // 清理失败不影响心跳与指派,吞掉即可
  182. _logger.LogWarning(ex, "[EtlInstanceRegistrar] 清理下线实例失败");
  183. }
  184. }
  185. }