| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409 |
- using System.Threading.Channels;
- using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
- using Admin.NET.Plugin.AiDOP.MaterialWarehouse;
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
- public sealed class ModuleRebuildQueue : ISingleton
- {
- private readonly Channel<byte> _channel = Channel.CreateBounded<byte>(new BoundedChannelOptions(8)
- {
- FullMode = BoundedChannelFullMode.DropOldest,
- SingleReader = true,
- SingleWriter = false
- });
- public void Pulse() => _channel.Writer.TryWrite(0);
- public ChannelReader<byte> Reader => _channel.Reader;
- }
- public sealed class ModuleRebuildJobAccepted
- {
- public bool Ok { get; set; }
- public string ModuleCode { get; set; } = string.Empty;
- public long? JobId { get; set; }
- public string Status { get; set; }
- public string Message { get; set; }
- }
- public sealed class ModuleRebuildJobDto
- {
- public bool Ok { get; set; } = true;
- public string ModuleCode { get; set; } = string.Empty;
- public long JobId { get; set; }
- public string Status { get; set; }
- public string CurrentStage { get; set; }
- public int StageIndex { get; set; }
- public int StageTotal { get; set; }
- public int ProgressPercent { get; set; }
- public string ProgressMessage { get; set; }
- public DateTime? LastProgressAt { get; set; }
- public DateTime? HeartbeatAt { get; set; }
- public string FailedStage { get; set; }
- public DateTime SubmittedAt { get; set; }
- public DateTime? StartedAt { get; set; }
- public DateTime? FinishedAt { get; set; }
- public int? DurationMs { get; set; }
- public string BatchId { get; set; }
- public int StageRows { get; set; }
- public int StandardRows { get; set; }
- public int DwdRows { get; set; }
- public int KpiRows { get; set; }
- public int AtomicRows { get; set; }
- public string DetailJson { get; set; }
- public string ErrorMessage { get; set; }
- public string Message { get; set; }
- }
- public interface IModuleRebuildJobStore
- {
- Task<AdoModuleDashboardRebuildJob> InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default);
- Task<AdoModuleDashboardRebuildJob> FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default);
- /// <summary>
- /// 把一条仍在 QUEUED 的任务收口为 CANCELLED,给人工重算腾位。
- ///
- /// <para>条件更新(<c>WHERE status='QUEUED'</c>):若 Worker 已把它抢成 RUNNING 则返回 <c>false</c>,
- /// 调用方必须据此退回 409,不得把正在跑的任务掀掉。</para>
- /// </summary>
- Task<bool> TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default);
- /// <summary>
- /// 批量把某个自动触发类型下**仍未开跑**的任务收口为 CANCELLED,返回收口行数。
- ///
- /// <para><b>只碰 <c>QUEUED</c> 且 <c>requested_by IS NULL</c> 的行</b>:
- /// 条件里带 <c>trigger_type</c> 是为了让调用方只清自己那一类,
- /// 带 <c>requested_by IS NULL</c> 是为了绝不误伤页面手工任务
- /// (<c>ModuleRebuildTriggerType.IsManualTrigger</c> 把带 <c>requested_by</c> 的行也算人工)。
- /// RUNNING 一律不碰。</para>
- ///
- /// <para><b>刻意跨租户</b>:调用方是夜间全量作业,它本来就一次性扇出全部启用 scope,
- /// 清的也正是自己上一轮那一批。不要给这里加 tenant 过滤,否则只能清掉一个租户。</para>
- /// </summary>
- /// <param name="triggerType">要清理的触发类型,如 <c>AUTO_NIGHTLY</c>。</param>
- /// <param name="olderThan">仅清理入队时间早于 now - olderThan 的行;<c>TimeSpan.Zero</c> 表示不限。</param>
- Task<int> CancelUnstartedByTriggerAsync(string triggerType, TimeSpan olderThan, string reason, CancellationToken ct = default);
- /// <summary>
- /// 把长期滞留在 QUEUED、已无消费可能的任务收口为 FAILED,返回收口行数。
- ///
- /// <para><b>为什么 <see cref="FailStaleRunningAsync"/> 不够</b>:它只管 RUNNING。
- /// QUEUED 的滞留有独立成因——非执行机只领手工任务,所以指派一丢,
- /// 自动任务在队列里没有任何消费者,却会一直挡着同 scope 的后续入队。</para>
- ///
- /// <para><b>阈值必须给足</b>:<c>GlobalMaxParallelScopes = 1</c> 且单 scope 跑 10~20 分钟,
- /// 夜间一次扇出 28 个 scope 时,最后一个**合法**等待可达 7 小时左右。
- /// 阈值定短了会把正常排队的任务误杀,见 <c>ModuleRebuildLock.QueuedStaleAfter</c>。</para>
- /// </summary>
- Task<int> FailStaleQueuedAsync(TimeSpan staleAfter, string reason, CancellationToken ct = default);
- /// <summary>该 scope 最近一条 SUCCESS,用于 AUTO 冷却判定。没有则返回 null。</summary>
- Task<AdoModuleDashboardRebuildJob> FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default);
- Task<AdoModuleDashboardRebuildJob> GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default);
- Task<AdoModuleDashboardRebuildJob> GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default);
- /// <param name="enabledModules">本实例启用的模块码集合,空集合直接返回 null。</param>
- /// <param name="globalMaxParallelScopes">
- /// 全库并发上限:统计**整个库**处于 RUNNING 的 job 条数,达到上限即不再抢占。
- /// 与单进程的 <c>MaxParallelScopes</c> 不同,它跨实例生效。
- /// </param>
- /// <param name="runner">
- /// 本实例是否为 ETL 执行机(<c>AidopJobGate.IsRunner</c>)。
- /// <c>false</c> 时只领取页面手工任务,不领 <c>AUTO</c> / <c>BOOTSTRAP</c>。
- /// <para><b>不要改成「非执行机干脆不启动 Worker」</b>:队列同时承载页面「数据重算」按钮,
- /// 停掉 Worker 会让开发机上的手工重算永远停在 QUEUED。</para>
- /// </param>
- /// <param name="ct">取消令牌。</param>
- Task<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(IReadOnlyCollection<string> enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default);
- Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default);
- Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default);
- Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default);
- Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default);
- Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default);
- /// <summary>置协作式取消标志。只对仍在 QUEUED / RUNNING 的行生效,返回是否改到行。</summary>
- Task<bool> RequestCancelAsync(long jobId, CancellationToken ct = default);
- /// <summary>阶段边界查询取消标志。查不到行视为未取消。</summary>
- Task<bool> IsCancelRequestedAsync(long jobId, CancellationToken ct = default);
- /// <summary>
- /// 找该模块、该租户在本次重算开始后仍为 RUNNING 的 transform run log。
- /// job_code 约定为 <c>{module}_MDP_SYNC_TRANSFORM</c>。找不到返回 0。
- /// </summary>
- Task<long> FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter);
- }
- public sealed class ModuleRebuildJobStore : IModuleRebuildJobStore, ITransient
- {
- private readonly ISqlSugarClient _db;
- private readonly ILogger _logger;
- public ModuleRebuildJobStore(ISqlSugarClient db, ILoggerFactory loggerFactory)
- {
- _db = db;
- _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildJobStore));
- }
- public async Task<AdoModuleDashboardRebuildJob> InsertQueuedAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default)
- {
- var id = await _db.Insertable(row).ExecuteReturnIdentityAsync();
- row.Id = id;
- return row;
- }
- public Task<AdoModuleDashboardRebuildJob> FindActiveAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) =>
- _db.Queryable<AdoModuleDashboardRebuildJob>()
- .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
- && (x.Status == ModuleRebuildStatus.Queued || x.Status == ModuleRebuildStatus.Running))
- .OrderBy(x => x.Id)
- .FirstAsync(ct);
- public async Task<bool> TrySupersedeQueuedAsync(long jobId, string reason, CancellationToken ct = default)
- {
- var now = DateTime.Now;
- var n = await _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob
- {
- Status = ModuleRebuildStatus.Cancelled,
- CurrentStage = ModuleRebuildStages.Cancelled,
- FinishedAt = now,
- ErrorMessage = reason,
- UpdateTime = now
- })
- .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Queued)
- .ExecuteCommandAsync(ct);
- return n > 0;
- }
- public async Task<int> CancelUnstartedByTriggerAsync(string triggerType, TimeSpan olderThan, string reason, CancellationToken ct = default)
- {
- if (string.IsNullOrWhiteSpace(triggerType)) return 0;
- var trigger = triggerType.Trim().ToUpperInvariant();
- var now = DateTime.Now;
- var cutoff = now - olderThan;
- return await _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob
- {
- Status = ModuleRebuildStatus.Cancelled,
- CurrentStage = ModuleRebuildStages.Cancelled,
- FinishedAt = now,
- ErrorMessage = reason,
- UpdateTime = now
- })
- // requested_by IS NULL 这一支不可省:带 requested_by 的行按
- // ModuleRebuildTriggerType.IsManualTrigger 算人工触发,清了就是误伤页面按钮。
- .Where(x => x.Status == ModuleRebuildStatus.Queued
- && x.TriggerType == trigger
- && x.RequestedBy == null)
- .WhereIF(olderThan > TimeSpan.Zero, x => x.SubmittedAt < cutoff)
- .ExecuteCommandAsync(ct);
- }
- public async Task<int> FailStaleQueuedAsync(TimeSpan staleAfter, string reason, CancellationToken ct = default)
- {
- var now = DateTime.Now;
- var cutoff = now - staleAfter;
- return await _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob
- {
- Status = ModuleRebuildStatus.Failed,
- CurrentStage = ModuleRebuildStages.Failed,
- FinishedAt = now,
- ErrorMessage = reason,
- UpdateTime = now
- })
- .Where(x => x.Status == ModuleRebuildStatus.Queued && x.SubmittedAt < cutoff)
- .ExecuteCommandAsync(ct);
- }
- public Task<AdoModuleDashboardRebuildJob> FindLastSuccessAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) =>
- _db.Queryable<AdoModuleDashboardRebuildJob>()
- .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId
- && x.Status == ModuleRebuildStatus.Success)
- .OrderBy(x => x.FinishedAt, OrderByType.Desc)
- .FirstAsync(ct);
- public Task<AdoModuleDashboardRebuildJob> GetByIdAsync(string moduleCode, long id, long tenantId, long factoryId, CancellationToken ct = default) =>
- _db.Queryable<AdoModuleDashboardRebuildJob>()
- .FirstAsync(x => x.Id == id && x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId, ct);
- public Task<AdoModuleDashboardRebuildJob> GetLatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default) =>
- _db.Queryable<AdoModuleDashboardRebuildJob>()
- .Where(x => x.ModuleCode == moduleCode && x.TenantId == tenantId && x.FactoryId == factoryId)
- .OrderBy(x => x.Id, OrderByType.Desc)
- .FirstAsync(ct);
- public async Task<bool> RequestCancelAsync(long jobId, CancellationToken ct = default)
- {
- var affected = await _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob { CancelRequestedFlag = true, UpdateTime = DateTime.Now })
- .Where(x => x.Id == jobId
- && (x.Status == ModuleRebuildStatus.Queued || x.Status == ModuleRebuildStatus.Running))
- .ExecuteCommandAsync(ct);
- return affected > 0;
- }
- public async Task<long> FindOpenRunLogIdAsync(string moduleCode, long tenantId, DateTime startedAfter)
- {
- var jobCode = $"{moduleCode}_MDP_SYNC_TRANSFORM";
- var id = await _db.Ado.GetLongAsync(
- """
- SELECT id FROM mdp_transform_run_log
- WHERE job_code=@JobCode AND tenant_id=@TenantId AND status='RUNNING' AND start_time>=@StartedAfter
- ORDER BY id DESC LIMIT 1
- """,
- new List<SugarParameter>
- {
- new("@JobCode", jobCode),
- new("@TenantId", tenantId),
- new("@StartedAfter", startedAfter.AddMinutes(-1))
- });
- return id;
- }
- public async Task<bool> IsCancelRequestedAsync(long jobId, CancellationToken ct = default)
- {
- var flag = await _db.Queryable<AdoModuleDashboardRebuildJob>()
- .Where(x => x.Id == jobId)
- .Select(x => x.CancelRequestedFlag)
- .FirstAsync(ct);
- return flag;
- }
- /// <summary>跨实例抢占互斥锁键。只保护「数 RUNNING + 置 RUNNING」这一小段,不覆盖重算本体。</summary>
- public const string ClaimLockKey = "aidop:mdp:rebuild:claim";
- public async Task<AdoModuleDashboardRebuildJob> ClaimNextQueuedAsync(
- IReadOnlyCollection<string> enabledModules, int globalMaxParallelScopes, bool runner, CancellationToken ct = default)
- {
- if (enabledModules == null || enabledModules.Count == 0)
- return null;
- // 计数与抢占必须在同一把跨实例锁内:否则 N 个实例会在同一瞬间都读到
- // 「RUNNING 还没到上限」,各自抢走一个 job,上限形同虚设。
- // 非阻塞取锁——抢不到说明别的实例正在抢占,直接让出,5 秒后下一拍再来。
- await using var claimGate = await InventoryInboundLockGuard.TryAcquireAsync(
- _db, _logger, ClaimLockKey, ct);
- if (!claimGate.Acquired)
- return null;
- var running = await _db.Queryable<AdoModuleDashboardRebuildJob>()
- .Where(x => x.Status == ModuleRebuildStatus.Running)
- .CountAsync(ct);
- if (running >= globalMaxParallelScopes)
- return null;
- var enabled = enabledModules.ToList();
- var row = await _db.Queryable<AdoModuleDashboardRebuildJob>()
- .Where(x => x.Status == ModuleRebuildStatus.Queued && enabled.Contains(x.ModuleCode))
- // 非执行机只领页面手工任务。用**白名单**而非「排除 AUTO/BOOTSTRAP」:
- // 后续新增的自动触发类型(如夜间兜底的 AUTO_NIGHTLY)天然落在白名单外,
- // 不会因为有人忘了往黑名单里补一项就被开发机领去跑全量。
- .WhereIF(!runner, x => x.TriggerType == ModuleRebuildTriggerType.Manual || x.RequestedBy != null)
- .Where(x => SqlFunc.Subqueryable<AdoModuleDashboardRebuildJob>()
- .Where(y => y.ModuleCode == x.ModuleCode && y.TenantId == x.TenantId && y.FactoryId == x.FactoryId
- && y.Status == ModuleRebuildStatus.Running)
- .NotAny())
- .OrderBy(x => x.Id)
- .FirstAsync(ct);
- if (row == null)
- return null;
- var now = DateTime.Now;
- var n = await _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob
- {
- Status = ModuleRebuildStatus.Running,
- CurrentStage = ModuleRebuildStages.AcquiringLock,
- StageIndex = 0,
- ProgressPercent = 2,
- ProgressMessage = $"等待现有 {row.ModuleCode} 全量任务完成",
- LastProgressAt = now,
- StartedAt = now,
- HeartbeatAt = now,
- UpdateTime = now
- })
- .Where(x => x.Id == row.Id && x.Status == ModuleRebuildStatus.Queued)
- .ExecuteCommandAsync(ct);
- if (n <= 0)
- return null;
- row.Status = ModuleRebuildStatus.Running;
- row.CurrentStage = ModuleRebuildStages.AcquiringLock;
- row.ProgressPercent = 2;
- row.StartedAt = now;
- row.HeartbeatAt = now;
- row.UpdateTime = now;
- return row;
- }
- public Task UpdateAsync(AdoModuleDashboardRebuildJob row, CancellationToken ct = default) =>
- _db.Updateable(row).ExecuteCommandAsync(ct);
- public Task UpdateProgressAsync(long jobId, string currentStage, int stageIndex, int progressPercent, string message, DateTime now, CancellationToken ct = default) =>
- _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob
- {
- CurrentStage = currentStage,
- StageIndex = stageIndex,
- ProgressPercent = progressPercent,
- ProgressMessage = message,
- LastProgressAt = now,
- HeartbeatAt = now,
- UpdateTime = now
- })
- .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Running)
- .ExecuteCommandAsync(ct);
- public Task UpdateStageResultAsync(long jobId, string completedStage, int rows, int progressPercent, string nextMessage, string detailJson, DateTime now, CancellationToken ct = default)
- {
- var column = completedStage switch
- {
- ModuleRebuildStages.Staging or ModuleRebuildStages.T8Inbound => "stage_rows",
- ModuleRebuildStages.Standard or ModuleRebuildStages.KpiPreparing => "standard_rows",
- ModuleRebuildStages.Dwd => "dwd_rows",
- ModuleRebuildStages.Kpi or ModuleRebuildStages.KpiCalculating => "kpi_rows",
- ModuleRebuildStages.Atomic => "atomic_rows",
- _ => null
- };
- if (column == null)
- return UpdateProgressAsync(jobId, completedStage, 0, progressPercent, nextMessage, now, ct);
- return _db.Ado.ExecuteCommandAsync(
- $"""
- UPDATE ado_module_dashboard_rebuild_job
- SET `{column}`=@Rows,
- progress_percent=@Pct,
- progress_message=@Msg,
- detail_json=IFNULL(@Detail, detail_json),
- last_progress_at=@Now,
- heartbeat_at=@Now,
- update_time=@Now
- WHERE id=@Id AND status='RUNNING'
- """,
- new SugarParameter("@Rows", rows),
- new SugarParameter("@Pct", progressPercent),
- new SugarParameter("@Msg", nextMessage),
- new SugarParameter("@Detail", (object?)detailJson ?? DBNull.Value),
- new SugarParameter("@Now", now),
- new SugarParameter("@Id", jobId));
- }
- public Task TouchHeartbeatAsync(long jobId, DateTime now, CancellationToken ct = default) =>
- _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob { HeartbeatAt = now, UpdateTime = now })
- .Where(x => x.Id == jobId && x.Status == ModuleRebuildStatus.Running)
- .ExecuteCommandAsync(ct);
- public Task FailStaleRunningAsync(TimeSpan staleAfter, CancellationToken ct = default)
- {
- var cutoff = DateTime.Now - staleAfter;
- var now = DateTime.Now;
- return _db.Updateable<AdoModuleDashboardRebuildJob>()
- .SetColumns(x => new AdoModuleDashboardRebuildJob
- {
- Status = ModuleRebuildStatus.Failed,
- CurrentStage = ModuleRebuildStages.Failed,
- FinishedAt = now,
- LastProgressAt = now,
- UpdateTime = now,
- ErrorMessage = "服务中断,任务未正常结束"
- })
- .Where(x => x.Status == ModuleRebuildStatus.Running && (x.HeartbeatAt == null || x.HeartbeatAt < cutoff))
- .ExecuteCommandAsync(ct);
- }
- }
|