using Admin.NET.Plugin.AiDOP.Entity.SmartOps; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.Infrastructure; /// /// ETL 执行机实例注册器:把本进程登记进 ado_etl_instance 并定期心跳, /// 同时把本实例是否为执行机发布到 供 读取。 /// /// 判定方式:槽位比对。本进程声明自己属于哪个部署槽位 /// (AidopInstanceIdentity.SlotCode),指派记在 ado_etl_runner_designation; /// 两者相等即为执行机。未声明槽位者永不是执行机。 /// 之所以不拿 instance_id 比对,是因为它含进程ID与启动时间, /// 指派挂上去每次重启都会丢——2026-09-28 的 28 条积压即由此而来。 /// /// 只读指派、从不自行指派:本注册器插入新行时 is_runner 恒为 0, /// 也绝不写指派表。指派是人工动作(超管在 MDP 运行监控页操作,或直接改库)。 /// 禁止在这里实现「抢锁自荐」或「心跳超时自动接管」—— /// 那会让开发机在某次重启后意外变成执行机,正是本次治理要消除的情形。 /// /// 同槽位撞车则两边都不跑:若同一槽位另有存活实例(滚动重启的重叠窗口, /// 或有人把两个进程配成同一槽位),本进程发布 false 而不是照常跑。 /// 没有租约就没有仲裁者,此时抢着跑等于双跑全量,那正是 2026-09-24 连接耗尽的成因; /// 代价只是重启重叠的几十秒空窗,定时作业不受影响。 /// /// 本服务不受 管辖,必须常开。 /// 它是闸门的数据来源,被闸门关掉就成了死锁:没人心跳 → 没人发布指派 → 永远不是执行机。 /// public sealed class EtlInstanceRegistrar : BackgroundService { /// 心跳周期。 按其 3 倍取值。 public static readonly TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(20); /// 启动延时:等宿主配置与建表脚本就绪,避免首拍必然失败刷一条 warn。 private static readonly TimeSpan StartupDelay = TimeSpan.FromSeconds(5); /// 下线实例行的保留期。超期即清理,防止表随重启无限增长。 private static readonly TimeSpan DeadInstanceRetention = TimeSpan.FromDays(7); /// 存活判据:连续三拍没心跳即视为下线,与 同宽。 private static readonly TimeSpan LivenessWindow = TimeSpan.FromSeconds(60); private readonly IServiceScopeFactory _scopeFactory; private readonly ILogger _logger; private int _consecutiveFailures; private int _guardrailsLogged; public EtlInstanceRegistrar(IServiceScopeFactory scopeFactory, ILoggerFactory loggerFactory) { _scopeFactory = scopeFactory; _logger = loggerFactory.CreateLogger(nameof(EtlInstanceRegistrar)); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { try { await Task.Delay(StartupDelay, stoppingToken); } catch (OperationCanceledException) { return; } _logger.LogInformation( "[EtlInstanceRegistrar] 本实例 instanceId={InstanceId} machine={Machine} pid={Pid} version={Version} slot={Slot}", AidopInstanceIdentity.InstanceId, AidopInstanceIdentity.MachineName, AidopInstanceIdentity.ProcessId, AidopInstanceIdentity.AppVersion, AidopInstanceIdentity.SlotCode ?? "(未声明)"); var nextCleanupAt = DateTimeOffset.UtcNow; var lastConflictLogged = false; while (!stoppingToken.IsCancellationRequested) { try { var beat = await HeartbeatAsync(stoppingToken); AidopRunnerState.Publish(beat.IsRunner); // 撞车是配置事故或滚动重启重叠,两者都必须让人看见; // 但每 20 秒一条会刷屏,故只在状态翻转时打。 if (beat.SlotConflict != lastConflictLogged) { if (beat.SlotConflict) _logger.LogWarning( "[EtlInstanceRegistrar] 槽位撞车:slot={Slot} 另有 {Count} 个存活实例," + "本实例暂不作为执行机(避免双跑全量)", AidopInstanceIdentity.SlotCode, beat.SameSlotLiveOthers); else _logger.LogInformation( "[EtlInstanceRegistrar] 槽位撞车已解除 slot={Slot} isRunner={IsRunner}", AidopInstanceIdentity.SlotCode, beat.IsRunner); lastConflictLogged = beat.SlotConflict; } if (Interlocked.Exchange(ref _guardrailsLogged, 1) == 0) await LogMySqlGuardrailsAsync(beat.IsRunner, stoppingToken); if (_consecutiveFailures > 0) { _logger.LogInformation( "[EtlInstanceRegistrar] 心跳恢复,连续失败 {Count} 次后已重新发布指派 isRunner={IsRunner}", _consecutiveFailures, beat.IsRunner); _consecutiveFailures = 0; } if (DateTimeOffset.UtcNow >= nextCleanupAt) { await CleanupDeadAsync(stoppingToken); nextCleanupAt = DateTimeOffset.UtcNow + TimeSpan.FromHours(6); } } catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { break; } catch (Exception ex) { // 查库失败必须打回未知,让闸门回退到配置值: // 指派可能已被运维改掉,继续沿用旧值会让两台实例同时自认执行机。 AidopRunnerState.Invalidate(); // 库不可用时本循环每 20 秒一次,不做节流会刷屏淹没真正的错误。 _consecutiveFailures++; if (_consecutiveFailures == 1 || _consecutiveFailures % 30 == 0) _logger.LogWarning(ex, "[EtlInstanceRegistrar] 心跳失败(连续第 {Count} 次),指派状态已置未知,闸门回退到配置值", _consecutiveFailures); } try { await Task.Delay(HeartbeatInterval, stoppingToken); } catch (OperationCanceledException) { break; } } } /// 一拍心跳的结果。 private readonly record struct HeartbeatResult(bool IsRunner, string DesignatedSlot, int SameSlotLiveOthers) { public bool SlotConflict => SameSlotLiveOthers > 0; } /// /// 心跳一拍:更新本实例行(不存在则插入),读取指派槽位并与本进程槽位比对。 /// private async Task HeartbeatAsync(CancellationToken ct) { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var instances = scope.ServiceProvider.GetRequiredService(); var instanceId = AidopInstanceIdentity.InstanceId; var slot = AidopInstanceIdentity.SlotCode; var now = DateTime.Now; // 刻意不更新 is_runner:那已是遗留列,指派在 ado_etl_runner_designation。 var affected = await db.Updateable() .SetColumns(x => new AdoEtlInstance { LastHeartbeatAt = now, AppVersion = AidopInstanceIdentity.AppVersion, SlotCode = slot, UpdateTime = now }) .Where(x => x.InstanceId == instanceId) .ExecuteCommandAsync(ct); if (affected <= 0) { await db.Insertable(new AdoEtlInstance { InstanceId = instanceId, MachineName = AidopInstanceIdentity.MachineName, ProcessId = AidopInstanceIdentity.ProcessId, AppVersion = AidopInstanceIdentity.AppVersion, SlotCode = slot, IsRunner = false, StartedAt = AidopInstanceIdentity.StartedAt, LastHeartbeatAt = now, CreateTime = now, UpdateTime = now }).ExecuteCommandAsync(ct); } var designated = await instances.GetDesignatedSlotAsync(ct); // 没声明槽位就不是候选,连撞车都不必查——省一次查询,也省得把「未声明」 // 误报成撞车(多台未配置的开发机 slot 全是 NULL,NULL 之间不构成同槽位)。 if (string.IsNullOrEmpty(slot) || !string.Equals(slot, designated, StringComparison.Ordinal)) return new HeartbeatResult(false, designated, 0); var liveCutoff = now - LivenessWindow; var others = await db.Queryable() .Where(x => x.SlotCode == slot && x.InstanceId != instanceId && x.LastHeartbeatAt >= liveCutoff) .CountAsync(ct); return new HeartbeatResult(others == 0, designated, others); } /// /// 首次心跳成功后打一条现场参数,全进程只打一次。 /// /// 为什么要打:2026-09-24 的 1040 事故里,165 的 /// innodb_buffer_pool_size 一直是默认 128MB、max_connections 是默认 151, /// 而没有任何人能从日志看出这一点——只能事后 SSH 上去查。 /// 把它连同执行机身份打进启动日志,下次排障第一眼就能看到。 /// /// 不打连接串、不打任何口令。失败只 warn,不影响心跳。 /// private async Task LogMySqlGuardrailsAsync(bool assigned, CancellationToken ct) { try { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var row = await db.Ado.SqlQuerySingleAsync( "SELECT @@max_connections AS MaxConnections, " + "@@innodb_buffer_pool_size AS BufferPoolBytes, " + "DATABASE() AS DbName"); _logger.LogInformation( "[EtlInstanceRegistrar] 现场参数 db={Db} maxConnections={MaxConn} bufferPoolMB={PoolMb} " + "isRunner={IsRunner} instanceId={InstanceId}", row?.DbName, row?.MaxConnections, row?.BufferPoolBytes / 1024 / 1024, assigned, AidopInstanceIdentity.InstanceId); } catch (Exception ex) { _logger.LogWarning(ex, "[EtlInstanceRegistrar] 读取现场 MySQL 参数失败,跳过"); } } private sealed class MySqlGuardrails { public long MaxConnections { get; set; } public long BufferPoolBytes { get; set; } public string DbName { get; set; } } /// /// 清理早已下线的实例行,不分槽位一视同仁。 /// /// 旧模型必须特意保住被指派实例的行,因为指派就寄生在 /// ado_etl_instance.is_runner 上,删行等于删指派。现在指派独立存在 /// ado_etl_runner_designation,这条例外不但没用,还有害:执行机槽位 /// 每重启一次留一行,恰恰是重启最勤的那个槽位无限堆积,页面上还会摆着一排 /// 早已死掉的同槽位实例。运维要看的「指派槽位上没有活实例」由监控页的 /// noLiveRunner 告警给出,不需要靠留尸体来表达。 /// private async Task CleanupDeadAsync(CancellationToken ct) { try { using var scope = _scopeFactory.CreateScope(); var db = scope.ServiceProvider.GetRequiredService(); var cutoff = DateTime.Now - DeadInstanceRetention; var removed = await db.Deleteable() .Where(x => x.LastHeartbeatAt < cutoff) .ExecuteCommandAsync(ct); if (removed > 0) _logger.LogInformation("[EtlInstanceRegistrar] 已清理下线实例 {Count} 行", removed); } catch (Exception ex) { // 清理失败不影响心跳与指派,吞掉即可 _logger.LogWarning(ex, "[EtlInstanceRegistrar] 清理下线实例失败"); } } }