| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183 |
- using SqlSugar;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
- /// <summary>
- /// source / staging / dim 三层对账。
- ///
- /// 分两级:
- /// <list type="bullet">
- /// <item><b>FAIL</b>(阻断)—— 由 <see cref="S0DimMaterializer"/> 在 <b>dim 事务内</b>调用,
- /// 失败即抛异常触发回滚,dim 保持上一轮的完整快照;</item>
- /// <item><b>WARN</b>(不阻断)—— 孤儿子行、源侧无租户行等数据质量项,只记录。</item>
- /// </list>
- ///
- /// 所有查询都显式带租户谓词:source 侧 <c>tenant_id=@TenantId</c>,
- /// staging 侧四段谓词,dim 侧 <c>tenant_id=@TenantId</c>。**不存在跨租户聚合。**
- /// </summary>
- public sealed class S0DimReconciler : ITransient
- {
- /// <summary>孤儿子行最多记录条数(避免 summary_json 膨胀)。</summary>
- public const int OrphanSampleLimit = 50;
- private readonly ISqlSugarClient _db;
- /// <summary>注入本库客户端(SqlSugarScope 单例,事务上下文随异步流传递)。</summary>
- public S0DimReconciler(ISqlSugarClient db) => _db = db;
- private sealed class BizKeyQualityRow
- {
- public int Total { get; set; }
- public int Distinct_Cnt { get; set; }
- public int Blank_Cnt { get; set; }
- }
- private sealed class ChecksumRow
- {
- public int Cnt { get; set; }
- public decimal Chk { get; set; }
- }
- private static List<SugarParameter> Params(S0DimDefinition def, long tenantId, string? batchId)
- {
- var ps = new List<SugarParameter>
- {
- new("@TenantId", tenantId),
- new("@SourceSystem", def.SourceSystem),
- new("@SourceTable", def.SourceTable)
- };
- if (batchId is not null) ps.Add(new SugarParameter("@BatchId", batchId));
- return ps;
- }
- /// <summary>
- /// 完整对账。<paramref name="batchId"/> 非空时按本批口径校验 staging 与 dim 纯度;
- /// 为空时按「当前状态」口径(用于独立的只读对账端点)。
- /// </summary>
- public async Task<S0DimReconcileResult> ReconcileAsync(
- S0DimDefinition def, long tenantId, string? batchId, CancellationToken ct = default)
- {
- if (tenantId <= 0) throw new InvalidOperationException($"[{def.Key}] 对账拒绝无效租户:{tenantId}");
- var withBatch = !string.IsNullOrWhiteSpace(batchId);
- var r = new S0DimReconcileResult { Key = def.Key, TenantId = tenantId, BatchId = batchId };
- var ps = Params(def, tenantId, batchId);
- // ── Layer 1 · 三层行数(全部带租户谓词)
- r.SourceCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildSourceCountSql(def), ps);
- r.StagingCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildStagingCountSql(def, withBatch), ps);
- r.DimCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildDimCountSql(def), ps);
- r.SourceUntenanted = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildSourceUntenantedCountSql(def), ps);
- if (r.StagingCount != r.SourceCount)
- r.Failures.Add($"stg_count({r.StagingCount}) != source_count({r.SourceCount})");
- if (r.DimCount != r.StagingCount)
- r.Failures.Add($"dim_count({r.DimCount}) != stg_count({r.StagingCount})");
- if (r.SourceUntenanted > 0)
- r.Warnings.Add($"源表有 {r.SourceUntenanted} 行无法归属租户(tenant_id IS NULL 或 <=0),整条链路不可见");
- if (withBatch)
- {
- r.DimBatchImpurity = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildDimBatchImpuritySql(def), ps);
- if (r.DimBatchImpurity > 0)
- r.Failures.Add($"dim 存在 {r.DimBatchImpurity} 行非本批残留(FULL Replace 未换干净)");
- }
- // ── Layer 2 · 业务键质量 + 四向集合等价
- var srcSet = S0DimSqlBuilder.SourceBizKeySetSql(def);
- var stgSet = S0DimSqlBuilder.StagingBizKeySetSql(def, withBatch);
- var dimSet = S0DimSqlBuilder.DimBizKeySetSql(def);
- foreach (var (layer, setSql) in new[] { ("source", srcSet), ("staging", stgSet), ("dim", dimSet) })
- {
- var q = (await _db.Ado.SqlQueryAsync<BizKeyQualityRow>(
- S0DimSqlBuilder.BuildBizKeyQualitySql(setSql), ps)).FirstOrDefault() ?? new BizKeyQualityRow();
- var quality = new S0DimBizKeyQuality
- { Layer = layer, Total = q.Total, Distinct = q.Distinct_Cnt, Blank = q.Blank_Cnt };
- r.BizKeyQuality.Add(quality);
- if (quality.Total != quality.Distinct)
- r.Failures.Add($"{layer} 业务键重复:total={quality.Total} distinct={quality.Distinct}");
- if (quality.Blank > 0)
- r.Failures.Add(layer == "staging"
- ? $"staging 有 {quality.Blank} 行 source_biz_key 为空 —— biz_key_expr 已回落到 source_row_id"
- : $"{layer} 有 {quality.Blank} 行业务键为空");
- }
- r.BizKeyDiff = new S0DimBizKeyDiff
- {
- SourceMinusStaging = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(srcSet, stgSet), ps),
- StagingMinusSource = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(stgSet, srcSet), ps),
- StagingMinusDim = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(stgSet, dimSet), ps),
- DimMinusStaging = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildAntiJoinCountSql(dimSet, stgSet), ps)
- };
- if (!r.BizKeyDiff.IsEqual)
- r.Failures.Add(
- $"业务键集合不等价:src\\stg={r.BizKeyDiff.SourceMinusStaging} stg\\src={r.BizKeyDiff.StagingMinusSource} " +
- $"stg\\dim={r.BizKeyDiff.StagingMinusDim} dim\\stg={r.BizKeyDiff.DimMinusStaging}");
- // ── Layer 3 · 属性校验和(source ↔ dim 端到端)
- var srcChk = (await _db.Ado.SqlQueryAsync<ChecksumRow>(S0DimSqlBuilder.BuildSourceChecksumSql(def), ps))
- .FirstOrDefault() ?? new ChecksumRow();
- var dimChk = (await _db.Ado.SqlQueryAsync<ChecksumRow>(S0DimSqlBuilder.BuildDimChecksumSql(def), ps))
- .FirstOrDefault() ?? new ChecksumRow();
- r.SourceChecksum = srcChk.Chk.ToString("F0");
- r.DimChecksum = dimChk.Chk.ToString("F0");
- if (srcChk.Cnt != dimChk.Cnt || srcChk.Chk != dimChk.Chk)
- r.Failures.Add($"属性校验和不一致:source=({srcChk.Cnt},{r.SourceChecksum}) dim=({dimChk.Cnt},{r.DimChecksum})");
- // ── Layer 4 · 结构断言
- r.DimDuplicateBizKeys = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildDimDuplicateBizKeySql(def), ps);
- if (r.DimDuplicateBizKeys > 0)
- r.Failures.Add($"dim 业务键重复组数={r.DimDuplicateBizKeys}");
- if (def.MirrorUniqueColumns is { Count: > 0 })
- {
- r.MirrorUniqueViolations = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildMirrorUniqueViolationSql(def), ps);
- if (r.MirrorUniqueViolations > 0)
- r.Failures.Add(
- $"镜像唯一性违例({string.Join(",", def.MirrorUniqueColumns)})组数={r.MirrorUniqueViolations} —— 源侧唯一性口径可能已放宽");
- }
- // ── 父子完整性(WARN,不阻断:dim 不得对源撒谎,也不静默丢行)
- if (def.ParentDimTable is not null)
- {
- // 真实总数与样本分开取:OrphanChildren 受 OrphanSampleLimit 截断,
- // 直接拿它的 Count 当告警数字会把「97 个孤儿」说成「50 个」——
- // 静默截断比不报还糟,因为读者会以为那就是全部。
- r.OrphanChildCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildOrphanChildCountSql(def), ps);
- var rows = await _db.Ado.GetDataTableAsync(
- S0DimSqlBuilder.BuildOrphanChildSql(def, OrphanSampleLimit), ps);
- foreach (System.Data.DataRow row in rows.Rows)
- {
- var parts = new List<string>();
- foreach (System.Data.DataColumn c in rows.Columns)
- parts.Add(row[c] == DBNull.Value ? "" : row[c].ToString() ?? "");
- r.OrphanChildren.Add(string.Join(S0DimSqlBuilder.BizKeySeparator, parts));
- }
- if (r.OrphanChildCount > 0)
- {
- var truncated = r.OrphanChildCount > r.OrphanChildren.Count
- ? $"(仅外化前 {r.OrphanChildren.Count} 条样本)" : "";
- r.Warnings.Add(
- $"{r.OrphanChildCount} 行子维度在父维度 {def.ParentDimTable} 中无父行(保留并外化,不丢弃){truncated}");
- }
- }
- return r;
- }
- /// <summary>
- /// 事务内阻断级断言:任何 FAIL 直接抛出,交由 <c>MdpStdFullReplace</c> 回滚。
- /// </summary>
- public async Task<S0DimReconcileResult> AssertBlockingAsync(
- S0DimDefinition def, long tenantId, string batchId, CancellationToken ct = default)
- {
- var r = await ReconcileAsync(def, tenantId, batchId, ct);
- if (!r.IsPass)
- throw new InvalidOperationException(
- $"[{def.Key}] tenant={tenantId} batch={batchId} 对账失败,已回滚:{string.Join(";", r.Failures)}");
- return r;
- }
- }
|