S0DimRefreshService.cs 10 KB

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