using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim; /// /// source / staging / dim 三层对账。 /// /// 分两级: /// /// FAIL(阻断)—— 由 dim 事务内调用, /// 失败即抛异常触发回滚,dim 保持上一轮的完整快照; /// WARN(不阻断)—— 孤儿子行、源侧无租户行等数据质量项,只记录。 /// /// /// 所有查询都显式带租户谓词:source 侧 tenant_id=@TenantId, /// staging 侧四段谓词,dim 侧 tenant_id=@TenantId。**不存在跨租户聚合。** /// public sealed class S0DimReconciler : ITransient { /// 孤儿子行最多记录条数(避免 summary_json 膨胀)。 public const int OrphanSampleLimit = 50; private readonly ISqlSugarClient _db; /// 注入本库客户端(SqlSugarScope 单例,事务上下文随异步流传递)。 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 Params(S0DimDefinition def, long tenantId, string? batchId) { var ps = new List { new("@TenantId", tenantId), new("@SourceSystem", def.SourceSystem), new("@SourceTable", def.SourceTable) }; if (batchId is not null) ps.Add(new SugarParameter("@BatchId", batchId)); return ps; } /// /// 完整对账。 非空时按本批口径校验 staging 与 dim 纯度; /// 为空时按「当前状态」口径(用于独立的只读对账端点)。 /// public async Task 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( 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(S0DimSqlBuilder.BuildSourceChecksumSql(def), ps)) .FirstOrDefault() ?? new ChecksumRow(); var dimChk = (await _db.Ado.SqlQueryAsync(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(); 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; } /// /// 事务内阻断级断言:任何 FAIL 直接抛出,交由 MdpStdFullReplace 回滚。 /// public async Task 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; } }