| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196 |
- using System.Text.Json;
- using Microsoft.Extensions.Logging;
- using SqlSugar;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
- /// <summary>
- /// S0 维度刷新编排:批次号、逐维度物化、run log 汇总。
- ///
- /// 一次刷新 = 一个 <c>batch_id</c>(同时是 staging / dim 的 <c>sync_batch_id</c> 与
- /// <c>mdp_transform_run_log.batch_id</c>;后者是 UNIQUE,故只写一条汇总行,
- /// 逐维度明细进 <c>summary_json</c>)。
- /// </summary>
- public sealed class S0DimRefreshService : ITransient
- {
- /// <summary>run log 的作业编码。</summary>
- 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<S0DimRefreshService> _logger;
- /// <summary>构造。</summary>
- public S0DimRefreshService(
- ISqlSugarClient db,
- S0DimMaterializer materializer,
- S0DimReconciler reconciler,
- ILogger<S0DimRefreshService> logger)
- {
- _db = db;
- _materializer = materializer;
- _reconciler = reconciler;
- _logger = logger;
- }
- /// <summary>
- /// 刷新指定租户的维度。<paramref name="key"/> 为空表示全部(按 <see cref="S0DimCatalog.All"/> 的依赖顺序)。
- /// </summary>
- public async Task<S0DimRefreshResult> 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;
- }
- /// <summary>只读对账(不改数据)。<paramref name="key"/> 为空表示全部。</summary>
- public async Task<List<S0DimReconcileResult>> 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<S0DimReconcileResult>();
- foreach (var def in defs)
- {
- ct.ThrowIfCancellationRequested();
- // batchId 传 null:按「当前状态」口径对账,不限定某一批
- list.Add(await _reconciler.ReconcileAsync(def, tenantId, null, ct));
- }
- return list;
- }
- /// <summary>最近一次运行记录(按租户)。</summary>
- public async Task<object?> LastRunAsync(long tenantId, CancellationToken ct = default)
- {
- var rows = await _db.Ado.SqlQueryAsync<S0DimRunLogRow>(
- """
- 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<SugarParameter> { 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<long> 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<SugarParameter>
- {
- 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<SugarParameter> { 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<SugarParameter>
- {
- 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; }
- }
- }
|