S0DimReconciler.cs 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220
  1. using SqlSugar;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
  3. /// <summary>
  4. /// source / staging / dim 三层对账。
  5. ///
  6. /// 分两级:
  7. /// <list type="bullet">
  8. /// <item><b>FAIL</b>(阻断)—— 由 <see cref="S0DimMaterializer"/> 在 <b>dim 事务内</b>调用,
  9. /// 失败即抛异常触发回滚,dim 保持上一轮的完整快照;</item>
  10. /// <item><b>WARN</b>(不阻断)—— 孤儿子行、源侧无租户行等数据质量项,只记录。</item>
  11. /// </list>
  12. ///
  13. /// 所有查询都显式带租户谓词:source 侧 <c>tenant_id=@TenantId</c>,
  14. /// staging 侧四段谓词,dim 侧 <c>tenant_id=@TenantId</c>。**不存在跨租户聚合。**
  15. /// </summary>
  16. public sealed class S0DimReconciler : ITransient
  17. {
  18. /// <summary>孤儿子行最多记录条数(避免 summary_json 膨胀)。</summary>
  19. public const int OrphanSampleLimit = 50;
  20. private readonly ISqlSugarClient _db;
  21. /// <summary>注入本库客户端(SqlSugarScope 单例,事务上下文随异步流传递)。</summary>
  22. public S0DimReconciler(ISqlSugarClient db) => _db = db;
  23. private sealed class BizKeyQualityRow
  24. {
  25. public int Total { get; set; }
  26. public int Distinct_Cnt { get; set; }
  27. public int Blank_Cnt { get; set; }
  28. }
  29. private sealed class ChecksumRow
  30. {
  31. public int Cnt { get; set; }
  32. public decimal Chk { get; set; }
  33. }
  34. private static List<SugarParameter> Params(S0DimDefinition def, long tenantId, string? batchId)
  35. {
  36. var ps = new List<SugarParameter>
  37. {
  38. new("@TenantId", tenantId),
  39. new("@SourceSystem", def.SourceSystem),
  40. new("@SourceTable", def.SourceTable)
  41. };
  42. if (batchId is not null) ps.Add(new SugarParameter("@BatchId", batchId));
  43. return ps;
  44. }
  45. /// <summary>
  46. /// 读取 source / dim 两侧的属性校验和读数(<c>count</c> + <c>checksum</c> 成对)。**纯只读,不写任何数据。**
  47. ///
  48. /// <para>两个调用点共用本方法,这是刻意的:</para>
  49. /// <list type="number">
  50. /// <item><see cref="ReconcileAsync"/> 的 Layer 3 —— dim 事务内的**写后**对账;</item>
  51. /// <item><see cref="S0DimMaterializer"/> 的**写前**闸门 —— 两侧一致即整轮跳过。</item>
  52. /// </list>
  53. ///
  54. /// <para>闸门的正确性依赖「写前判据 ⊆ 写后判据」:写后对账断言了 source == dim,
  55. /// 故任何一轮 SUCCESS 之后两者必然相等;若下一轮开跑时**仍**相等,
  56. /// 说明源侧自上轮以来没有任何会影响 dim 的变化。
  57. /// 一旦两处各写一份读取逻辑,这条推理就会在某次单边修改后静默失效 ——
  58. /// 所以只允许存在这一份。</para>
  59. ///
  60. /// <para>只用到 <c>@TenantId</c>(两条 SQL 的租户谓词),不涉及批次,
  61. /// 因此可以在 purge / pull 之前、staging 尚未装载时安全调用。</para>
  62. /// </summary>
  63. public async Task<S0DimChecksumPair> ReadChecksumPairAsync(
  64. S0DimDefinition def, long tenantId, CancellationToken ct = default)
  65. {
  66. if (tenantId <= 0) throw new InvalidOperationException($"[{def.Key}] 校验和读取拒绝无效租户:{tenantId}");
  67. var ps = Params(def, tenantId, null);
  68. var srcChk = (await _db.Ado.SqlQueryAsync<ChecksumRow>(S0DimSqlBuilder.BuildSourceChecksumSql(def), ps))
  69. .FirstOrDefault() ?? new ChecksumRow();
  70. var dimChk = (await _db.Ado.SqlQueryAsync<ChecksumRow>(S0DimSqlBuilder.BuildDimChecksumSql(def), ps))
  71. .FirstOrDefault() ?? new ChecksumRow();
  72. return new S0DimChecksumPair
  73. {
  74. SourceCount = srcChk.Cnt,
  75. SourceChecksumValue = srcChk.Chk,
  76. DimCount = dimChk.Cnt,
  77. DimChecksumValue = dimChk.Chk
  78. };
  79. }
  80. /// <summary>
  81. /// 完整对账。<paramref name="batchId"/> 非空时按本批口径校验 staging 与 dim 纯度;
  82. /// 为空时按「当前状态」口径(用于独立的只读对账端点)。
  83. /// </summary>
  84. public async Task<S0DimReconcileResult> ReconcileAsync(
  85. S0DimDefinition def, long tenantId, string? batchId, CancellationToken ct = default)
  86. {
  87. if (tenantId <= 0) throw new InvalidOperationException($"[{def.Key}] 对账拒绝无效租户:{tenantId}");
  88. var withBatch = !string.IsNullOrWhiteSpace(batchId);
  89. var r = new S0DimReconcileResult { Key = def.Key, TenantId = tenantId, BatchId = batchId };
  90. var ps = Params(def, tenantId, batchId);
  91. // ── Layer 1 · 三层行数(全部带租户谓词)
  92. r.SourceCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildSourceCountSql(def), ps);
  93. r.StagingCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildStagingCountSql(def, withBatch), ps);
  94. r.DimCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildDimCountSql(def), ps);
  95. r.SourceUntenanted = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildSourceUntenantedCountSql(def), ps);
  96. if (r.StagingCount != r.SourceCount)
  97. r.Failures.Add($"stg_count({r.StagingCount}) != source_count({r.SourceCount})");
  98. if (r.DimCount != r.StagingCount)
  99. r.Failures.Add($"dim_count({r.DimCount}) != stg_count({r.StagingCount})");
  100. if (r.SourceUntenanted > 0)
  101. r.Warnings.Add($"源表有 {r.SourceUntenanted} 行无法归属租户(tenant_id IS NULL 或 <=0),整条链路不可见");
  102. if (withBatch)
  103. {
  104. r.DimBatchImpurity = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildDimBatchImpuritySql(def), ps);
  105. if (r.DimBatchImpurity > 0)
  106. r.Failures.Add($"dim 存在 {r.DimBatchImpurity} 行非本批残留(FULL Replace 未换干净)");
  107. }
  108. // ── Layer 2 · 业务键质量 + 四向集合等价
  109. var srcSet = S0DimSqlBuilder.SourceBizKeySetSql(def);
  110. var stgSet = S0DimSqlBuilder.StagingBizKeySetSql(def, withBatch);
  111. var dimSet = S0DimSqlBuilder.DimBizKeySetSql(def);
  112. foreach (var (layer, setSql) in new[] { ("source", srcSet), ("staging", stgSet), ("dim", dimSet) })
  113. {
  114. var q = (await _db.Ado.SqlQueryAsync<BizKeyQualityRow>(
  115. S0DimSqlBuilder.BuildBizKeyQualitySql(setSql), ps)).FirstOrDefault() ?? new BizKeyQualityRow();
  116. var quality = new S0DimBizKeyQuality
  117. { Layer = layer, Total = q.Total, Distinct = q.Distinct_Cnt, Blank = q.Blank_Cnt };
  118. r.BizKeyQuality.Add(quality);
  119. if (quality.Total != quality.Distinct)
  120. r.Failures.Add($"{layer} 业务键重复:total={quality.Total} distinct={quality.Distinct}");
  121. if (quality.Blank > 0)
  122. r.Failures.Add(layer == "staging"
  123. ? $"staging 有 {quality.Blank} 行 source_biz_key 为空 —— biz_key_expr 已回落到 source_row_id"
  124. : $"{layer} 有 {quality.Blank} 行业务键为空");
  125. }
  126. r.BizKeyDiff = new S0DimBizKeyDiff
  127. {
  128. SourceMinusStaging = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(srcSet, stgSet), ps),
  129. StagingMinusSource = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(stgSet, srcSet), ps),
  130. StagingMinusDim = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(stgSet, dimSet), ps),
  131. DimMinusStaging = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(dimSet, stgSet), ps)
  132. };
  133. if (!r.BizKeyDiff.IsEqual)
  134. r.Failures.Add(
  135. $"业务键集合不等价:src\\stg={r.BizKeyDiff.SourceMinusStaging} stg\\src={r.BizKeyDiff.StagingMinusSource} " +
  136. $"stg\\dim={r.BizKeyDiff.StagingMinusDim} dim\\stg={r.BizKeyDiff.DimMinusStaging}");
  137. // ── Layer 3 · 属性校验和(source ↔ dim 端到端)
  138. // 读取走 ReadChecksumPairAsync —— 与 S0DimMaterializer 的写前闸门是**同一条路径**,
  139. // 不存在「写前一套算法、写后另一套」的漂移空间。
  140. var chk = await ReadChecksumPairAsync(def, tenantId, ct);
  141. r.SourceChecksum = chk.SourceChecksum;
  142. r.DimChecksum = chk.DimChecksum;
  143. if (!chk.IsMatch)
  144. r.Failures.Add($"属性校验和不一致:source=({chk.SourceCount},{r.SourceChecksum}) dim=({chk.DimCount},{r.DimChecksum})");
  145. // ── Layer 4 · 结构断言
  146. r.DimDuplicateBizKeys = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildDimDuplicateBizKeySql(def), ps);
  147. if (r.DimDuplicateBizKeys > 0)
  148. r.Failures.Add($"dim 业务键重复组数={r.DimDuplicateBizKeys}");
  149. if (def.MirrorUniqueColumns is { Count: > 0 })
  150. {
  151. r.MirrorUniqueViolations = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildMirrorUniqueViolationSql(def), ps);
  152. if (r.MirrorUniqueViolations > 0)
  153. r.Failures.Add(
  154. $"镜像唯一性违例({string.Join(",", def.MirrorUniqueColumns)})组数={r.MirrorUniqueViolations} —— 源侧唯一性口径可能已放宽");
  155. }
  156. // ── 父子完整性(WARN,不阻断:dim 不得对源撒谎,也不静默丢行)
  157. if (def.ParentDimTable is not null)
  158. {
  159. // 真实总数与样本分开取:OrphanChildren 受 OrphanSampleLimit 截断,
  160. // 直接拿它的 Count 当告警数字会把「97 个孤儿」说成「50 个」——
  161. // 静默截断比不报还糟,因为读者会以为那就是全部。
  162. r.OrphanChildCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildOrphanChildCountSql(def), ps);
  163. var rows = await _db.Ado.GetDataTableAsync(
  164. S0DimSqlBuilder.BuildOrphanChildSql(def, OrphanSampleLimit), ps);
  165. foreach (System.Data.DataRow row in rows.Rows)
  166. {
  167. var parts = new List<string>();
  168. foreach (System.Data.DataColumn c in rows.Columns)
  169. parts.Add(row[c] == DBNull.Value ? "" : row[c].ToString() ?? "");
  170. r.OrphanChildren.Add(string.Join(S0DimSqlBuilder.BizKeySeparator, parts));
  171. }
  172. if (r.OrphanChildCount > 0)
  173. {
  174. var truncated = r.OrphanChildCount > r.OrphanChildren.Count
  175. ? $"(仅外化前 {r.OrphanChildren.Count} 条样本)" : "";
  176. r.Warnings.Add(
  177. $"{r.OrphanChildCount} 行子维度在父维度 {def.ParentDimTable} 中无父行(保留并外化,不丢弃){truncated}");
  178. }
  179. }
  180. return r;
  181. }
  182. /// <summary>
  183. /// 事务内阻断级断言:任何 FAIL 直接抛出,交由 <c>MdpStdFullReplace</c> 回滚。
  184. /// </summary>
  185. public async Task<S0DimReconcileResult> AssertBlockingAsync(
  186. S0DimDefinition def, long tenantId, string batchId, CancellationToken ct = default)
  187. {
  188. var r = await ReconcileAsync(def, tenantId, batchId, ct);
  189. if (!r.IsPass)
  190. throw new InvalidOperationException(
  191. $"[{def.Key}] tenant={tenantId} batch={batchId} 对账失败,已回滚:{string.Join(";", r.Failures)}");
  192. return r;
  193. }
  194. }