S0DimRefreshService.cs 9.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204
  1. using System.Text.Json;
  2. using Microsoft.Extensions.Logging;
  3. using SqlSugar;
  4. namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
  5. /// <summary>
  6. /// S0 维度刷新编排:批次号、逐维度物化、run log 汇总。
  7. ///
  8. /// 一次刷新 = 一个 <c>batch_id</c>(同时是 staging / dim 的 <c>sync_batch_id</c> 与
  9. /// <c>mdp_transform_run_log.batch_id</c>;后者是 UNIQUE,故只写一条汇总行,
  10. /// 逐维度明细进 <c>summary_json</c>)。
  11. /// </summary>
  12. public sealed class S0DimRefreshService : ITransient
  13. {
  14. /// <summary>run log 的作业编码。</summary>
  15. public const string JobCode = "S0_DIM_REFRESH";
  16. private static readonly JsonSerializerOptions JsonOpts = new()
  17. {
  18. Encoder = System.Text.Encodings.Web.JavaScriptEncoder.UnsafeRelaxedJsonEscaping
  19. };
  20. private readonly ISqlSugarClient _db;
  21. private readonly S0DimMaterializer _materializer;
  22. private readonly S0DimReconciler _reconciler;
  23. private readonly ILogger<S0DimRefreshService> _logger;
  24. /// <summary>构造。</summary>
  25. public S0DimRefreshService(
  26. ISqlSugarClient db,
  27. S0DimMaterializer materializer,
  28. S0DimReconciler reconciler,
  29. ILogger<S0DimRefreshService> logger)
  30. {
  31. _db = db;
  32. _materializer = materializer;
  33. _reconciler = reconciler;
  34. _logger = logger;
  35. }
  36. /// <summary>
  37. /// 刷新指定租户的维度。<paramref name="key"/> 为空表示全部(按 <see cref="S0DimCatalog.All"/> 的依赖顺序)。
  38. /// </summary>
  39. public async Task<S0DimRefreshResult> RefreshAsync(
  40. long tenantId, string? key = null, string triggerType = "MANUAL", CancellationToken ct = default)
  41. {
  42. if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}");
  43. S0DimCatalog.ValidateAll();
  44. var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)];
  45. var startedAt = DateTime.Now;
  46. var batchId = $"S0DIM_{tenantId}_{startedAt:yyyyMMddHHmmssfff}";
  47. var result = new S0DimRefreshResult { BatchId = batchId, TenantId = tenantId };
  48. result.RunLogId = await InsertRunLogAsync(tenantId, batchId, startedAt, triggerType);
  49. try
  50. {
  51. foreach (var def in defs)
  52. {
  53. ct.ThrowIfCancellationRequested();
  54. result.Items.Add(await _materializer.RunAsync(def, tenantId, batchId, ct));
  55. }
  56. }
  57. catch (OperationCanceledException)
  58. {
  59. // run log 没有 stale 回收机制,而 LastRunAsync 会把 RUNNING 行当成「最近一次结果」返回。
  60. // 若取消时直接抛出,这条记录会永久停在 RUNNING,读者无法区分「正在跑」与「早就断了」。
  61. // 故先收口再重抛。CompleteRunLogAsync 不接收 ct —— 否则它自己也会被同一个 token 取消。
  62. result.Status = "CANCELED";
  63. result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds;
  64. await CompleteRunLogAsync(result, startedAt);
  65. _logger.LogWarning("[S0Dim] refresh canceled tenant={Tenant} batch={Batch} 已完成 {Done}/{Total} 个维度",
  66. tenantId, batchId, result.Items.Count, defs.Count);
  67. throw;
  68. }
  69. // 归类顺序:FAILED > 有告警 > 全部跳过 > SUCCESS。
  70. // SKIPPED_NO_CHANGE 排在 SUCCESS 之前、告警之后:它是「非失败」,但**不得伪装成 SUCCESS** ——
  71. // 否则「自动刷新是否真的在跑」将不可观测(全跳过与全刷新在 run log 里长得一模一样)。
  72. // 只有**所有**维度都被写前闸门跳过才整轮记 SKIPPED_NO_CHANGE;只要有一个真刷过就是 SUCCESS。
  73. // 被跳过的维度 Reconcile 为 null,故不会贡献告警,与上一档不冲突。
  74. result.Status = result.Items.Any(i => i.Status == S0DimMaterializeResult.StatusFailed)
  75. ? "FAILED"
  76. : result.Items.Any(i => i.Reconcile?.Warnings.Count > 0)
  77. ? "SUCCESS_WITH_WARNING"
  78. : result.Items.Count > 0
  79. && result.Items.All(i => i.Status == S0DimMaterializeResult.StatusSkippedNoChange)
  80. ? S0DimMaterializeResult.StatusSkippedNoChange
  81. : "SUCCESS";
  82. result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds;
  83. await CompleteRunLogAsync(result, startedAt);
  84. _logger.LogInformation("[S0Dim] refresh done tenant={Tenant} batch={Batch} status={Status}",
  85. tenantId, batchId, result.Status);
  86. return result;
  87. }
  88. /// <summary>只读对账(不改数据)。<paramref name="key"/> 为空表示全部。</summary>
  89. public async Task<List<S0DimReconcileResult>> ReconcileAsync(
  90. long tenantId, string? key = null, CancellationToken ct = default)
  91. {
  92. if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}");
  93. S0DimCatalog.ValidateAll();
  94. var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)];
  95. var list = new List<S0DimReconcileResult>();
  96. foreach (var def in defs)
  97. {
  98. ct.ThrowIfCancellationRequested();
  99. // batchId 传 null:按「当前状态」口径对账,不限定某一批
  100. list.Add(await _reconciler.ReconcileAsync(def, tenantId, null, ct));
  101. }
  102. return list;
  103. }
  104. /// <summary>最近一次运行记录(按租户)。</summary>
  105. public async Task<object?> LastRunAsync(long tenantId, CancellationToken ct = default)
  106. {
  107. var rows = await _db.Ado.SqlQueryAsync<S0DimRunLogRow>(
  108. """
  109. SELECT batch_id AS BatchId, status AS Status, start_time AS StartTime, end_time AS EndTime,
  110. duration_ms AS DurationMs, stage_rows AS StageRows, standard_rows AS StandardRows,
  111. error_message AS ErrorMessage, summary_json AS SummaryJson
  112. FROM mdp_transform_run_log
  113. WHERE job_code = @JobCode AND tenant_id = @TenantId
  114. ORDER BY id DESC LIMIT 1
  115. """,
  116. new List<SugarParameter> { new("@JobCode", JobCode), new("@TenantId", tenantId) });
  117. return rows.FirstOrDefault();
  118. }
  119. private static S0DimDefinition ResolveOrThrow(string key) =>
  120. S0DimCatalog.Find(key)
  121. ?? throw new InvalidOperationException(
  122. $"未知维度 key:{key}(可用:{string.Join(" / ", S0DimCatalog.All.Select(d => d.Key))})");
  123. private async Task<long> InsertRunLogAsync(long tenantId, string batchId, DateTime startedAt, string triggerType)
  124. {
  125. var trigger = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant();
  126. await _db.Ado.ExecuteCommandAsync(
  127. """
  128. INSERT INTO mdp_transform_run_log
  129. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  130. VALUES
  131. (@TenantId, @JobCode, 'S0 主数据维度刷新(source → mdp_stg_s0_* → dim_*)', @Trigger, @BatchId, 'RUNNING', @StartTime)
  132. """,
  133. new List<SugarParameter>
  134. {
  135. new("@TenantId", tenantId),
  136. new("@JobCode", JobCode),
  137. new("@Trigger", trigger),
  138. new("@BatchId", batchId),
  139. new("@StartTime", startedAt)
  140. });
  141. // 带上 tenant_id:batch_id 虽是 UNIQUE 且内嵌 tenantId,但那是命名约定 + 表级约束两个**外部**事实;
  142. // 语句自身带租户谓词才符合 D5,也避免将来 batch_id 命名变更时静默取到别人的行。
  143. return await _db.Ado.GetLongAsync(
  144. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId AND tenant_id=@TenantId ORDER BY id DESC LIMIT 1",
  145. new List<SugarParameter> { new("@BatchId", batchId), new("@TenantId", tenantId) });
  146. }
  147. private async Task CompleteRunLogAsync(S0DimRefreshResult result, DateTime startedAt)
  148. {
  149. var errors = result.Items.Where(i => i.Status == S0DimMaterializeResult.StatusFailed)
  150. .Select(i => $"{i.Key}: {i.Error}").ToList();
  151. await _db.Ado.ExecuteCommandAsync(
  152. """
  153. UPDATE mdp_transform_run_log
  154. SET status=@Status, end_time=@EndTime, duration_ms=@DurationMs,
  155. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=0,
  156. error_message=@Error, summary_json=@Summary, update_time=CURRENT_TIMESTAMP
  157. WHERE id=@Id AND tenant_id=@TenantId
  158. """,
  159. new List<SugarParameter>
  160. {
  161. new("@TenantId", result.TenantId),
  162. new("@Status", result.Status),
  163. new("@EndTime", DateTime.Now),
  164. new("@DurationMs", result.DurationMs),
  165. new("@StageRows", result.Items.Sum(i => i.StagingWritten)),
  166. new("@StandardRows", result.Items.Sum(i => i.DimRows)),
  167. new("@Error", errors.Count == 0 ? null : string.Join(" | ", errors)),
  168. new("@Summary", JsonSerializer.Serialize(result.Items, JsonOpts)),
  169. new("@Id", result.RunLogId)
  170. });
  171. }
  172. private sealed class S0DimRunLogRow
  173. {
  174. public string BatchId { get; set; } = "";
  175. public string Status { get; set; } = "";
  176. public DateTime StartTime { get; set; }
  177. public DateTime? EndTime { get; set; }
  178. public int? DurationMs { get; set; }
  179. public int StageRows { get; set; }
  180. public int StandardRows { get; set; }
  181. public string? ErrorMessage { get; set; }
  182. public string? SummaryJson { get; set; }
  183. }
  184. }