using System.Text;
namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
///
/// 由 生成 SQL 的**唯一**位置(纯函数,不接触 DB,可单测)。
///
/// 生成的所有语句都必须能逐条证明租户范围:
/// staging 侧固定四段谓词 tenant_id / source_system / source_table / sync_batch_id;
/// dim 侧固定 tenant_id;source 侧固定 tenant_id。
///
public static class S0DimSqlBuilder
{
/// staging 别名。
public const string StgAlias = "s";
/// dim 别名。
public const string DimAlias = "d";
/// source 别名。
public const string SrcAlias = "a";
/// 业务键各段之间的分隔符,与 MdpStagingWriter.BuildBizKey 的 string.Join("#") 一致。
public const string BizKeySeparator = "#";
/// checksum 行内字段分隔(US,业务值不可能包含)。
private const string FieldSep = "CHAR(31 USING utf8mb4)";
/// checksum 的 NULL 哨兵(RS),与空串可区分。
private const string NullSentinel = "CHAR(30 USING utf8mb4)";
private static string Q(string ident)
{
if (!S0DimDefinition.IdentifierRe.IsMatch(ident))
throw new InvalidOperationException($"非法标识符:{ident}");
return $"`{ident}`";
}
/// 列取值表达式(作用于 staging 的 raw_data;TenantIdColumn 例外)。
private static string ValueExpr(S0DimColumn c) => c.Kind switch
{
S0DimValueKind.TenantIdColumn => $"{StgAlias}.{Q(S0DimDefinition.TenantColumn)}",
S0DimValueKind.Str => MdpJsonSql.Str(StgAlias, c.JsonPath!),
S0DimValueKind.Int => MdpJsonSql.Int(StgAlias, c.JsonPath!),
S0DimValueKind.Dec => MdpJsonSql.Dec(StgAlias, c.JsonPath!, c.Precision, c.Scale),
S0DimValueKind.DateTimeSec => MdpJsonSql.DateTimeSec(StgAlias, c.JsonPath!),
S0DimValueKind.BoolTrue => MdpJsonSql.BoolTrue(StgAlias, c.JsonPath!),
_ => throw new InvalidOperationException($"未支持的取值方式:{c.Kind}")
};
///
/// staging 四段固定谓词。 为 false 时用于「不限本批」的只读对账。
///
public static string StagingFilter(bool withBatch = true)
{
var sb = new StringBuilder();
sb.Append($"{StgAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId");
sb.Append($" AND {StgAlias}.`source_system` = @SourceSystem");
sb.Append($" AND {StgAlias}.`source_table` = @SourceTable");
if (withBatch) sb.Append($" AND {StgAlias}.`sync_batch_id` = @BatchId");
return sb.ToString();
}
/// 行级过滤:Required 列非空 + 布尔列可辨识(不满足的行不进 dim,随后被行数对账捕获)。
private static IReadOnlyList RowGuards(S0DimDefinition def)
{
var guards = new List();
foreach (var c in def.Columns)
{
if (c.Kind == S0DimValueKind.TenantIdColumn) continue;
if (c.Required)
{
guards.Add($"({ValueExpr(c)}) IS NOT NULL");
// IS NOT NULL 挡不住空串。业务键走 CONCAT_WS 拼接,空段会被**保留**,
// 于是 Domain='' 会产出 '#LOC01' 这种「看起来合法」的键,而
// AssertNoSourceDuplicateAsync 的 blank_cnt 只捕获整串全空的情况。
// 当前 Catalog 中 Required 的 Str 列恰好就是全部业务键分量
// (domain_code / work_center_code / department_code / location_code / shelf_code),
// 故此守卫命中面 = 业务键,无附带影响。
if (c.Kind == S0DimValueKind.Str)
guards.Add($"({ValueExpr(c)}) <> ''");
}
if (c.Kind == S0DimValueKind.BoolTrue)
{
// MdpJsonSql.BoolTrue 对 null / 不可识别值静默返回 0;此处强制要求原始值可辨识,
// 否则整行被过滤 → dim_count < stg_count → 对账 FAIL(暴露而非吞掉)
var raw = MdpJsonSql.Raw(StgAlias, c.JsonPath!);
guards.Add($"LOWER({raw}) IN ('1','0','true','false')");
}
}
return guards;
}
///
/// dim 物化语句:INSERT INTO dim (...) SELECT ... FROM staging。
/// **刻意不带 ON DUPLICATE KEY UPDATE** —— 源侧业务键重复时直接撞 dim 唯一键报错,
/// 由外层事务回滚(transform fail),绝不 last-wins。
/// 参数:@TenantId @SourceSystem @SourceTable @BatchId @Now
///
public static string BuildInsertSql(S0DimDefinition def)
{
def.Validate();
var cols = new List();
var vals = new List();
foreach (var c in def.Columns)
{
cols.Add(Q(c.TargetColumn));
vals.Add(ValueExpr(c));
}
cols.AddRange(["`source_system`", "`source_biz_key`", "`sync_batch_id`", "`sync_time`"]);
vals.AddRange([$"{StgAlias}.`source_system`", $"{StgAlias}.`source_biz_key`", "@BatchId", "@Now"]);
var where = new List { StagingFilter() };
where.AddRange(RowGuards(def));
return $"""
INSERT INTO {Q(def.DimTable)}
({string.Join(", ", cols)})
SELECT
{string.Join(",\n ", vals)}
FROM {Q(def.StagingTable)} {StgAlias}
WHERE {string.Join("\n AND ", where)}
""";
}
/// staging 分区清理(purge)。参数:@TenantId @SourceSystem @SourceTable
public static string BuildPurgeStagingSql(S0DimDefinition def) =>
$"""
DELETE {StgAlias} FROM {Q(def.StagingTable)} {StgAlias}
WHERE {StagingFilter(withBatch: false)}
""";
/// 源表行数(按租户)。参数:@TenantId
public static string BuildSourceCountSql(S0DimDefinition def) =>
$"SELECT COUNT(*) FROM {Q(def.SourceTable)} {SrcAlias} " +
$"WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
/// 源表中无法归属租户的行数(tenant_id IS NULL 或 <= 0)—— 这些行整条链路都看不见。
public static string BuildSourceUntenantedCountSql(S0DimDefinition def) =>
$"SELECT COUNT(*) FROM {Q(def.SourceTable)} {SrcAlias} " +
$"WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} IS NULL " +
$"OR {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} <= 0";
/// staging 行数。参数:@TenantId @SourceSystem @SourceTable [@BatchId]
public static string BuildStagingCountSql(S0DimDefinition def, bool withBatch = true) =>
$"SELECT COUNT(*) FROM {Q(def.StagingTable)} {StgAlias} WHERE {StagingFilter(withBatch)}";
/// dim 行数(按租户)。参数:@TenantId
public static string BuildDimCountSql(S0DimDefinition def) =>
$"SELECT COUNT(*) FROM {Q(def.DimTable)} {DimAlias} " +
$"WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
/// dim 中不属于本批的残留行数(FULL Replace 之后必须为 0)。参数:@TenantId @BatchId
public static string BuildDimBatchImpuritySql(S0DimDefinition def) =>
$"SELECT COUNT(*) FROM {Q(def.DimTable)} {DimAlias} " +
$"WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId " +
$"AND {DimAlias}.`sync_batch_id` <> @BatchId";
/// 源侧业务键表达式(物理列,与 mdp_entity.biz_key_expr 同序同分隔)。
public static string SourceBizKeyExpr(S0DimDefinition def) =>
$"CONCAT_WS('{BizKeySeparator}', {string.Join(", ", def.SourceBizKeyColumns.Select(c => $"{SrcAlias}.{Q(c)}"))})";
/// dim 侧业务键表达式(去掉 tenant_id 后的业务键列)。
public static string DimBizKeyExpr(S0DimDefinition def) =>
$"CONCAT_WS('{BizKeySeparator}', {string.Join(", ", def.BusinessKeyColumns.Skip(1).Select(c => $"{DimAlias}.{Q(c)}"))})";
/// 三层业务键集合来源子查询。
public static string SourceBizKeySetSql(S0DimDefinition def) =>
$"SELECT {SourceBizKeyExpr(def)} AS bk FROM {Q(def.SourceTable)} {SrcAlias} " +
$"WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
///
public static string StagingBizKeySetSql(S0DimDefinition def, bool withBatch = true) =>
$"SELECT {StgAlias}.`source_biz_key` AS bk FROM {Q(def.StagingTable)} {StgAlias} WHERE {StagingFilter(withBatch)}";
///
public static string DimBizKeySetSql(S0DimDefinition def) =>
$"SELECT {DimBizKeyExpr(def)} AS bk FROM {Q(def.DimTable)} {DimAlias} " +
$"WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
///
/// 集合差集计数: \ ,
/// 返回**左侧无匹配的行数**(不是去重后的键数)。
///
/// 🔴 不要写成 LEFT JOIN … ON y.bk <=> x.bk WHERE y.bk IS NULL(2026-09-07 实测):
/// 四个方向里总有一侧的 bk 是 CONCAT_WS(...) 计算表达式,**既不能建索引也不能哈希**。
/// 当它落在 JOIN 右侧时,优化器只能退化成 nested loop —— 16150×16150 行的
/// stg \ src 方向实测 130.3s(EXPLAIN: rows=32,762,720),
/// 直接撞穿 SqlSugarSetup.CommandTimeOut = 30,而对称的 src \ stg 方向因右侧走
/// idx_stg_*_biz 只要 0.5s —— 同一个缺陷只在一半方向上显形,小表永远发现不了。
/// 把 <=> 换成 = 无用(实测 127.7s),NOT EXISTS 也只到 69.2s。
///
/// 本写法把两侧各扫一遍后按 bk 聚合,复杂度 O(N+M),与方向无关:
/// 同一 stg \ src 实测 0.5s(260×),EXPLAIN 里 nested loop 消失。
///
/// 语义与旧写法逐条对齐:① 未匹配的左行在 LEFT JOIN 下恰好产出 1 行,
/// 故 SUM(l_cnt) 与 COUNT(*) 等价,**左右两侧的重复行都不会改变结果**;
/// ② GROUP BY 把所有 NULL 归为同一组,天然等价于 <=> 的 NULL 相等语义。
///
public static string BuildAntiJoinCountSql(string leftSql, string rightSql) =>
$"""
SELECT COALESCE(SUM(CASE WHEN g.r_cnt = 0 THEN g.l_cnt ELSE 0 END), 0) FROM (
SELECT u.bk AS bk, SUM(u.l) AS l_cnt, SUM(u.r) AS r_cnt FROM (
SELECT bk, 1 AS l, 0 AS r FROM ({leftSql}) la
UNION ALL
SELECT bk, 0 AS l, 1 AS r FROM ({rightSql}) ra
) u GROUP BY u.bk
) g
""";
/// 集合内是否有重复 / 空业务键:返回 (总数, 去重数, 空值数)。
public static string BuildBizKeyQualitySql(string setSql) =>
$"""
SELECT COUNT(*) AS total, COUNT(DISTINCT x.bk) AS distinct_cnt,
SUM(CASE WHEN x.bk IS NULL OR x.bk = '' THEN 1 ELSE 0 END) AS blank_cnt
FROM ({setSql}) x
""";
/// NULL 归一(checksum 用)。
private static string Norm(string expr, S0DimValueKind kind) => kind switch
{
S0DimValueKind.DateTimeSec => $"IFNULL(DATE_FORMAT({expr}, '%Y-%m-%d %H:%i:%s'), {NullSentinel})",
S0DimValueKind.BoolTrue => $"IFNULL(CAST(CAST({expr} AS UNSIGNED) AS CHAR), {NullSentinel})",
_ => $"IFNULL(CAST({expr} AS CHAR), {NullSentinel})"
};
private static IReadOnlyList ChecksumColumns(S0DimDefinition def) =>
def.Columns.Where(c => c.Kind != S0DimValueKind.TenantIdColumn).ToList();
///
/// 单行哈希:MD5 前 15 个十六进制位 → 60 bit 无符号整数。
///
/// 🔴 CAST(... AS UNSIGNED) **不可省略**:CONV() 返回的是**字符串**,
/// 而 MySQL 的 SUM(字符串) 会按 DOUBLE 累加 —— 只有约 16 位有效数字,
/// 19 位的和末几位不可靠,且**两侧扫描顺序不同会产生不同舍入**,
/// 于是完全相同的数据也会算出不同校验和(2026-09-07 Batch 2 实测:
/// 同一租户 source=…980000 / dim=…980700,count 相同却误报不一致)。
/// 转成整数后 SUM 走 DECIMAL 精确累加,与顺序无关。
///
private static string RowHash(string rowStringExpr) =>
$"CAST(CONV(SUBSTRING(MD5({rowStringExpr}),1,15),16,10) AS UNSIGNED)";
///
/// 属性校验和:SUM()。
/// 用 SUM 而非 GROUP_CONCAT/BIT_XOR:交换律使**排序与结果无关**(无需 ORDER BY,也不受 group_concat_max_len 截断);
/// XOR 会让成对重复相互抵消从而掩盖重复。必须与 COUNT 成对使用。
///
public static string BuildSourceChecksumSql(S0DimDefinition def)
{
var parts = ChecksumColumns(def).Select(c => Norm($"{SrcAlias}.{Q(c.JsonPath!)}", c.Kind));
var rowStr = $"CONCAT_WS({FieldSep}, {Norm(SourceBizKeyExpr(def), S0DimValueKind.Str)}, {string.Join(", ", parts)})";
return $"SELECT COUNT(*) AS cnt, CAST(IFNULL(SUM({RowHash(rowStr)}),0) AS DECIMAL(40,0)) AS chk " +
$"FROM {Q(def.SourceTable)} {SrcAlias} WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
}
///
public static string BuildDimChecksumSql(S0DimDefinition def)
{
var parts = ChecksumColumns(def).Select(c => Norm($"{DimAlias}.{Q(c.TargetColumn)}", c.Kind));
var rowStr = $"CONCAT_WS({FieldSep}, {Norm(DimBizKeyExpr(def), S0DimValueKind.Str)}, {string.Join(", ", parts)})";
return $"SELECT COUNT(*) AS cnt, CAST(IFNULL(SUM({RowHash(rowStr)}),0) AS DECIMAL(40,0)) AS chk " +
$"FROM {Q(def.DimTable)} {DimAlias} WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
}
/// dim 业务键重复计数(唯一键之外的兜底断言)。参数:@TenantId
public static string BuildDimDuplicateBizKeySql(S0DimDefinition def) =>
$"""
SELECT COUNT(*) FROM (
SELECT {DimBizKeyExpr(def)} AS bk
FROM {Q(def.DimTable)} {DimAlias}
WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId
GROUP BY bk HAVING COUNT(*) > 1
) z
""";
/// 镜像唯一性断言(如 dim_location 的 (tenant_id, location_code))。参数:@TenantId
public static string BuildMirrorUniqueViolationSql(S0DimDefinition def)
{
if (def.MirrorUniqueColumns is not { Count: > 0 }) throw new InvalidOperationException($"[{def.Key}] 未声明 MirrorUniqueColumns");
var cols = string.Join(", ", def.MirrorUniqueColumns.Select(c => $"{DimAlias}.{Q(c)}"));
return $"""
SELECT COUNT(*) FROM (
SELECT {cols}
FROM {Q(def.DimTable)} {DimAlias}
WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId
GROUP BY {cols} HAVING COUNT(*) > 1
) z
""";
}
/// 孤儿子行(父维度中找不到父行)。参数:@TenantId
///
/// 孤儿判定的公共 FROM/WHERE —— 样本查询与计数查询**共用同一份**,避免两处口径漂移。
///
/// 「未建立关系」不是孤儿:父键列为 NULL 表示该行压根没声明父(如员工未分配部门),
/// 与「声明了父但父不存在」(真悬挂)是两种语义。<=> 是 NULL 安全等值,
/// 而父维度的业务键列均为 NOT NULL,故不加 IS NOT NULL 过滤会把全部未分配行
/// 误报成孤儿,淹没真正的悬挂引用。对父键全部 Required 的维度(如 LocationShelf)本条件恒真,行为不变。
///
private static string OrphanFromWhere(S0DimDefinition def)
{
var on = string.Join(" AND ", def.ParentKeyColumns!.Select(c => $"p.{Q(c)} <=> {DimAlias}.{Q(c)}"));
var declared = string.Join(" AND ", def.ParentKeyColumns!.Select(c => $"{DimAlias}.{Q(c)} IS NOT NULL"));
return $"""
FROM {Q(def.DimTable)} {DimAlias}
LEFT JOIN {Q(def.ParentDimTable!)} p
ON p.{Q(S0DimDefinition.TenantColumn)} = {DimAlias}.{Q(S0DimDefinition.TenantColumn)} AND {on}
WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId AND {declared} AND p.`id` IS NULL
""";
}
/// 孤儿子行**样本**(受 limit 截断)。参数:@TenantId
public static string BuildOrphanChildSql(S0DimDefinition def, int limit)
{
AssertParentDeclared(def);
var sel = string.Join(", ", def.BusinessKeyColumns.Skip(1).Select(c => $"{DimAlias}.{Q(c)}"));
return $"SELECT {sel} {OrphanFromWhere(def)} LIMIT {limit}";
}
/// 孤儿子行**真实总数**(不截断)。参数:@TenantId
public static string BuildOrphanChildCountSql(S0DimDefinition def)
{
AssertParentDeclared(def);
return $"SELECT COUNT(*) {OrphanFromWhere(def)}";
}
private static void AssertParentDeclared(S0DimDefinition def)
{
if (def.ParentDimTable is null || def.ParentKeyColumns is not { Count: > 0 })
throw new InvalidOperationException($"[{def.Key}] 未声明父维度");
}
// ════════════════════════════════════════════════════════════════════════════════════
// 同库集合式贴源(source → staging)
//
// 背景:共享的 MdpStagingWriter 对**每一行**发一条 INSERT ... ON DUPLICATE KEY UPDATE
// (MdpStagingWriter.cs 的 ExecuteCommandAsync),是为跨源(T8 SQLServer / dopdemorq /
// Excel / API)设计的通用最小公倍数。当源与目标恰好是同一个物理库时,这些往返是纯浪费:
// 实测 aidopdev 上约 65ms/行、稳定 ~15 行/s,ItemMaster 单租户 16150 行需 ≈20 分钟。
//
// 本段只**新增**一条同库快路径,不改 MdpStagingWriter —— 它被 187 个 entity 共用
// (S0=20 / S1=34 / S2=16 / S3=44 / S4=8 / S5=29 / S6=15 / S7=12 / T8=9)。
// 不满足同库判定时由调用方原样回退旧路径。
// ════════════════════════════════════════════════════════════════════════════════════
///
/// source_row_id 的候选列,**顺序必须与 MdpDbPullExecutor.ResolveSourceRowId 逐字一致**,
/// 否则同一源行在两条路径下会得到不同的 source_row_id,贴源唯一键随之漂移。
///
public static readonly IReadOnlyList SourceRowIdCandidates =
[
"RecID", "RecId", "recid", "Recid",
"id", "Id", "ID",
"noid", "billno", "BillNo", "djbh", "Djbh"
];
///
/// 快路径支持的 DATA_TYPE。**白名单而非黑名单**:出现任何未列出的类型
/// (binary / blob / json / enum / set / time / geometry …)
/// 一律判定不适用并回退旧路径,绝不猜测其序列化形态 ——
/// 例如 blob 在旧路径是 byte[] → base64 字符串,而 JSON_OBJECT 会直出原始字节,
/// 两者不等价且不会报错,属静默失真。
/// (实测:S0 全部 19 张源表只用到 varchar/bigint/datetime/tinyint/decimal/int/bit 七类,全部在列。)
///
public static readonly IReadOnlySet SameDbSupportedDataTypes =
new HashSet(StringComparer.OrdinalIgnoreCase)
{
"tinyint", "smallint", "mediumint", "int", "integer", "bigint", "year",
"decimal", "numeric", "float", "double",
"datetime", "timestamp", "date",
"char", "varchar", "tinytext", "text", "mediumtext", "longtext",
"bit"
};
private static readonly IReadOnlySet IntegerDataTypes =
new HashSet(StringComparer.OrdinalIgnoreCase)
{ "tinyint", "smallint", "mediumint", "int", "integer", "bigint", "year" };
///
/// 单列在 JSON_OBJECT 中的表达式,产出的 JSON_TYPE 必须与旧
/// JsonSerializer.Serialize(dict, RawDataJsonOptions) **逐类型一致**。
/// 下表每一行都在 aidopdev 上用同一批真实行对拍过(旧 raw_data vs 本表达式):
///
///
/// - tinyint(1)(恰好)BOOLEAN true/false
/// - tinyint(1) unsigned zerofillINTEGER 0(TreatTinyAsBoolean 不生效)
/// - bit(1)INTEGER 0
/// - decimal(18,5)DOUBLE 0.0(**不是** DECIMAL 0.00000)
/// - datetimeSTRING 2023-12-18 13:57:30.000000(**不是** JSON DATETIME)
/// - 字符类STRING 直出
///
/// SQL NULL 在各分支下都产出 JSON null,与旧路径的 C# null 一致。
///
private static string JsonValueExpr(S0DimSourceColumn c)
{
var col = $"{SrcAlias}.{Q(c.Name)}";
if (string.Equals(c.ColumnType, "tinyint(1)", StringComparison.OrdinalIgnoreCase))
return $"CAST(IF({col} IS NULL, NULL, IF({col} <> 0, 'true', 'false')) AS JSON)";
if (string.Equals(c.DataType, "bit", StringComparison.OrdinalIgnoreCase))
return $"CAST({col} AS UNSIGNED)";
if (c.DataType is "decimal" or "numeric" or "float" or "double"
|| string.Equals(c.DataType, "decimal", StringComparison.OrdinalIgnoreCase)
|| string.Equals(c.DataType, "numeric", StringComparison.OrdinalIgnoreCase)
|| string.Equals(c.DataType, "float", StringComparison.OrdinalIgnoreCase)
|| string.Equals(c.DataType, "double", StringComparison.OrdinalIgnoreCase))
return $"CAST({col} AS DOUBLE)";
if (string.Equals(c.DataType, "datetime", StringComparison.OrdinalIgnoreCase)
|| string.Equals(c.DataType, "timestamp", StringComparison.OrdinalIgnoreCase)
|| string.Equals(c.DataType, "date", StringComparison.OrdinalIgnoreCase))
return $"DATE_FORMAT({col}, '%Y-%m-%d %H:%i:%s.%f')";
// 整数类与字符类都由 JSON_OBJECT 直出(INTEGER / STRING),与旧路径一致
return col;
}
///
/// 单个候选列的「非空且非纯空白」取值,语义对齐旧路径的
/// string.IsNullOrWhiteSpace(v.ToString()) 跳过逻辑。
/// 注意旧路径返回的是**未 TRIM 的原值**,故此处只用 TRIM 判空、取值仍取原值。
///
private static string NonBlankChar(string column) =>
$"IF({SrcAlias}.{Q(column)} IS NULL OR TRIM(CAST({SrcAlias}.{Q(column)} AS CHAR)) = '', " +
$"NULL, CAST({SrcAlias}.{Q(column)} AS CHAR))";
///
/// source_row_id 表达式:按候选顺序取第一个非空白值。
/// 旧路径在全部候选都空白时兜底 Guid.NewGuid();本表达式**不生成 UUID()** ——
/// 调用方须先确保存在一个 NOT NULL 的整数候选列(见 的前置条件),
/// 否则判定不适用、回退旧路径。这样既避免了在 INSERT ... SELECT 里引入非确定函数,
/// 也保证两条路径对同一行给出同一个 source_row_id。
///
private static string RowIdExpr(IReadOnlyList presentCandidates) =>
presentCandidates.Count == 1
? NonBlankChar(presentCandidates[0])
: $"COALESCE({string.Join(", ", presentCandidates.Select(NonBlankChar))})";
///
/// source_biz_key 表达式,语义对齐 MdpStagingWriter.BuildBizKey:
/// **任一分量为 NULL 或空串 → 整个业务键作废,回落 source_row_id**。
///
/// 🔴 这里不能直接写 CONCAT_WS('#', …) —— CONCAT_WS 会**跳过** NULL 分量,
/// 于是 (NULL, 'B') 会得到 'B',而旧路径给的是 source_row_id。
/// 两条路径由此产出不同的 source_biz_key,且四向集合比对会在**下一次**刷新才炸,
/// 极难定位。
///
private static string BizKeyExpr(S0DimDefinition def, string rowIdExpr)
{
var guard = string.Join(" AND ", def.SourceBizKeyColumns.Select(
c => $"({SrcAlias}.{Q(c)} IS NOT NULL AND CAST({SrcAlias}.{Q(c)} AS CHAR) <> '')"));
var parts = string.Join(", ", def.SourceBizKeyColumns.Select(
c => $"CAST({SrcAlias}.{Q(c)} AS CHAR)"));
return $"CASE WHEN {guard} THEN CONCAT_WS('{BizKeySeparator}', {parts}) ELSE {rowIdExpr} END";
}
///
/// 同库集合式贴源语句:INSERT INTO staging (...) SELECT ... FROM source WHERE tenant_id = @TenantId。
/// 一条语句完成整租户装载,DB 往返 O(1),取代旧路径的 O(N)。
///
/// 租户硬门:WHERE 结构性包含 {SrcAlias}.tenant_id = @TenantId,
/// 源行无有效租户即不进结果集 —— 与 RequireMatchingSourceTenant=true 的 fail-closed 语义一致,
/// **绝不拿 ctx 兜底、绝不落 0**。
///
/// 为什么保留 ON DUPLICATE KEY UPDATE:现有 FULL 契约确实保证
/// S0DimMaterializer 在 purge(三段谓词,本租户+本源系统+本源表)之后才装载,
/// 单纯 INSERT 本不会撞键;但旧 MdpStagingWriter 是 upsert,
/// 源侧出现重复 source_row_id 时它 last-wins 而不报错。保留 upsert 才能让两条路径
/// 在**异常数据**下也同构,不因换路径而改变 FULL 语义。
///
/// 维度定义。
/// 源表全部列(决定 raw_data 的 JSON 形态)。
/// staging 表实际存在的列(factory_id / company_id 按存在与否可选)。
/// 源表中实际存在的 source_row_id 候选列,按候选顺序。
public static string BuildSameDbStagingInsertSql(
S0DimDefinition def,
IReadOnlyList sourceColumns,
IReadOnlySet stagingColumns,
IReadOnlyList rowIdCandidates)
{
if (sourceColumns.Count == 0)
throw new InvalidOperationException($"[{def.Key}] 源表无列,无法生成同库贴源语句");
if (rowIdCandidates.Count == 0)
throw new InvalidOperationException($"[{def.Key}] 源表无 source_row_id 候选列");
var rowId = RowIdExpr(rowIdCandidates);
var json = "JSON_OBJECT(" + string.Join(", ", sourceColumns.Select(
c => $"'{c.Name}', {JsonValueExpr(c)}")) + ")";
var cols = new List { S0DimDefinition.TenantColumn };
var vals = new List { $"{SrcAlias}.{Q(S0DimDefinition.TenantColumn)}" };
// 与旧 MdpStagingWriter 对齐:staging 有该列才写,且 factory 走 NormalizeFactoryId 口径
// (>0 且 != tenant 才保留,否则 NULL)。ItemMaster 无 factory_id/FactoryId 列 →
// 表达式恒 NULL,与旧路径实测的 16150/16150 全 NULL 一致。
if (stagingColumns.Contains("factory_id"))
{
var f = SourceFactoryExpr(sourceColumns);
cols.Add("factory_id");
vals.Add(f is null
? "NULL"
: $"IF({f} > 0 AND {f} <> {SrcAlias}.{Q(S0DimDefinition.TenantColumn)}, {f}, NULL)");
}
if (stagingColumns.Contains("company_id"))
{
var c = FindColumn(sourceColumns, "company_id");
cols.Add("company_id");
vals.Add(c is null ? "NULL" : $"{SrcAlias}.{Q(c)}");
}
cols.AddRange(["source_system", "source_table", "source_row_id", "source_biz_key",
"raw_data", "sync_batch_id", "sync_time", "process_status", "create_time"]);
vals.AddRange(["@SourceSystem", "@SourceTable", rowId, BizKeyExpr(def, rowId),
json, "@BatchId", "@Now", "'PENDING'", "@Now"]);
var updates = new List
{
"`source_row_id`=VALUES(`source_row_id`)",
"`raw_data`=VALUES(`raw_data`)",
"`sync_batch_id`=VALUES(`sync_batch_id`)",
"`sync_time`=VALUES(`sync_time`)",
"`process_status`='PENDING'",
$"`{S0DimDefinition.TenantColumn}`=VALUES(`{S0DimDefinition.TenantColumn}`)"
};
if (stagingColumns.Contains("factory_id")) updates.Add("`factory_id`=VALUES(`factory_id`)");
if (stagingColumns.Contains("company_id")) updates.Add("`company_id`=VALUES(`company_id`)");
updates.Add("`update_time`=@Now");
return $"""
INSERT INTO {Q(def.StagingTable)}
({string.Join(", ", cols.Select(Q))})
SELECT
{string.Join(",\n ", vals)}
FROM {Q(def.SourceTable)} {SrcAlias}
WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId
ON DUPLICATE KEY UPDATE
{string.Join(",\n ", updates)}
""";
}
/// 源侧工厂列:对齐旧路径的 factory_id ?? FactoryId 顺序;都不存在返回 null。
private static string? SourceFactoryExpr(IReadOnlyList cols)
{
var hit = FindColumn(cols, "factory_id") ?? FindColumn(cols, "FactoryId");
return hit is null ? null : $"{SrcAlias}.{Q(hit)}";
}
private static string? FindColumn(IReadOnlyList cols, string name) =>
cols.FirstOrDefault(c => string.Equals(c.Name, name, StringComparison.OrdinalIgnoreCase))?.Name;
}