EtlInstanceRegistrar.cs 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263
  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>。本进程声明自己属于哪个部署槽位
  11. /// (<c>AidopInstanceIdentity.SlotCode</c>),指派记在 <c>ado_etl_runner_designation</c>;
  12. /// 两者相等即为执行机。<b>未声明槽位者永不是执行机</b>。
  13. /// 之所以不拿 <c>instance_id</c> 比对,是因为它含进程ID与启动时间,
  14. /// 指派挂上去每次重启都会丢——2026-09-28 的 28 条积压即由此而来。</para>
  15. ///
  16. /// <para><b>只读指派、从不自行指派</b>:本注册器插入新行时 <c>is_runner</c> 恒为 0,
  17. /// 也绝不写指派表。指派是人工动作(超管在 MDP 运行监控页操作,或直接改库)。
  18. /// 禁止在这里实现「抢锁自荐」或「心跳超时自动接管」——
  19. /// 那会让开发机在某次重启后意外变成执行机,正是本次治理要消除的情形。</para>
  20. ///
  21. /// <para><b>同槽位撞车则两边都不跑</b>:若同一槽位另有存活实例(滚动重启的重叠窗口,
  22. /// 或有人把两个进程配成同一槽位),本进程发布 false 而不是照常跑。
  23. /// 没有租约就没有仲裁者,此时抢着跑等于双跑全量,那正是 2026-09-24 连接耗尽的成因;
  24. /// 代价只是重启重叠的几十秒空窗,定时作业不受影响。</para>
  25. ///
  26. /// <para><b>本服务不受 <see cref="AidopJobGate"/> 管辖</b>,必须常开。
  27. /// 它是闸门的数据来源,被闸门关掉就成了死锁:没人心跳 → 没人发布指派 → 永远不是执行机。</para>
  28. /// </summary>
  29. public sealed class EtlInstanceRegistrar : BackgroundService
  30. {
  31. /// <summary>心跳周期。<see cref="AidopRunnerState.FreshnessWindow"/> 按其 3 倍取值。</summary>
  32. public static readonly TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(20);
  33. /// <summary>启动延时:等宿主配置与建表脚本就绪,避免首拍必然失败刷一条 warn。</summary>
  34. private static readonly TimeSpan StartupDelay = TimeSpan.FromSeconds(5);
  35. /// <summary>下线实例行的保留期。超期即清理,防止表随重启无限增长。</summary>
  36. private static readonly TimeSpan DeadInstanceRetention = TimeSpan.FromDays(7);
  37. /// <summary>存活判据:连续三拍没心跳即视为下线,与 <see cref="AidopRunnerState.FreshnessWindow"/> 同宽。</summary>
  38. private static readonly TimeSpan LivenessWindow = TimeSpan.FromSeconds(60);
  39. private readonly IServiceScopeFactory _scopeFactory;
  40. private readonly ILogger _logger;
  41. private int _consecutiveFailures;
  42. private int _guardrailsLogged;
  43. public EtlInstanceRegistrar(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory)
  44. {
  45. _scopeFactory = scopeFactory;
  46. _logger = loggerFactory.CreateLogger(nameof(EtlInstanceRegistrar));
  47. }
  48. protected override async Task ExecuteAsync(CancellationToken stoppingToken)
  49. {
  50. try { await Task.Delay(StartupDelay, stoppingToken); }
  51. catch (OperationCanceledException) { return; }
  52. _logger.LogInformation(
  53. "[EtlInstanceRegistrar] 本实例 instanceId={InstanceId} machine={Machine} pid={Pid} version={Version} slot={Slot}",
  54. AidopInstanceIdentity.InstanceId, AidopInstanceIdentity.MachineName,
  55. AidopInstanceIdentity.ProcessId, AidopInstanceIdentity.AppVersion,
  56. AidopInstanceIdentity.SlotCode ?? "(未声明)");
  57. var nextCleanupAt = DateTimeOffset.UtcNow;
  58. var lastConflictLogged = false;
  59. while (!stoppingToken.IsCancellationRequested)
  60. {
  61. try
  62. {
  63. var beat = await HeartbeatAsync(stoppingToken);
  64. AidopRunnerState.Publish(beat.IsRunner);
  65. // 撞车是配置事故或滚动重启重叠,两者都必须让人看见;
  66. // 但每 20 秒一条会刷屏,故只在状态翻转时打。
  67. if (beat.SlotConflict != lastConflictLogged)
  68. {
  69. if (beat.SlotConflict)
  70. _logger.LogWarning(
  71. "[EtlInstanceRegistrar] 槽位撞车:slot={Slot} 另有 {Count} 个存活实例,"
  72. + "本实例暂不作为执行机(避免双跑全量)",
  73. AidopInstanceIdentity.SlotCode, beat.SameSlotLiveOthers);
  74. else
  75. _logger.LogInformation(
  76. "[EtlInstanceRegistrar] 槽位撞车已解除 slot={Slot} isRunner={IsRunner}",
  77. AidopInstanceIdentity.SlotCode, beat.IsRunner);
  78. lastConflictLogged = beat.SlotConflict;
  79. }
  80. if (Interlocked.Exchange(ref _guardrailsLogged, 1) == 0)
  81. await LogMySqlGuardrailsAsync(beat.IsRunner, stoppingToken);
  82. if (_consecutiveFailures > 0)
  83. {
  84. _logger.LogInformation(
  85. "[EtlInstanceRegistrar] 心跳恢复,连续失败 {Count} 次后已重新发布指派 isRunner={IsRunner}",
  86. _consecutiveFailures, beat.IsRunner);
  87. _consecutiveFailures = 0;
  88. }
  89. if (DateTimeOffset.UtcNow >= nextCleanupAt)
  90. {
  91. await CleanupDeadAsync(stoppingToken);
  92. nextCleanupAt = DateTimeOffset.UtcNow + TimeSpan.FromHours(6);
  93. }
  94. }
  95. catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
  96. {
  97. break;
  98. }
  99. catch (Exception ex)
  100. {
  101. // 查库失败必须打回未知,让闸门回退到配置值:
  102. // 指派可能已被运维改掉,继续沿用旧值会让两台实例同时自认执行机。
  103. AidopRunnerState.Invalidate();
  104. // 库不可用时本循环每 20 秒一次,不做节流会刷屏淹没真正的错误。
  105. _consecutiveFailures++;
  106. if (_consecutiveFailures == 1 || _consecutiveFailures % 30 == 0)
  107. _logger.LogWarning(ex,
  108. "[EtlInstanceRegistrar] 心跳失败(连续第 {Count} 次),指派状态已置未知,闸门回退到配置值",
  109. _consecutiveFailures);
  110. }
  111. try { await Task.Delay(HeartbeatInterval, stoppingToken); }
  112. catch (OperationCanceledException) { break; }
  113. }
  114. }
  115. /// <summary>一拍心跳的结果。</summary>
  116. private readonly record struct HeartbeatResult(bool IsRunner, string DesignatedSlot, int SameSlotLiveOthers)
  117. {
  118. public bool SlotConflict => SameSlotLiveOthers > 0;
  119. }
  120. /// <summary>
  121. /// 心跳一拍:更新本实例行(不存在则插入),读取指派槽位并与本进程槽位比对。
  122. /// </summary>
  123. private async Task<HeartbeatResult> HeartbeatAsync(CancellationToken ct)
  124. {
  125. using var scope = _scopeFactory.CreateScope();
  126. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  127. var instances = scope.ServiceProvider.GetRequiredService<IEtlInstanceStore>();
  128. var instanceId = AidopInstanceIdentity.InstanceId;
  129. var slot = AidopInstanceIdentity.SlotCode;
  130. var now = DateTime.Now;
  131. // 刻意不更新 is_runner:那已是遗留列,指派在 ado_etl_runner_designation。
  132. var affected = await db.Updateable<AdoEtlInstance>()
  133. .SetColumns(x => new AdoEtlInstance
  134. {
  135. LastHeartbeatAt = now,
  136. AppVersion = AidopInstanceIdentity.AppVersion,
  137. SlotCode = slot,
  138. UpdateTime = now
  139. })
  140. .Where(x => x.InstanceId == instanceId)
  141. .ExecuteCommandAsync(ct);
  142. if (affected <= 0)
  143. {
  144. await db.Insertable(new AdoEtlInstance
  145. {
  146. InstanceId = instanceId,
  147. MachineName = AidopInstanceIdentity.MachineName,
  148. ProcessId = AidopInstanceIdentity.ProcessId,
  149. AppVersion = AidopInstanceIdentity.AppVersion,
  150. SlotCode = slot,
  151. IsRunner = false,
  152. StartedAt = AidopInstanceIdentity.StartedAt,
  153. LastHeartbeatAt = now,
  154. CreateTime = now,
  155. UpdateTime = now
  156. }).ExecuteCommandAsync(ct);
  157. }
  158. var designated = await instances.GetDesignatedSlotAsync(ct);
  159. // 没声明槽位就不是候选,连撞车都不必查——省一次查询,也省得把「未声明」
  160. // 误报成撞车(多台未配置的开发机 slot 全是 NULL,NULL 之间不构成同槽位)。
  161. if (string.IsNullOrEmpty(slot) || !string.Equals(slot, designated, StringComparison.Ordinal))
  162. return new HeartbeatResult(false, designated, 0);
  163. var liveCutoff = now - LivenessWindow;
  164. var others = await db.Queryable<AdoEtlInstance>()
  165. .Where(x => x.SlotCode == slot && x.InstanceId != instanceId && x.LastHeartbeatAt >= liveCutoff)
  166. .CountAsync(ct);
  167. return new HeartbeatResult(others == 0, designated, others);
  168. }
  169. /// <summary>
  170. /// 首次心跳成功后打一条现场参数,全进程只打一次。
  171. ///
  172. /// <para><b>为什么要打</b>:2026-09-24 的 1040 事故里,165 的
  173. /// <c>innodb_buffer_pool_size</c> 一直是默认 128MB、<c>max_connections</c> 是默认 151,
  174. /// 而没有任何人能从日志看出这一点——只能事后 SSH 上去查。
  175. /// 把它连同执行机身份打进启动日志,下次排障第一眼就能看到。</para>
  176. ///
  177. /// <para>不打连接串、不打任何口令。失败只 warn,不影响心跳。</para>
  178. /// </summary>
  179. private async Task LogMySqlGuardrailsAsync(bool assigned, CancellationToken ct)
  180. {
  181. try
  182. {
  183. using var scope = _scopeFactory.CreateScope();
  184. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  185. var row = await db.Ado.SqlQuerySingleAsync<MySqlGuardrails>(
  186. "SELECT @@max_connections AS MaxConnections, "
  187. + "@@innodb_buffer_pool_size AS BufferPoolBytes, "
  188. + "DATABASE() AS DbName");
  189. _logger.LogInformation(
  190. "[EtlInstanceRegistrar] 现场参数 db={Db} maxConnections={MaxConn} bufferPoolMB={PoolMb} "
  191. + "isRunner={IsRunner} instanceId={InstanceId}",
  192. row?.DbName, row?.MaxConnections, row?.BufferPoolBytes / 1024 / 1024,
  193. assigned, AidopInstanceIdentity.InstanceId);
  194. }
  195. catch (Exception ex)
  196. {
  197. _logger.LogWarning(ex, "[EtlInstanceRegistrar] 读取现场 MySQL 参数失败,跳过");
  198. }
  199. }
  200. private sealed class MySqlGuardrails
  201. {
  202. public long MaxConnections { get; set; }
  203. public long BufferPoolBytes { get; set; }
  204. public string DbName { get; set; }
  205. }
  206. /// <summary>
  207. /// 清理早已下线的实例行,不分槽位一视同仁。
  208. ///
  209. /// <para>旧模型必须特意保住被指派实例的行,因为指派就寄生在
  210. /// <c>ado_etl_instance.is_runner</c> 上,删行等于删指派。现在指派独立存在
  211. /// <c>ado_etl_runner_designation</c>,这条例外不但没用,还有害:执行机槽位
  212. /// 每重启一次留一行,恰恰是重启最勤的那个槽位无限堆积,页面上还会摆着一排
  213. /// 早已死掉的同槽位实例。运维要看的「指派槽位上没有活实例」由监控页的
  214. /// noLiveRunner 告警给出,不需要靠留尸体来表达。</para>
  215. /// </summary>
  216. private async Task CleanupDeadAsync(CancellationToken ct)
  217. {
  218. try
  219. {
  220. using var scope = _scopeFactory.CreateScope();
  221. var db = scope.ServiceProvider.GetRequiredService<ISqlSugarClient>();
  222. var cutoff = DateTime.Now - DeadInstanceRetention;
  223. var removed = await db.Deleteable<AdoEtlInstance>()
  224. .Where(x => x.LastHeartbeatAt < cutoff)
  225. .ExecuteCommandAsync(ct);
  226. if (removed > 0)
  227. _logger.LogInformation("[EtlInstanceRegistrar] 已清理下线实例 {Count} 行", removed);
  228. }
  229. catch (Exception ex)
  230. {
  231. // 清理失败不影响心跳与指派,吞掉即可
  232. _logger.LogWarning(ex, "[EtlInstanceRegistrar] 清理下线实例失败");
  233. }
  234. }
  235. }