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;
}
///
/// 读取 source / dim 两侧的属性校验和读数(count + checksum 成对)。**纯只读,不写任何数据。**
///
/// 两个调用点共用本方法,这是刻意的:
///
/// - 的 Layer 3 —— dim 事务内的**写后**对账;
/// - 的**写前**闸门 —— 两侧一致即整轮跳过。
///
///
/// 闸门的正确性依赖「写前判据 ⊆ 写后判据」:写后对账断言了 source == dim,
/// 故任何一轮 SUCCESS 之后两者必然相等;若下一轮开跑时**仍**相等,
/// 说明源侧自上轮以来没有任何会影响 dim 的变化。
/// 一旦两处各写一份读取逻辑,这条推理就会在某次单边修改后静默失效 ——
/// 所以只允许存在这一份。
///
/// 只用到 @TenantId(两条 SQL 的租户谓词),不涉及批次,
/// 因此可以在 purge / pull 之前、staging 尚未装载时安全调用。
///
public async Task ReadChecksumPairAsync(
S0DimDefinition def, long tenantId, CancellationToken ct = default)
{
if (tenantId <= 0) throw new InvalidOperationException($"[{def.Key}] 校验和读取拒绝无效租户:{tenantId}");
var ps = Params(def, tenantId, null);
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();
return new S0DimChecksumPair
{
SourceCount = srcChk.Cnt,
SourceChecksumValue = srcChk.Chk,
DimCount = dimChk.Cnt,
DimChecksumValue = dimChk.Chk
};
}
///
/// 完整对账。 非空时按本批口径校验 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 端到端)
// 读取走 ReadChecksumPairAsync —— 与 S0DimMaterializer 的写前闸门是**同一条路径**,
// 不存在「写前一套算法、写后另一套」的漂移空间。
var chk = await ReadChecksumPairAsync(def, tenantId, ct);
r.SourceChecksum = chk.SourceChecksum;
r.DimChecksum = chk.DimChecksum;
if (!chk.IsMatch)
r.Failures.Add($"属性校验和不一致:source=({chk.SourceCount},{r.SourceChecksum}) dim=({chk.DimCount},{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;
}
}