using System.Text.Json; using Microsoft.Extensions.Logging; using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim; /// /// S0 维度刷新编排:批次号、逐维度物化、run log 汇总。 /// /// 一次刷新 = 一个 batch_id(同时是 staging / dim 的 sync_batch_id 与 /// mdp_transform_run_log.batch_id;后者是 UNIQUE,故只写一条汇总行, /// 逐维度明细进 summary_json)。 /// public sealed class S0DimRefreshService : ITransient { /// run log 的作业编码。 public const string JobCode = "S0_DIM_REFRESH"; private static readonly JsonSerializerOptions JsonOpts = new() { Encoder = System.Text.Encodings.Web.JavaScriptEncoder.UnsafeRelaxedJsonEscaping }; private readonly ISqlSugarClient _db; private readonly S0DimMaterializer _materializer; private readonly S0DimReconciler _reconciler; private readonly ILogger _logger; /// 构造。 public S0DimRefreshService( ISqlSugarClient db, S0DimMaterializer materializer, S0DimReconciler reconciler, ILogger logger) { _db = db; _materializer = materializer; _reconciler = reconciler; _logger = logger; } /// /// 刷新指定租户的维度。 为空表示全部(按 的依赖顺序)。 /// public async Task RefreshAsync( long tenantId, string? key = null, string triggerType = "MANUAL", CancellationToken ct = default) { if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}"); S0DimCatalog.ValidateAll(); var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)]; var startedAt = DateTime.Now; var batchId = $"S0DIM_{tenantId}_{startedAt:yyyyMMddHHmmssfff}"; var result = new S0DimRefreshResult { BatchId = batchId, TenantId = tenantId }; result.RunLogId = await InsertRunLogAsync(tenantId, batchId, startedAt, triggerType); try { foreach (var def in defs) { ct.ThrowIfCancellationRequested(); result.Items.Add(await _materializer.RunAsync(def, tenantId, batchId, ct)); } } catch (OperationCanceledException) { // run log 没有 stale 回收机制,而 LastRunAsync 会把 RUNNING 行当成「最近一次结果」返回。 // 若取消时直接抛出,这条记录会永久停在 RUNNING,读者无法区分「正在跑」与「早就断了」。 // 故先收口再重抛。CompleteRunLogAsync 不接收 ct —— 否则它自己也会被同一个 token 取消。 result.Status = "CANCELED"; result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds; await CompleteRunLogAsync(result, startedAt); _logger.LogWarning("[S0Dim] refresh canceled tenant={Tenant} batch={Batch} 已完成 {Done}/{Total} 个维度", tenantId, batchId, result.Items.Count, defs.Count); throw; } result.Status = result.Items.Any(i => i.Status == "FAILED") ? "FAILED" : result.Items.Any(i => i.Reconcile?.Warnings.Count > 0) ? "SUCCESS_WITH_WARNING" : "SUCCESS"; result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds; await CompleteRunLogAsync(result, startedAt); _logger.LogInformation("[S0Dim] refresh done tenant={Tenant} batch={Batch} status={Status}", tenantId, batchId, result.Status); return result; } /// 只读对账(不改数据)。 为空表示全部。 public async Task> ReconcileAsync( long tenantId, string? key = null, CancellationToken ct = default) { if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}"); S0DimCatalog.ValidateAll(); var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)]; var list = new List(); foreach (var def in defs) { ct.ThrowIfCancellationRequested(); // batchId 传 null:按「当前状态」口径对账,不限定某一批 list.Add(await _reconciler.ReconcileAsync(def, tenantId, null, ct)); } return list; } /// 最近一次运行记录(按租户)。 public async Task LastRunAsync(long tenantId, CancellationToken ct = default) { var rows = await _db.Ado.SqlQueryAsync( """ SELECT batch_id AS BatchId, status AS Status, start_time AS StartTime, end_time AS EndTime, duration_ms AS DurationMs, stage_rows AS StageRows, standard_rows AS StandardRows, error_message AS ErrorMessage, summary_json AS SummaryJson FROM mdp_transform_run_log WHERE job_code = @JobCode AND tenant_id = @TenantId ORDER BY id DESC LIMIT 1 """, new List { new("@JobCode", JobCode), new("@TenantId", tenantId) }); return rows.FirstOrDefault(); } private static S0DimDefinition ResolveOrThrow(string key) => S0DimCatalog.Find(key) ?? throw new InvalidOperationException( $"未知维度 key:{key}(可用:{string.Join(" / ", S0DimCatalog.All.Select(d => d.Key))})"); private async Task InsertRunLogAsync(long tenantId, string batchId, DateTime startedAt, string triggerType) { var trigger = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant(); await _db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_transform_run_log (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time) VALUES (@TenantId, @JobCode, 'S0 主数据维度刷新(source → mdp_stg_s0_* → dim_*)', @Trigger, @BatchId, 'RUNNING', @StartTime) """, new List { new("@TenantId", tenantId), new("@JobCode", JobCode), new("@Trigger", trigger), new("@BatchId", batchId), new("@StartTime", startedAt) }); // 带上 tenant_id:batch_id 虽是 UNIQUE 且内嵌 tenantId,但那是命名约定 + 表级约束两个**外部**事实; // 语句自身带租户谓词才符合 D5,也避免将来 batch_id 命名变更时静默取到别人的行。 return await _db.Ado.GetLongAsync( "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId AND tenant_id=@TenantId ORDER BY id DESC LIMIT 1", new List { new("@BatchId", batchId), new("@TenantId", tenantId) }); } private async Task CompleteRunLogAsync(S0DimRefreshResult result, DateTime startedAt) { var errors = result.Items.Where(i => i.Status == "FAILED") .Select(i => $"{i.Key}: {i.Error}").ToList(); await _db.Ado.ExecuteCommandAsync( """ UPDATE mdp_transform_run_log SET status=@Status, end_time=@EndTime, duration_ms=@DurationMs, stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=0, error_message=@Error, summary_json=@Summary, update_time=CURRENT_TIMESTAMP WHERE id=@Id AND tenant_id=@TenantId """, new List { new("@TenantId", result.TenantId), new("@Status", result.Status), new("@EndTime", DateTime.Now), new("@DurationMs", result.DurationMs), new("@StageRows", result.Items.Sum(i => i.StagingWritten)), new("@StandardRows", result.Items.Sum(i => i.DimRows)), new("@Error", errors.Count == 0 ? null : string.Join(" | ", errors)), new("@Summary", JsonSerializer.Serialize(result.Items, JsonOpts)), new("@Id", result.RunLogId) }); } private sealed class S0DimRunLogRow { public string BatchId { get; set; } = ""; public string Status { get; set; } = ""; public DateTime StartTime { get; set; } public DateTime? EndTime { get; set; } public int? DurationMs { get; set; } public int StageRows { get; set; } public int StandardRows { get; set; } public string? ErrorMessage { get; set; } public string? SummaryJson { get; set; } } }