using Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
using Microsoft.Extensions.Logging;
using SqlSugar;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
///
/// 单个 S0 维度的物化:source → staging → dim,FULL REPLACE 语义。
///
/// 时序与事务边界:
///
/// [0] 前置校验(只读):mdp_entity 配置与 definition 一致 + 源侧业务键无重复 + 数据量未超单页
/// [1] source_count(只读,按租户)
/// ── 阶段 I:staging(非事务,可重跑)
/// [2] purge 本租户+本源系统+本源表 分区
/// [3] FULL pull(源侧/写入侧双重租户过滤)
/// [4] staging_count(本批四段谓词)
/// ── 阶段 II:闸门
/// [5] stg_count != source_count → 本维度 FAILED,**dim 不动**(保持上一轮完整快照)
/// ── 阶段 III:dim(单事务)
/// [6] DELETE tenant 行 + INSERT(无 ON DUPLICATE KEY)+ 事务内阻断级对账 → 任一失败整体回滚
///
///
/// 失败时 dim 永远停在**上一轮的完整快照**,不会半新半旧。
///
public sealed class S0DimMaterializer : ITransient
{
private readonly ISqlSugarClient _db;
private readonly MdpSourcePullDispatcher _pullDispatcher;
private readonly S0DimStagingPurge _purge;
private readonly S0DimReconciler _reconciler;
private readonly ILogger _logger;
/// 构造。
public S0DimMaterializer(
ISqlSugarClient db,
MdpSourcePullDispatcher pullDispatcher,
S0DimStagingPurge purge,
S0DimReconciler reconciler,
ILogger logger)
{
_db = db;
_pullDispatcher = pullDispatcher;
_purge = purge;
_reconciler = reconciler;
_logger = logger;
}
/// 物化一个维度。
public async Task RunAsync(
S0DimDefinition def, long tenantId, string batchId, CancellationToken ct = default)
{
def.Validate();
if (tenantId <= 0) throw new InvalidOperationException($"[{def.Key}] 拒绝无效租户:{tenantId}");
if (string.IsNullOrWhiteSpace(batchId)) throw new InvalidOperationException($"[{def.Key}] batchId 不得为空");
var result = new S0DimMaterializeResult { Key = def.Key, TenantId = tenantId, BatchId = batchId };
var ps = new List
{
new("@TenantId", tenantId),
new("@SourceSystem", def.SourceSystem),
new("@SourceTable", def.SourceTable),
new("@BatchId", batchId)
};
try
{
// [0] 前置校验:配置漂移与源侧数据质量,全部在动任何数据之前
var batchSize = await AssertEntityContractAsync(def, ct);
var sourceCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildSourceCountSql(def), ps);
await AssertNoSourceDuplicateAsync(def, tenantId, sourceCount, batchSize, ps);
// [2] purge:三段谓词精确到本租户 + 本源系统 + 本源表
result.StagingPurged = await _purge.PurgeAsync(def, tenantId, ct);
// [3] FULL pull
// 租户三重防护:① MdpDbPullExecutor 对含租户列的源表加 WHERE tenant_id=@scopeTenantId;
// ② MdpStagingWriter 对 tenantValue != ctx.TenantId 的行直接跳过;
// ③ RequireMatchingSourceTenant=true → 源行无租户时不拿 ctx 兜底,直接跳过(绝不落 0)
var pullCtx = new MdpPullContext
{
TenantId = tenantId,
BatchId = batchId,
FullRefresh = true,
RequireMatchingSourceTenant = true
};
var pull = await _pullDispatcher.PullAllByEntityCodeAsync(def.EntityCode, pullCtx, ct);
result.SourcePulled = pull.RowsPulled;
result.StagingWritten = pull.RowsWritten;
// [4][5] 闸门:本批 staging 必须与源侧行数完全一致,否则不碰 dim
var stagingCount = await _db.Ado.GetIntAsync(S0DimSqlBuilder.BuildStagingCountSql(def), ps);
if (stagingCount != sourceCount)
throw new InvalidOperationException(
$"[{def.Key}] 贴源闸门未通过:source_count={sourceCount} 本批 stg_count={stagingCount}(dim 未改动)");
// [6] dim 单事务:DELETE 本租户 → INSERT → 事务内阻断级对账
result.DimRows = await MdpStdFullReplace.ReplaceAsync(
_db, def.DimTable, tenantId, extraWhere: null,
insertScopedAsync: async () =>
{
var now = DateTime.Now;
var insertPs = new List(ps) { new("@Now", now) };
var rows = await _db.Ado.ExecuteCommandAsync(S0DimSqlBuilder.BuildInsertSql(def), insertPs);
result.Reconcile = await _reconciler.AssertBlockingAsync(def, tenantId, batchId, ct);
return rows;
},
ct);
result.Status = "SUCCESS";
_logger.LogInformation(
"[S0Dim] {Key} tenant={Tenant} batch={Batch} purged={Purged} pulled={Pulled} stg={Stg} dim={Dim}",
def.Key, tenantId, batchId, result.StagingPurged, result.SourcePulled, stagingCount, result.DimRows);
}
catch (OperationCanceledException)
{
// 取消不是本维度的「失败」:吞掉它会把客户端断开记成 FAILED,
// 且外层 RefreshAsync 的 ct.ThrowIfCancellationRequested() 会在下一轮抛出,
// 导致 CompleteRunLogAsync 永不执行、run log 永久停在 RUNNING。
// 事务侧无需担心:异常已在 MdpStdFullReplace 的 catch 中回滚,dim 未改动。
throw;
}
catch (Exception ex)
{
result.Status = "FAILED";
result.Error = ex.Message;
_logger.LogError(ex, "[S0Dim] {Key} tenant={Tenant} batch={Batch} 物化失败(dim 未改动)",
def.Key, tenantId, batchId);
}
return result;
}
///
/// 断言 mdp_entity 的运行时配置与 definition 完全一致,返回 batch_size。
/// 这一步把「配置漂移」变成显式失败 —— 否则 PullAll 可能悄悄拉到别的表 / 别的源。
///
private async Task AssertEntityContractAsync(S0DimDefinition def, CancellationToken ct)
{
var entity = await _db.Queryable()
.Where(x => x.EntityCode == def.EntityCode)
.FirstAsync(ct)
?? throw new InvalidOperationException($"[{def.Key}] mdp_entity 未登记:{def.EntityCode}(migration 未执行?)");
if (entity.Status != 1)
throw new InvalidOperationException($"[{def.Key}] mdp_entity.{def.EntityCode} 已停用(status={entity.Status})");
var source = await _db.Queryable().Where(x => x.Id == entity.SourceId).FirstAsync(ct)
?? throw new InvalidOperationException($"[{def.Key}] mdp_source id={entity.SourceId} 不存在");
void Expect(string what, string? actual, string expected)
{
if (!string.Equals(actual, expected, StringComparison.Ordinal))
throw new InvalidOperationException($"[{def.Key}] mdp_entity 配置漂移:{what} 实际='{actual}' 期望='{expected}'");
}
Expect("source_code", source.SourceCode, def.SourceSystem);
Expect("source_table_name", entity.SourceTableName, def.SourceTable);
Expect("target_table_name", entity.TargetTableName, def.StagingTable);
Expect("biz_key_expr", entity.BizKeyExpr, string.Join(",", def.SourceBizKeyColumns));
// ── 执行器选路的三个开关:不查会被静默改道 ──
// MdpSourcePullDispatcher 按
// !IsNullOrWhiteSpace(entity.SourceApiPath) || source.SourceType == "API" → _apiExecutor
// source.SourceType == "FILE_EXCEL" → Excel 分支
// 选择执行器。上面那几条 Expect 全部**照样通过**,但拉取已不再是「本库 SELECT」。
// 比较口径与 dispatcher 一致(OrdinalIgnoreCase),否则大小写差异会造成假通过。
if (!string.IsNullOrWhiteSpace(entity.SourceApiPath))
throw new InvalidOperationException(
$"[{def.Key}] mdp_entity.source_api_path 必须为空,实际='{entity.SourceApiPath}':" +
"非空会让 MdpSourcePullDispatcher 改走 API 执行器,绕开本库 DB 拉取契约");
if (!string.Equals(source.SourceType, "DB", StringComparison.OrdinalIgnoreCase))
throw new InvalidOperationException(
$"[{def.Key}] mdp_source.source_type 必须为 DB,实际='{source.SourceType}'");
if (!string.Equals(source.DbType, "MySQL", StringComparison.OrdinalIgnoreCase))
throw new InvalidOperationException(
$"[{def.Key}] mdp_source.db_type 必须为 MySQL,实际='{source.DbType}'");
// 本阶段只做 FULL:留了 incr_column 会让 BuildSelectSql 生成 "col > @cursor",
// 使源侧水位为 NULL 的行(如 DepartmentMaster 的 UATDEMO/未分配)永久不可达
if (!string.Equals(entity.SyncMode, "FULL", StringComparison.OrdinalIgnoreCase))
throw new InvalidOperationException($"[{def.Key}] sync_mode 必须为 FULL,实际='{entity.SyncMode}'");
if (!string.IsNullOrWhiteSpace(entity.IncrColumn))
throw new InvalidOperationException($"[{def.Key}] incr_column 必须为空(FULL 语义),实际='{entity.IncrColumn}'");
return entity.BatchSize > 0 ? entity.BatchSize : 1000;
}
///
/// 源侧业务键重复探针 + 单页容量检查。两者都在 purge 之前,失败时**任何数据都未被触碰**。
///
private async Task AssertNoSourceDuplicateAsync(
S0DimDefinition def, long tenantId, int sourceCount, int batchSize, List ps)
{
// 分页依赖 MdpDbPullExecutor 无 incr_column 时的 "ORDER BY 1"(按第 1 个物理列),
// 而 LocationMaster 的第 1 列是 Capacity(非唯一)→ 一旦真的翻页,OFFSET 结果不稳定。
// 因此要求单页装得下;超出时显式失败并提示调大 mdp_entity.batch_size。
if (sourceCount >= batchSize)
throw new InvalidOperationException(
$"[{def.Key}] 源行数 {sourceCount} 已达单页容量 {batchSize}:" +
"通用执行器在无 incr_column 时按 'ORDER BY 1' 分页,对本表非稳定序,请调大 mdp_entity.batch_size");
var q = await _db.Ado.SqlQueryAsync(
S0DimSqlBuilder.BuildBizKeyQualitySql(S0DimSqlBuilder.SourceBizKeySetSql(def)), ps);
var probe = q.FirstOrDefault();
if (probe is null) return;
if (probe.Blank_Cnt > 0)
throw new InvalidOperationException(
$"[{def.Key}] tenant={tenantId} 源侧有 {probe.Blank_Cnt} 行业务键为空,拒绝物化");
if (probe.Total != probe.Distinct_Cnt)
throw new InvalidOperationException(
$"[{def.Key}] tenant={tenantId} 源侧业务键重复(total={probe.Total} distinct={probe.Distinct_Cnt}):" +
"主数据唯一性已被破坏,transform 显式失败 —— 不做 first-wins / last-wins 静默取舍");
}
private sealed class SourceBizKeyProbe
{
public int Total { get; set; }
public int Distinct_Cnt { get; set; }
public int Blank_Cnt { get; set; }
}
}