| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399 |
- using Admin.NET.Plugin.AiDOP.Entity.SmartOps;
- using Microsoft.Extensions.Logging;
- using Microsoft.Extensions.Logging.Abstractions;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.MdpRebuild;
- public sealed class ModuleRebuildService : ITransient
- {
- private readonly IModuleRebuildJobStore _jobs;
- private readonly ModuleRebuildQueue _queue;
- private readonly IModuleRebuildCapability _capability;
- private readonly ILogger _logger;
- public ModuleRebuildService(IModuleRebuildJobStore jobs, ModuleRebuildQueue queue, IModuleRebuildCapability capability)
- : this(jobs, queue, capability, NullLoggerFactory.Instance)
- {
- }
- private readonly MdpNeutralSourceCleanup? _cleanup;
- private readonly TransformRunLogFinalizer? _runLogFinalizer;
- public ModuleRebuildService(
- IModuleRebuildJobStore jobs,
- ModuleRebuildQueue queue,
- IModuleRebuildCapability capability,
- ILoggerFactory loggerFactory)
- : this(jobs, queue, capability, loggerFactory, null)
- {
- }
- public ModuleRebuildService(
- IModuleRebuildJobStore jobs,
- ModuleRebuildQueue queue,
- IModuleRebuildCapability capability,
- ILoggerFactory loggerFactory,
- MdpNeutralSourceCleanup? cleanup,
- TransformRunLogFinalizer? runLogFinalizer = null)
- {
- _jobs = jobs;
- _queue = queue;
- _capability = capability;
- _cleanup = cleanup;
- _runLogFinalizer = runLogFinalizer;
- _logger = loggerFactory.CreateLogger(nameof(ModuleRebuildService));
- }
- public async Task<(int StatusCode, ModuleRebuildJobAccepted Body)> EnqueueAsync(
- string moduleCode,
- long tenantId,
- long factoryId,
- long? requestedBy,
- string triggerType = "MANUAL",
- CancellationToken ct = default)
- {
- var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
- if (!_capability.IsEnabled(scope.ModuleCode))
- {
- return (404, new ModuleRebuildJobAccepted
- {
- Ok = false,
- ModuleCode = scope.ModuleCode,
- Status = "DISABLED",
- Message = $"{scope.ModuleCode} 数据重算尚未启用"
- });
- }
- var normalizedTrigger = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant();
- var manualRequest = ModuleRebuildTriggerType.IsManualTrigger(normalizedTrigger, requestedBy);
- var active = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
- if (active != null)
- {
- // 人工重算可以取代同 scope 里「还没开跑的自动任务」:二者都是同 scope 全量,取代不丢工作量。
- // 不让位则会死锁——ClaimNextQueuedAsync 只让执行机领自动任务,非执行机上这条 QUEUED 永远
- // 不会消失,页面「数据重算」按钮就被它永久挡住(执行机指派还会随进程重启失效)。
- var supersedable = manualRequest
- && active.Status == ModuleRebuildStatus.Queued
- && !ModuleRebuildTriggerType.IsManualTrigger(active.TriggerType, active.RequestedBy);
- if (!supersedable
- || !await _jobs.TrySupersedeQueuedAsync(active.Id, $"已被人工重算取代(原 {active.TriggerType})", ct))
- {
- return (409, Conflict(scope.ModuleCode, active));
- }
- _logger.LogInformation(
- "[ModuleRebuild] {Module} tenant={Tenant} factory={Factory} 人工重算取代待执行任务 job={JobId} trigger={Trigger}",
- scope.ModuleCode, scope.TenantId, scope.FactoryId, active.Id, active.TriggerType);
- }
- // AUTO / BOOTSTRAP 才冷却。MANUAL 是页面「数据重算」按钮,AUTO_NIGHTLY 是全量兜底,两者都必须放行。
- var cooldownHours = _capability.AutoMinIntervalHours;
- if (cooldownHours > 0
- && !manualRequest
- && normalizedTrigger != ModuleRebuildTriggerType.AutoNightly)
- {
- var lastSuccess = await _jobs.FindLastSuccessAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
- if (lastSuccess?.FinishedAt is { } finishedAt
- && DateTime.Now - finishedAt < TimeSpan.FromHours(cooldownHours))
- {
- return (409, new ModuleRebuildJobAccepted
- {
- Ok = false,
- ModuleCode = scope.ModuleCode,
- JobId = lastSuccess.Id,
- Status = "COOLDOWN",
- Message = $"{scope.ModuleCode} 距上次成功重算不足 {cooldownHours} 小时,本次跳过"
- });
- }
- }
- var now = DateTime.Now;
- var created = await _jobs.InsertQueuedAsync(new AdoModuleDashboardRebuildJob
- {
- ModuleCode = scope.ModuleCode,
- TenantId = scope.TenantId,
- FactoryId = scope.FactoryId,
- Status = ModuleRebuildStatus.Queued,
- CurrentStage = ModuleRebuildStages.Queued,
- StageIndex = 0,
- StageTotal = ModuleRebuildStages.StageTotal(scope.ModuleCode),
- ProgressPercent = 0,
- ProgressMessage = "已入队",
- LastProgressAt = now,
- TriggerType = normalizedTrigger,
- RequestedBy = requestedBy,
- SubmittedAt = now,
- CreateTime = now,
- UpdateTime = now
- }, ct);
- var oldest = await _jobs.FindActiveAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
- if (oldest != null && oldest.Id != created.Id)
- {
- created.Status = ModuleRebuildStatus.Failed;
- created.CurrentStage = ModuleRebuildStages.Failed;
- created.FinishedAt = DateTime.Now;
- created.ErrorMessage = $"{scope.ModuleCode} 数据重算正在执行,请勿重复提交";
- created.UpdateTime = DateTime.Now;
- await _jobs.UpdateAsync(created, ct);
- return (409, Conflict(scope.ModuleCode, oldest));
- }
- _queue.Pulse();
- return (202, new ModuleRebuildJobAccepted
- {
- Ok = true,
- ModuleCode = scope.ModuleCode,
- JobId = created.Id,
- Status = ModuleRebuildStatus.Queued,
- Message = $"{scope.ModuleCode} 数据重算已排队"
- });
- }
- public async Task<ModuleRebuildJobDto> GetAsync(string moduleCode, long jobId, long tenantId, long factoryId, CancellationToken ct = default)
- {
- var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
- var row = await _jobs.GetByIdAsync(scope.ModuleCode, jobId, scope.TenantId, scope.FactoryId, ct);
- if (row == null)
- throw Oops.Oh("任务不存在");
- return ToDto(row);
- }
- public async Task<ModuleRebuildJobDto> LatestAsync(string moduleCode, long tenantId, long factoryId, CancellationToken ct = default)
- {
- var scope = MdpRebuildScope.Create(moduleCode, tenantId, factoryId);
- var row = await _jobs.GetLatestAsync(scope.ModuleCode, scope.TenantId, scope.FactoryId, ct);
- return row == null ? null : ToDto(row);
- }
- public async Task FailStaleAsync(CancellationToken ct = default) =>
- await _jobs.FailStaleRunningAsync(ModuleRebuildLock.StaleAfter, ct);
- public async Task TouchHeartbeatAsync(long jobId, CancellationToken ct = default)
- {
- try
- {
- await _jobs.TouchHeartbeatAsync(jobId, DateTime.Now, ct);
- }
- catch (Exception ex)
- {
- _logger.LogWarning(ex, "[MdpRebuild] heartbeat failed jobId={JobId}", jobId);
- }
- }
- public async Task ApplyProgressAsync(long jobId, string moduleCode, ModuleProgressUpdate update, CancellationToken ct = default)
- {
- try
- {
- var now = DateTime.Now;
- if (update.CompletedStage != null && update.Rows.HasValue)
- {
- await _jobs.UpdateStageResultAsync(
- jobId, update.CompletedStage, update.Rows.Value, update.ProgressPercent, update.Message, update.DetailJson, now, ct);
- }
- await _jobs.UpdateProgressAsync(
- jobId, update.Stage, update.StageIndex, update.ProgressPercent, update.Message, now, ct);
- }
- catch (Exception ex)
- {
- _logger.LogWarning(ex, "[MdpRebuild] progress write failed jobId={JobId} module={Module} stage={Stage}", jobId, moduleCode, update.Stage);
- }
- }
- public async Task RunClaimedAsync(
- AdoModuleDashboardRebuildJob job,
- IModuleRebuildHandler handler,
- IModuleRebuildLock runLock,
- CancellationToken stoppingToken)
- {
- var failedStage = job.CurrentStage;
- var scope = MdpRebuildScope.Create(job.ModuleCode, job.TenantId, job.FactoryId);
- try
- {
- await ApplyProgressAsync(job.Id, scope.ModuleCode, new ModuleProgressUpdate(
- ModuleRebuildStages.AcquiringLock, 0, 2, $"等待现有 {scope.ModuleCode} 全量任务完成"), CancellationToken.None);
- var holderId = $"{job.TriggerType}:{scope.ScopeKey}:{Environment.MachineName}:{Guid.NewGuid():N}";
- var lease = await runLock.TryAcquireAsync(scope, holderId, job.Id, stoppingToken);
- if (lease == null)
- throw new ModuleRebuildAlreadyRunningException(scope.ModuleCode);
- await using (lease)
- {
- using var lockHeartbeatCts = CancellationTokenSource.CreateLinkedTokenSource(stoppingToken);
- var lockHeartbeat = KeepLeaseAliveAsync(lease, lockHeartbeatCts.Token);
- try
- {
- if (_cleanup != null)
- await _cleanup.PurgeModuleAsync(scope.TenantId, scope.ModuleCode, stoppingToken);
- var result = await handler.RunAsync(
- scope,
- job.TriggerType,
- job.Id,
- async update =>
- {
- failedStage = update.Stage;
- await ApplyProgressAsync(job.Id, scope.ModuleCode, update, CancellationToken.None);
- // 阶段边界:每个 handler 每跨一个阶段都会走到这里,是唯一的通用取消点。
- if (await _jobs.IsCancelRequestedAsync(job.Id, CancellationToken.None))
- throw new ModuleRebuildCancelledException();
- },
- stoppingToken);
- var now = DateTime.Now;
- job.Status = ModuleRebuildStatus.Success;
- job.CurrentStage = ModuleRebuildStages.Success;
- job.StageIndex = ModuleRebuildStages.StageTotal(scope.ModuleCode);
- job.ProgressPercent = 100;
- job.ProgressMessage = $"{scope.ModuleCode} 数据重算已完成";
- job.FailedStage = null;
- job.LastProgressAt = now;
- job.BatchId = result.BatchId;
- job.TransformRunLogId = result.RunLogId;
- job.StageRows = result.StageRows;
- job.StandardRows = result.StandardRows;
- job.DwdRows = result.DwdRows;
- job.KpiRows = result.KpiRows;
- job.AtomicRows = result.AtomicRows;
- job.DetailJson = result.DetailJson;
- job.FinishedAt = now;
- job.HeartbeatAt = now;
- job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null;
- job.UpdateTime = now;
- job.ErrorMessage = null;
- await _jobs.UpdateAsync(job, CancellationToken.None);
- }
- finally
- {
- lockHeartbeatCts.Cancel();
- try { await lockHeartbeat; } catch (OperationCanceledException) { }
- }
- }
- }
- catch (ModuleRebuildAlreadyRunningException)
- {
- var now = DateTime.Now;
- job.Status = ModuleRebuildStatus.Queued;
- job.CurrentStage = ModuleRebuildStages.AcquiringLock;
- job.StageIndex = 0;
- job.ProgressPercent = 2;
- job.ProgressMessage = $"等待现有 {scope.ModuleCode} 全量任务完成";
- job.StartedAt = null;
- job.HeartbeatAt = now;
- job.LastProgressAt = now;
- job.UpdateTime = now;
- job.ErrorMessage = null;
- await _jobs.UpdateAsync(job, CancellationToken.None);
- _queue.Pulse();
- }
- catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
- {
- await FailJobAsync(job, "服务停止,任务已取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage);
- }
- catch (ModuleRebuildCancelledException)
- {
- await FailJobAsync(job, "已由管理员取消", ModuleRebuildStatus.Cancelled, ModuleRebuildStages.Cancelled, failedStage);
- if (_runLogFinalizer != null)
- {
- var runLogId = await _jobs.FindOpenRunLogIdAsync(scope.ModuleCode, scope.TenantId, job.StartedAt ?? job.SubmittedAt);
- if (runLogId > 0)
- await _runLogFinalizer.FinalizeAsAdminCancelledAsync(runLogId, job.StartedAt ?? DateTime.Now);
- }
- }
- catch (Exception ex)
- {
- await FailJobAsync(
- job,
- Truncate(ex.Message, 2000),
- ModuleRebuildStatus.Failed,
- ModuleRebuildStages.Failed,
- failedStage);
- }
- }
- private async Task FailJobAsync(AdoModuleDashboardRebuildJob job, string message, string status, string currentStage, string failedStage)
- {
- var now = DateTime.Now;
- job.Status = status;
- job.CurrentStage = currentStage;
- job.FailedStage = status == ModuleRebuildStatus.Failed ? failedStage : job.FailedStage;
- job.ProgressMessage = status == ModuleRebuildStatus.Failed
- ? $"在【{ModuleRebuildStages.ToChinese(failedStage)}】失败"
- : message;
- job.FinishedAt = now;
- job.HeartbeatAt = now;
- job.LastProgressAt = now;
- job.DurationMs = job.StartedAt.HasValue ? ToDurationMs(now - job.StartedAt.Value) : null;
- job.ErrorMessage = message;
- job.UpdateTime = now;
- await _jobs.UpdateAsync(job, CancellationToken.None);
- }
- public static ModuleRebuildJobDto ToDto(AdoModuleDashboardRebuildJob row) => new()
- {
- Ok = true,
- ModuleCode = row.ModuleCode,
- JobId = row.Id,
- Status = row.Status,
- CurrentStage = row.CurrentStage,
- StageIndex = row.StageIndex,
- StageTotal = row.StageTotal <= 0 ? ModuleRebuildStages.StageTotal(row.ModuleCode) : row.StageTotal,
- ProgressPercent = row.ProgressPercent,
- ProgressMessage = row.ProgressMessage,
- LastProgressAt = row.LastProgressAt,
- HeartbeatAt = row.HeartbeatAt,
- FailedStage = row.FailedStage,
- SubmittedAt = row.SubmittedAt,
- StartedAt = row.StartedAt,
- FinishedAt = row.FinishedAt,
- DurationMs = row.DurationMs,
- BatchId = row.BatchId,
- StageRows = row.StageRows,
- StandardRows = row.StandardRows,
- DwdRows = row.DwdRows,
- KpiRows = row.KpiRows,
- AtomicRows = row.AtomicRows,
- DetailJson = row.DetailJson,
- ErrorMessage = row.ErrorMessage
- };
- private static ModuleRebuildJobAccepted Conflict(string moduleCode, AdoModuleDashboardRebuildJob active) => new()
- {
- Ok = false,
- ModuleCode = moduleCode,
- JobId = active.Id,
- Status = active.Status,
- Message = $"{moduleCode} 数据重算正在执行,请勿重复提交"
- };
- private static async Task KeepLeaseAliveAsync(IModuleRebuildLease lease, CancellationToken ct)
- {
- while (!ct.IsCancellationRequested)
- {
- try
- {
- await Task.Delay(TimeSpan.FromSeconds(30), ct);
- await lease.HeartbeatAsync(CancellationToken.None);
- }
- catch (OperationCanceledException)
- {
- return;
- }
- catch
- {
- // 单例 SqlSugar 连接被业务阶段占用时跳过本次锁心跳,避免打断 MDP。
- }
- }
- }
- private static int ToDurationMs(TimeSpan elapsed)
- {
- var ms = elapsed.TotalMilliseconds;
- if (double.IsNaN(ms) || ms <= 0) return 0;
- return ms >= int.MaxValue ? int.MaxValue : (int)ms;
- }
- private static string Truncate(string value, int max) =>
- string.IsNullOrEmpty(value) || value.Length <= max ? value : value[..max];
- }
|