S0DimRefreshService.cs 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196
  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. result.Status = result.Items.Any(i => i.Status == "FAILED")
  70. ? "FAILED"
  71. : result.Items.Any(i => i.Reconcile?.Warnings.Count > 0)
  72. ? "SUCCESS_WITH_WARNING"
  73. : "SUCCESS";
  74. result.DurationMs = (int)(DateTime.Now - startedAt).TotalMilliseconds;
  75. await CompleteRunLogAsync(result, startedAt);
  76. _logger.LogInformation("[S0Dim] refresh done tenant={Tenant} batch={Batch} status={Status}",
  77. tenantId, batchId, result.Status);
  78. return result;
  79. }
  80. /// <summary>只读对账(不改数据)。<paramref name="key"/> 为空表示全部。</summary>
  81. public async Task<List<S0DimReconcileResult>> ReconcileAsync(
  82. long tenantId, string? key = null, CancellationToken ct = default)
  83. {
  84. if (tenantId <= 0) throw new InvalidOperationException($"拒绝无效租户:{tenantId}");
  85. S0DimCatalog.ValidateAll();
  86. var defs = key is null or "" ? S0DimCatalog.All : [ResolveOrThrow(key)];
  87. var list = new List<S0DimReconcileResult>();
  88. foreach (var def in defs)
  89. {
  90. ct.ThrowIfCancellationRequested();
  91. // batchId 传 null:按「当前状态」口径对账,不限定某一批
  92. list.Add(await _reconciler.ReconcileAsync(def, tenantId, null, ct));
  93. }
  94. return list;
  95. }
  96. /// <summary>最近一次运行记录(按租户)。</summary>
  97. public async Task<object?> LastRunAsync(long tenantId, CancellationToken ct = default)
  98. {
  99. var rows = await _db.Ado.SqlQueryAsync<S0DimRunLogRow>(
  100. """
  101. SELECT batch_id AS BatchId, status AS Status, start_time AS StartTime, end_time AS EndTime,
  102. duration_ms AS DurationMs, stage_rows AS StageRows, standard_rows AS StandardRows,
  103. error_message AS ErrorMessage, summary_json AS SummaryJson
  104. FROM mdp_transform_run_log
  105. WHERE job_code = @JobCode AND tenant_id = @TenantId
  106. ORDER BY id DESC LIMIT 1
  107. """,
  108. new List<SugarParameter> { new("@JobCode", JobCode), new("@TenantId", tenantId) });
  109. return rows.FirstOrDefault();
  110. }
  111. private static S0DimDefinition ResolveOrThrow(string key) =>
  112. S0DimCatalog.Find(key)
  113. ?? throw new InvalidOperationException(
  114. $"未知维度 key:{key}(可用:{string.Join(" / ", S0DimCatalog.All.Select(d => d.Key))})");
  115. private async Task<long> InsertRunLogAsync(long tenantId, string batchId, DateTime startedAt, string triggerType)
  116. {
  117. var trigger = string.IsNullOrWhiteSpace(triggerType) ? "MANUAL" : triggerType.Trim().ToUpperInvariant();
  118. await _db.Ado.ExecuteCommandAsync(
  119. """
  120. INSERT INTO mdp_transform_run_log
  121. (tenant_id, job_code, job_name, trigger_type, batch_id, status, start_time)
  122. VALUES
  123. (@TenantId, @JobCode, 'S0 主数据维度刷新(source → mdp_stg_s0_* → dim_*)', @Trigger, @BatchId, 'RUNNING', @StartTime)
  124. """,
  125. new List<SugarParameter>
  126. {
  127. new("@TenantId", tenantId),
  128. new("@JobCode", JobCode),
  129. new("@Trigger", trigger),
  130. new("@BatchId", batchId),
  131. new("@StartTime", startedAt)
  132. });
  133. // 带上 tenant_id:batch_id 虽是 UNIQUE 且内嵌 tenantId,但那是命名约定 + 表级约束两个**外部**事实;
  134. // 语句自身带租户谓词才符合 D5,也避免将来 batch_id 命名变更时静默取到别人的行。
  135. return await _db.Ado.GetLongAsync(
  136. "SELECT id FROM mdp_transform_run_log WHERE batch_id=@BatchId AND tenant_id=@TenantId ORDER BY id DESC LIMIT 1",
  137. new List<SugarParameter> { new("@BatchId", batchId), new("@TenantId", tenantId) });
  138. }
  139. private async Task CompleteRunLogAsync(S0DimRefreshResult result, DateTime startedAt)
  140. {
  141. var errors = result.Items.Where(i => i.Status == "FAILED")
  142. .Select(i => $"{i.Key}: {i.Error}").ToList();
  143. await _db.Ado.ExecuteCommandAsync(
  144. """
  145. UPDATE mdp_transform_run_log
  146. SET status=@Status, end_time=@EndTime, duration_ms=@DurationMs,
  147. stage_rows=@StageRows, standard_rows=@StandardRows, dwd_rows=0,
  148. error_message=@Error, summary_json=@Summary, update_time=CURRENT_TIMESTAMP
  149. WHERE id=@Id AND tenant_id=@TenantId
  150. """,
  151. new List<SugarParameter>
  152. {
  153. new("@TenantId", result.TenantId),
  154. new("@Status", result.Status),
  155. new("@EndTime", DateTime.Now),
  156. new("@DurationMs", result.DurationMs),
  157. new("@StageRows", result.Items.Sum(i => i.StagingWritten)),
  158. new("@StandardRows", result.Items.Sum(i => i.DimRows)),
  159. new("@Error", errors.Count == 0 ? null : string.Join(" | ", errors)),
  160. new("@Summary", JsonSerializer.Serialize(result.Items, JsonOpts)),
  161. new("@Id", result.RunLogId)
  162. });
  163. }
  164. private sealed class S0DimRunLogRow
  165. {
  166. public string BatchId { get; set; } = "";
  167. public string Status { get; set; } = "";
  168. public DateTime StartTime { get; set; }
  169. public DateTime? EndTime { get; set; }
  170. public int? DurationMs { get; set; }
  171. public int StageRows { get; set; }
  172. public int StandardRows { get; set; }
  173. public string? ErrorMessage { get; set; }
  174. public string? SummaryJson { get; set; }
  175. }
  176. }