S0DimRefreshService.cs 11 KB

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