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;
}
}