S0DimSqlBuilder.cs 31 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558
  1. using System.Text;
  2. namespace Admin.NET.Plugin.AiDOP.DataPlatform.S0Dim;
  3. /// <summary>
  4. /// 由 <see cref="S0DimDefinition"/> 生成 SQL 的**唯一**位置(纯函数,不接触 DB,可单测)。
  5. ///
  6. /// 生成的所有语句都必须能逐条证明租户范围:
  7. /// staging 侧固定四段谓词 <c>tenant_id / source_system / source_table / sync_batch_id</c>;
  8. /// dim 侧固定 <c>tenant_id</c>;source 侧固定 <c>tenant_id</c>。
  9. /// </summary>
  10. public static class S0DimSqlBuilder
  11. {
  12. /// <summary>staging 别名。</summary>
  13. public const string StgAlias = "s";
  14. /// <summary>dim 别名。</summary>
  15. public const string DimAlias = "d";
  16. /// <summary>source 别名。</summary>
  17. public const string SrcAlias = "a";
  18. /// <summary>业务键各段之间的分隔符,与 <c>MdpStagingWriter.BuildBizKey</c> 的 <c>string.Join("#")</c> 一致。</summary>
  19. public const string BizKeySeparator = "#";
  20. /// <summary>checksum 行内字段分隔(US,业务值不可能包含)。</summary>
  21. private const string FieldSep = "CHAR(31 USING utf8mb4)";
  22. /// <summary>checksum 的 NULL 哨兵(RS),与空串可区分。</summary>
  23. private const string NullSentinel = "CHAR(30 USING utf8mb4)";
  24. private static string Q(string ident)
  25. {
  26. if (!S0DimDefinition.IdentifierRe.IsMatch(ident))
  27. throw new InvalidOperationException($"非法标识符:{ident}");
  28. return $"`{ident}`";
  29. }
  30. /// <summary>列取值表达式(作用于 staging 的 raw_data;TenantIdColumn 例外)。</summary>
  31. private static string ValueExpr(S0DimColumn c) => c.Kind switch
  32. {
  33. S0DimValueKind.TenantIdColumn => $"{StgAlias}.{Q(S0DimDefinition.TenantColumn)}",
  34. S0DimValueKind.Str => MdpJsonSql.Str(StgAlias, c.JsonPath!),
  35. S0DimValueKind.Int => MdpJsonSql.Int(StgAlias, c.JsonPath!),
  36. S0DimValueKind.Dec => MdpJsonSql.Dec(StgAlias, c.JsonPath!, c.Precision, c.Scale),
  37. S0DimValueKind.DateTimeSec => MdpJsonSql.DateTimeSec(StgAlias, c.JsonPath!),
  38. S0DimValueKind.BoolTrue => MdpJsonSql.BoolTrue(StgAlias, c.JsonPath!),
  39. _ => throw new InvalidOperationException($"未支持的取值方式:{c.Kind}")
  40. };
  41. /// <summary>
  42. /// staging 四段固定谓词。<paramref name="withBatch"/> 为 false 时用于「不限本批」的只读对账。
  43. /// </summary>
  44. public static string StagingFilter(bool withBatch = true)
  45. {
  46. var sb = new StringBuilder();
  47. sb.Append($"{StgAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId");
  48. sb.Append($" AND {StgAlias}.`source_system` = @SourceSystem");
  49. sb.Append($" AND {StgAlias}.`source_table` = @SourceTable");
  50. if (withBatch) sb.Append($" AND {StgAlias}.`sync_batch_id` = @BatchId");
  51. return sb.ToString();
  52. }
  53. /// <summary>行级过滤:Required 列非空 + 布尔列可辨识(不满足的行不进 dim,随后被行数对账捕获)。</summary>
  54. private static IReadOnlyList<string> RowGuards(S0DimDefinition def)
  55. {
  56. var guards = new List<string>();
  57. foreach (var c in def.Columns)
  58. {
  59. if (c.Kind == S0DimValueKind.TenantIdColumn) continue;
  60. if (c.Required)
  61. {
  62. guards.Add($"({ValueExpr(c)}) IS NOT NULL");
  63. // IS NOT NULL 挡不住空串。业务键走 CONCAT_WS 拼接,空段会被**保留**,
  64. // 于是 Domain='' 会产出 '#LOC01' 这种「看起来合法」的键,而
  65. // AssertNoSourceDuplicateAsync 的 blank_cnt 只捕获整串全空的情况。
  66. // 当前 Catalog 中 Required 的 Str 列恰好就是全部业务键分量
  67. // (domain_code / work_center_code / department_code / location_code / shelf_code),
  68. // 故此守卫命中面 = 业务键,无附带影响。
  69. if (c.Kind == S0DimValueKind.Str)
  70. guards.Add($"({ValueExpr(c)}) <> ''");
  71. }
  72. if (c.Kind == S0DimValueKind.BoolTrue)
  73. {
  74. // MdpJsonSql.BoolTrue 对 null / 不可识别值静默返回 0;此处强制要求原始值可辨识,
  75. // 否则整行被过滤 → dim_count < stg_count → 对账 FAIL(暴露而非吞掉)
  76. var raw = MdpJsonSql.Raw(StgAlias, c.JsonPath!);
  77. guards.Add($"LOWER({raw}) IN ('1','0','true','false')");
  78. }
  79. }
  80. return guards;
  81. }
  82. /// <summary>
  83. /// dim 物化语句:<c>INSERT INTO dim (...) SELECT ... FROM staging</c>。
  84. /// **刻意不带 ON DUPLICATE KEY UPDATE** —— 源侧业务键重复时直接撞 dim 唯一键报错,
  85. /// 由外层事务回滚(transform fail),绝不 last-wins。
  86. /// 参数:@TenantId @SourceSystem @SourceTable @BatchId @Now
  87. /// </summary>
  88. public static string BuildInsertSql(S0DimDefinition def)
  89. {
  90. def.Validate();
  91. var cols = new List<string>();
  92. var vals = new List<string>();
  93. foreach (var c in def.Columns)
  94. {
  95. cols.Add(Q(c.TargetColumn));
  96. vals.Add(ValueExpr(c));
  97. }
  98. cols.AddRange(["`source_system`", "`source_biz_key`", "`sync_batch_id`", "`sync_time`"]);
  99. vals.AddRange([$"{StgAlias}.`source_system`", $"{StgAlias}.`source_biz_key`", "@BatchId", "@Now"]);
  100. var where = new List<string> { StagingFilter() };
  101. where.AddRange(RowGuards(def));
  102. return $"""
  103. INSERT INTO {Q(def.DimTable)}
  104. ({string.Join(", ", cols)})
  105. SELECT
  106. {string.Join(",\n ", vals)}
  107. FROM {Q(def.StagingTable)} {StgAlias}
  108. WHERE {string.Join("\n AND ", where)}
  109. """;
  110. }
  111. /// <summary>staging 分区清理(purge)。参数:@TenantId @SourceSystem @SourceTable</summary>
  112. public static string BuildPurgeStagingSql(S0DimDefinition def) =>
  113. $"""
  114. DELETE {StgAlias} FROM {Q(def.StagingTable)} {StgAlias}
  115. WHERE {StagingFilter(withBatch: false)}
  116. """;
  117. /// <summary>源表行数(按租户)。参数:@TenantId</summary>
  118. public static string BuildSourceCountSql(S0DimDefinition def) =>
  119. $"SELECT COUNT(*) FROM {Q(def.SourceTable)} {SrcAlias} " +
  120. $"WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
  121. /// <summary>源表中无法归属租户的行数(tenant_id IS NULL 或 &lt;= 0)—— 这些行整条链路都看不见。</summary>
  122. public static string BuildSourceUntenantedCountSql(S0DimDefinition def) =>
  123. $"SELECT COUNT(*) FROM {Q(def.SourceTable)} {SrcAlias} " +
  124. $"WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} IS NULL " +
  125. $"OR {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} <= 0";
  126. /// <summary>staging 行数。参数:@TenantId @SourceSystem @SourceTable [@BatchId]</summary>
  127. public static string BuildStagingCountSql(S0DimDefinition def, bool withBatch = true) =>
  128. $"SELECT COUNT(*) FROM {Q(def.StagingTable)} {StgAlias} WHERE {StagingFilter(withBatch)}";
  129. /// <summary>dim 行数(按租户)。参数:@TenantId</summary>
  130. public static string BuildDimCountSql(S0DimDefinition def) =>
  131. $"SELECT COUNT(*) FROM {Q(def.DimTable)} {DimAlias} " +
  132. $"WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
  133. /// <summary>dim 中不属于本批的残留行数(FULL Replace 之后必须为 0)。参数:@TenantId @BatchId</summary>
  134. public static string BuildDimBatchImpuritySql(S0DimDefinition def) =>
  135. $"SELECT COUNT(*) FROM {Q(def.DimTable)} {DimAlias} " +
  136. $"WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId " +
  137. $"AND {DimAlias}.`sync_batch_id` <> @BatchId";
  138. /// <summary>源侧业务键表达式(物理列,与 mdp_entity.biz_key_expr 同序同分隔)。</summary>
  139. public static string SourceBizKeyExpr(S0DimDefinition def) =>
  140. $"CONCAT_WS('{BizKeySeparator}', {string.Join(", ", def.SourceBizKeyColumns.Select(c => $"{SrcAlias}.{Q(c)}"))})";
  141. /// <summary>dim 侧业务键表达式(去掉 tenant_id 后的业务键列)。</summary>
  142. public static string DimBizKeyExpr(S0DimDefinition def) =>
  143. $"CONCAT_WS('{BizKeySeparator}', {string.Join(", ", def.BusinessKeyColumns.Skip(1).Select(c => $"{DimAlias}.{Q(c)}"))})";
  144. /// <summary>三层业务键集合来源子查询。</summary>
  145. public static string SourceBizKeySetSql(S0DimDefinition def) =>
  146. $"SELECT {SourceBizKeyExpr(def)} AS bk FROM {Q(def.SourceTable)} {SrcAlias} " +
  147. $"WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
  148. /// <inheritdoc cref="SourceBizKeySetSql"/>
  149. public static string StagingBizKeySetSql(S0DimDefinition def, bool withBatch = true) =>
  150. $"SELECT {StgAlias}.`source_biz_key` AS bk FROM {Q(def.StagingTable)} {StgAlias} WHERE {StagingFilter(withBatch)}";
  151. /// <inheritdoc cref="SourceBizKeySetSql"/>
  152. public static string DimBizKeySetSql(S0DimDefinition def) =>
  153. $"SELECT {DimBizKeyExpr(def)} AS bk FROM {Q(def.DimTable)} {DimAlias} " +
  154. $"WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
  155. /// <summary>
  156. /// 集合差集计数:<paramref name="leftSql"/> \ <paramref name="rightSql"/>,
  157. /// 返回**左侧无匹配的行数**(不是去重后的键数)。
  158. ///
  159. /// <para>🔴 <b>不要写成 <c>LEFT JOIN … ON y.bk &lt;=&gt; x.bk WHERE y.bk IS NULL</c></b>(2026-09-07 实测):
  160. /// 四个方向里总有一侧的 <c>bk</c> 是 <c>CONCAT_WS(...)</c> 计算表达式,**既不能建索引也不能哈希**。
  161. /// 当它落在 JOIN 右侧时,优化器只能退化成 nested loop —— 16150×16150 行的
  162. /// <c>stg \ src</c> 方向实测 <b>130.3s</b>(EXPLAIN: <c>rows=32,762,720</c>),
  163. /// 直接撞穿 <c>SqlSugarSetup.CommandTimeOut = 30</c>,而对称的 <c>src \ stg</c> 方向因右侧走
  164. /// <c>idx_stg_*_biz</c> 只要 0.5s —— <b>同一个缺陷只在一半方向上显形</b>,小表永远发现不了。
  165. /// 把 <c>&lt;=&gt;</c> 换成 <c>=</c> 无用(实测 127.7s),<c>NOT EXISTS</c> 也只到 69.2s。</para>
  166. ///
  167. /// <para>本写法把两侧各扫一遍后按 <c>bk</c> 聚合,复杂度 O(N+M),与方向无关:
  168. /// 同一 <c>stg \ src</c> 实测 <b>0.5s</b>(260×),EXPLAIN 里 nested loop 消失。</para>
  169. ///
  170. /// <para>语义与旧写法逐条对齐:① 未匹配的左行在 LEFT JOIN 下恰好产出 1 行,
  171. /// 故 <c>SUM(l_cnt)</c> 与 <c>COUNT(*)</c> 等价,**左右两侧的重复行都不会改变结果**;
  172. /// ② <c>GROUP BY</c> 把所有 NULL 归为同一组,天然等价于 <c>&lt;=&gt;</c> 的 NULL 相等语义。</para>
  173. /// </summary>
  174. public static string BuildAntiJoinCountSql(string leftSql, string rightSql) =>
  175. $"""
  176. SELECT COALESCE(SUM(CASE WHEN g.r_cnt = 0 THEN g.l_cnt ELSE 0 END), 0) FROM (
  177. SELECT u.bk AS bk, SUM(u.l) AS l_cnt, SUM(u.r) AS r_cnt FROM (
  178. SELECT bk, 1 AS l, 0 AS r FROM ({leftSql}) la
  179. UNION ALL
  180. SELECT bk, 0 AS l, 1 AS r FROM ({rightSql}) ra
  181. ) u GROUP BY u.bk
  182. ) g
  183. """;
  184. /// <summary>集合内是否有重复 / 空业务键:返回 (总数, 去重数, 空值数)。</summary>
  185. public static string BuildBizKeyQualitySql(string setSql) =>
  186. $"""
  187. SELECT COUNT(*) AS total, COUNT(DISTINCT x.bk) AS distinct_cnt,
  188. SUM(CASE WHEN x.bk IS NULL OR x.bk = '' THEN 1 ELSE 0 END) AS blank_cnt
  189. FROM ({setSql}) x
  190. """;
  191. /// <summary>NULL 归一(checksum 用)。</summary>
  192. private static string Norm(string expr, S0DimValueKind kind) => kind switch
  193. {
  194. S0DimValueKind.DateTimeSec => $"IFNULL(DATE_FORMAT({expr}, '%Y-%m-%d %H:%i:%s'), {NullSentinel})",
  195. S0DimValueKind.BoolTrue => $"IFNULL(CAST(CAST({expr} AS UNSIGNED) AS CHAR), {NullSentinel})",
  196. _ => $"IFNULL(CAST({expr} AS CHAR), {NullSentinel})"
  197. };
  198. private static IReadOnlyList<S0DimColumn> ChecksumColumns(S0DimDefinition def) =>
  199. def.Columns.Where(c => c.Kind != S0DimValueKind.TenantIdColumn).ToList();
  200. /// <summary>
  201. /// 单行哈希:MD5 前 15 个十六进制位 → 60 bit 无符号整数。
  202. ///
  203. /// 🔴 <c>CAST(... AS UNSIGNED)</c> **不可省略**:<c>CONV()</c> 返回的是**字符串**,
  204. /// 而 MySQL 的 <c>SUM(字符串)</c> 会按 DOUBLE 累加 —— 只有约 16 位有效数字,
  205. /// 19 位的和末几位不可靠,且**两侧扫描顺序不同会产生不同舍入**,
  206. /// 于是完全相同的数据也会算出不同校验和(2026-09-07 Batch 2 实测:
  207. /// 同一租户 source=…980000 / dim=…980700,count 相同却误报不一致)。
  208. /// 转成整数后 <c>SUM</c> 走 DECIMAL 精确累加,与顺序无关。
  209. /// </summary>
  210. private static string RowHash(string rowStringExpr) =>
  211. $"CAST(CONV(SUBSTRING(MD5({rowStringExpr}),1,15),16,10) AS UNSIGNED)";
  212. /// <summary>
  213. /// 属性校验和:<c>SUM(<see cref="RowHash"/>)</c>。
  214. /// 用 SUM 而非 GROUP_CONCAT/BIT_XOR:交换律使**排序与结果无关**(无需 ORDER BY,也不受 group_concat_max_len 截断);
  215. /// XOR 会让成对重复相互抵消从而掩盖重复。必须与 COUNT 成对使用。
  216. /// </summary>
  217. public static string BuildSourceChecksumSql(S0DimDefinition def)
  218. {
  219. var parts = ChecksumColumns(def).Select(c => Norm($"{SrcAlias}.{Q(c.JsonPath!)}", c.Kind));
  220. var rowStr = $"CONCAT_WS({FieldSep}, {Norm(SourceBizKeyExpr(def), S0DimValueKind.Str)}, {string.Join(", ", parts)})";
  221. return $"SELECT COUNT(*) AS cnt, CAST(IFNULL(SUM({RowHash(rowStr)}),0) AS DECIMAL(40,0)) AS chk " +
  222. $"FROM {Q(def.SourceTable)} {SrcAlias} WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
  223. }
  224. /// <inheritdoc cref="BuildSourceChecksumSql"/>
  225. public static string BuildDimChecksumSql(S0DimDefinition def)
  226. {
  227. var parts = ChecksumColumns(def).Select(c => Norm($"{DimAlias}.{Q(c.TargetColumn)}", c.Kind));
  228. var rowStr = $"CONCAT_WS({FieldSep}, {Norm(DimBizKeyExpr(def), S0DimValueKind.Str)}, {string.Join(", ", parts)})";
  229. return $"SELECT COUNT(*) AS cnt, CAST(IFNULL(SUM({RowHash(rowStr)}),0) AS DECIMAL(40,0)) AS chk " +
  230. $"FROM {Q(def.DimTable)} {DimAlias} WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId";
  231. }
  232. /// <summary>dim 业务键重复计数(唯一键之外的兜底断言)。参数:@TenantId</summary>
  233. public static string BuildDimDuplicateBizKeySql(S0DimDefinition def) =>
  234. $"""
  235. SELECT COUNT(*) FROM (
  236. SELECT {DimBizKeyExpr(def)} AS bk
  237. FROM {Q(def.DimTable)} {DimAlias}
  238. WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId
  239. GROUP BY bk HAVING COUNT(*) > 1
  240. ) z
  241. """;
  242. /// <summary>镜像唯一性断言(如 dim_location 的 (tenant_id, location_code))。参数:@TenantId</summary>
  243. public static string BuildMirrorUniqueViolationSql(S0DimDefinition def)
  244. {
  245. if (def.MirrorUniqueColumns is not { Count: > 0 }) throw new InvalidOperationException($"[{def.Key}] 未声明 MirrorUniqueColumns");
  246. var cols = string.Join(", ", def.MirrorUniqueColumns.Select(c => $"{DimAlias}.{Q(c)}"));
  247. return $"""
  248. SELECT COUNT(*) FROM (
  249. SELECT {cols}
  250. FROM {Q(def.DimTable)} {DimAlias}
  251. WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId
  252. GROUP BY {cols} HAVING COUNT(*) > 1
  253. ) z
  254. """;
  255. }
  256. /// <summary>孤儿子行(父维度中找不到父行)。参数:@TenantId</summary>
  257. /// <summary>
  258. /// 孤儿判定的公共 FROM/WHERE —— 样本查询与计数查询**共用同一份**,避免两处口径漂移。
  259. ///
  260. /// 「未建立关系」不是孤儿:父键列为 NULL 表示该行压根没声明父(如员工未分配部门),
  261. /// 与「声明了父但父不存在」(真悬挂)是两种语义。<c>&lt;=&gt;</c> 是 NULL 安全等值,
  262. /// 而父维度的业务键列均为 NOT NULL,故不加 <c>IS NOT NULL</c> 过滤会把全部未分配行
  263. /// 误报成孤儿,淹没真正的悬挂引用。对父键全部 Required 的维度(如 LocationShelf)本条件恒真,行为不变。
  264. /// </summary>
  265. private static string OrphanFromWhere(S0DimDefinition def)
  266. {
  267. var on = string.Join(" AND ", def.ParentKeyColumns!.Select(c => $"p.{Q(c)} <=> {DimAlias}.{Q(c)}"));
  268. var declared = string.Join(" AND ", def.ParentKeyColumns!.Select(c => $"{DimAlias}.{Q(c)} IS NOT NULL"));
  269. return $"""
  270. FROM {Q(def.DimTable)} {DimAlias}
  271. LEFT JOIN {Q(def.ParentDimTable!)} p
  272. ON p.{Q(S0DimDefinition.TenantColumn)} = {DimAlias}.{Q(S0DimDefinition.TenantColumn)} AND {on}
  273. WHERE {DimAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId AND {declared} AND p.`id` IS NULL
  274. """;
  275. }
  276. /// <summary>孤儿子行**样本**(受 limit 截断)。参数:@TenantId</summary>
  277. public static string BuildOrphanChildSql(S0DimDefinition def, int limit)
  278. {
  279. AssertParentDeclared(def);
  280. var sel = string.Join(", ", def.BusinessKeyColumns.Skip(1).Select(c => $"{DimAlias}.{Q(c)}"));
  281. return $"SELECT {sel} {OrphanFromWhere(def)} LIMIT {limit}";
  282. }
  283. /// <summary>孤儿子行**真实总数**(不截断)。参数:@TenantId</summary>
  284. public static string BuildOrphanChildCountSql(S0DimDefinition def)
  285. {
  286. AssertParentDeclared(def);
  287. return $"SELECT COUNT(*) {OrphanFromWhere(def)}";
  288. }
  289. private static void AssertParentDeclared(S0DimDefinition def)
  290. {
  291. if (def.ParentDimTable is null || def.ParentKeyColumns is not { Count: > 0 })
  292. throw new InvalidOperationException($"[{def.Key}] 未声明父维度");
  293. }
  294. // ════════════════════════════════════════════════════════════════════════════════════
  295. // 同库集合式贴源(source → staging)
  296. //
  297. // 背景:共享的 MdpStagingWriter 对**每一行**发一条 INSERT ... ON DUPLICATE KEY UPDATE
  298. // (MdpStagingWriter.cs 的 ExecuteCommandAsync),是为跨源(T8 SQLServer / dopdemorq /
  299. // Excel / API)设计的通用最小公倍数。当源与目标恰好是同一个物理库时,这些往返是纯浪费:
  300. // 实测 aidopdev 上约 65ms/行、稳定 ~15 行/s,ItemMaster 单租户 16150 行需 ≈20 分钟。
  301. //
  302. // 本段只**新增**一条同库快路径,不改 MdpStagingWriter —— 它被 187 个 entity 共用
  303. // (S0=20 / S1=34 / S2=16 / S3=44 / S4=8 / S5=29 / S6=15 / S7=12 / T8=9)。
  304. // 不满足同库判定时由调用方原样回退旧路径。
  305. // ════════════════════════════════════════════════════════════════════════════════════
  306. /// <summary>
  307. /// <c>source_row_id</c> 的候选列,**顺序必须与 <c>MdpDbPullExecutor.ResolveSourceRowId</c> 逐字一致**,
  308. /// 否则同一源行在两条路径下会得到不同的 <c>source_row_id</c>,贴源唯一键随之漂移。
  309. /// </summary>
  310. public static readonly IReadOnlyList<string> SourceRowIdCandidates =
  311. [
  312. "RecID", "RecId", "recid", "Recid",
  313. "id", "Id", "ID",
  314. "noid", "billno", "BillNo", "djbh", "Djbh"
  315. ];
  316. /// <summary>
  317. /// 快路径支持的 <c>DATA_TYPE</c>。**白名单而非黑名单**:出现任何未列出的类型
  318. /// (<c>binary</c> / <c>blob</c> / <c>json</c> / <c>enum</c> / <c>set</c> / <c>time</c> / <c>geometry</c> …)
  319. /// 一律判定不适用并回退旧路径,绝不猜测其序列化形态 ——
  320. /// 例如 <c>blob</c> 在旧路径是 <c>byte[]</c> → base64 字符串,而 <c>JSON_OBJECT</c> 会直出原始字节,
  321. /// 两者不等价且不会报错,属静默失真。
  322. /// (实测:S0 全部 19 张源表只用到 varchar/bigint/datetime/tinyint/decimal/int/bit 七类,全部在列。)
  323. /// </summary>
  324. public static readonly IReadOnlySet<string> SameDbSupportedDataTypes =
  325. new HashSet<string>(StringComparer.OrdinalIgnoreCase)
  326. {
  327. "tinyint", "smallint", "mediumint", "int", "integer", "bigint", "year",
  328. "decimal", "numeric", "float", "double",
  329. "datetime", "timestamp", "date",
  330. "char", "varchar", "tinytext", "text", "mediumtext", "longtext",
  331. "bit"
  332. };
  333. private static readonly IReadOnlySet<string> IntegerDataTypes =
  334. new HashSet<string>(StringComparer.OrdinalIgnoreCase)
  335. { "tinyint", "smallint", "mediumint", "int", "integer", "bigint", "year" };
  336. /// <summary>
  337. /// 单列在 <c>JSON_OBJECT</c> 中的表达式,产出的 <c>JSON_TYPE</c> 必须与旧
  338. /// <c>JsonSerializer.Serialize(dict, RawDataJsonOptions)</c> **逐类型一致**。
  339. /// 下表每一行都在 aidopdev 上用同一批真实行对拍过(旧 raw_data vs 本表达式):
  340. ///
  341. /// <list type="table">
  342. /// <item><term><c>tinyint(1)</c>(恰好)</term><description>BOOLEAN <c>true</c>/<c>false</c></description></item>
  343. /// <item><term><c>tinyint(1) unsigned zerofill</c></term><description>INTEGER <c>0</c>(TreatTinyAsBoolean 不生效)</description></item>
  344. /// <item><term><c>bit(1)</c></term><description>INTEGER <c>0</c></description></item>
  345. /// <item><term><c>decimal(18,5)</c></term><description>DOUBLE <c>0.0</c>(**不是** DECIMAL <c>0.00000</c>)</description></item>
  346. /// <item><term><c>datetime</c></term><description>STRING <c>2023-12-18 13:57:30.000000</c>(**不是** JSON DATETIME)</description></item>
  347. /// <item><term>字符类</term><description>STRING 直出</description></item>
  348. /// </list>
  349. /// SQL NULL 在各分支下都产出 JSON <c>null</c>,与旧路径的 C# <c>null</c> 一致。
  350. /// </summary>
  351. private static string JsonValueExpr(S0DimSourceColumn c)
  352. {
  353. var col = $"{SrcAlias}.{Q(c.Name)}";
  354. if (string.Equals(c.ColumnType, "tinyint(1)", StringComparison.OrdinalIgnoreCase))
  355. return $"CAST(IF({col} IS NULL, NULL, IF({col} <> 0, 'true', 'false')) AS JSON)";
  356. if (string.Equals(c.DataType, "bit", StringComparison.OrdinalIgnoreCase))
  357. return $"CAST({col} AS UNSIGNED)";
  358. if (c.DataType is "decimal" or "numeric" or "float" or "double"
  359. || string.Equals(c.DataType, "decimal", StringComparison.OrdinalIgnoreCase)
  360. || string.Equals(c.DataType, "numeric", StringComparison.OrdinalIgnoreCase)
  361. || string.Equals(c.DataType, "float", StringComparison.OrdinalIgnoreCase)
  362. || string.Equals(c.DataType, "double", StringComparison.OrdinalIgnoreCase))
  363. return $"CAST({col} AS DOUBLE)";
  364. if (string.Equals(c.DataType, "datetime", StringComparison.OrdinalIgnoreCase)
  365. || string.Equals(c.DataType, "timestamp", StringComparison.OrdinalIgnoreCase)
  366. || string.Equals(c.DataType, "date", StringComparison.OrdinalIgnoreCase))
  367. return $"DATE_FORMAT({col}, '%Y-%m-%d %H:%i:%s.%f')";
  368. // 整数类与字符类都由 JSON_OBJECT 直出(INTEGER / STRING),与旧路径一致
  369. return col;
  370. }
  371. /// <summary>
  372. /// 单个候选列的「非空且非纯空白」取值,语义对齐旧路径的
  373. /// <c>string.IsNullOrWhiteSpace(v.ToString())</c> 跳过逻辑。
  374. /// 注意旧路径返回的是**未 TRIM 的原值**,故此处只用 TRIM 判空、取值仍取原值。
  375. /// </summary>
  376. private static string NonBlankChar(string column) =>
  377. $"IF({SrcAlias}.{Q(column)} IS NULL OR TRIM(CAST({SrcAlias}.{Q(column)} AS CHAR)) = '', " +
  378. $"NULL, CAST({SrcAlias}.{Q(column)} AS CHAR))";
  379. /// <summary>
  380. /// <c>source_row_id</c> 表达式:按候选顺序取第一个非空白值。
  381. /// 旧路径在全部候选都空白时兜底 <c>Guid.NewGuid()</c>;本表达式**不生成 <c>UUID()</c>** ——
  382. /// 调用方须先确保存在一个 NOT NULL 的整数候选列(见 <see cref="BuildSameDbStagingInsertSql"/> 的前置条件),
  383. /// 否则判定不适用、回退旧路径。这样既避免了在 <c>INSERT ... SELECT</c> 里引入非确定函数,
  384. /// 也保证两条路径对同一行给出同一个 <c>source_row_id</c>。
  385. /// </summary>
  386. private static string RowIdExpr(IReadOnlyList<string> presentCandidates) =>
  387. presentCandidates.Count == 1
  388. ? NonBlankChar(presentCandidates[0])
  389. : $"COALESCE({string.Join(", ", presentCandidates.Select(NonBlankChar))})";
  390. /// <summary>
  391. /// <c>source_biz_key</c> 表达式,语义对齐 <c>MdpStagingWriter.BuildBizKey</c>:
  392. /// **任一分量为 NULL 或空串 → 整个业务键作废,回落 <c>source_row_id</c>**。
  393. ///
  394. /// <para>🔴 这里不能直接写 <c>CONCAT_WS('#', …)</c> —— <c>CONCAT_WS</c> 会**跳过** NULL 分量,
  395. /// 于是 <c>(NULL, 'B')</c> 会得到 <c>'B'</c>,而旧路径给的是 <c>source_row_id</c>。
  396. /// 两条路径由此产出不同的 <c>source_biz_key</c>,且四向集合比对会在**下一次**刷新才炸,
  397. /// 极难定位。</para>
  398. /// </summary>
  399. private static string BizKeyExpr(S0DimDefinition def, string rowIdExpr)
  400. {
  401. var guard = string.Join(" AND ", def.SourceBizKeyColumns.Select(
  402. c => $"({SrcAlias}.{Q(c)} IS NOT NULL AND CAST({SrcAlias}.{Q(c)} AS CHAR) <> '')"));
  403. var parts = string.Join(", ", def.SourceBizKeyColumns.Select(
  404. c => $"CAST({SrcAlias}.{Q(c)} AS CHAR)"));
  405. return $"CASE WHEN {guard} THEN CONCAT_WS('{BizKeySeparator}', {parts}) ELSE {rowIdExpr} END";
  406. }
  407. /// <summary>
  408. /// 同库集合式贴源语句:<c>INSERT INTO staging (...) SELECT ... FROM source WHERE tenant_id = @TenantId</c>。
  409. /// 一条语句完成整租户装载,DB 往返 <b>O(1)</b>,取代旧路径的 O(N)。
  410. ///
  411. /// <para><b>租户硬门</b>:<c>WHERE</c> 结构性包含 <c>{SrcAlias}.tenant_id = @TenantId</c>,
  412. /// 源行无有效租户即不进结果集 —— 与 <c>RequireMatchingSourceTenant=true</c> 的 fail-closed 语义一致,
  413. /// **绝不拿 ctx 兜底、绝不落 0**。</para>
  414. ///
  415. /// <para><b>为什么保留 <c>ON DUPLICATE KEY UPDATE</c></b>:现有 FULL 契约确实保证
  416. /// <c>S0DimMaterializer</c> 在 purge(三段谓词,本租户+本源系统+本源表)之后才装载,
  417. /// 单纯 <c>INSERT</c> 本不会撞键;但旧 <c>MdpStagingWriter</c> 是 upsert,
  418. /// 源侧出现重复 <c>source_row_id</c> 时它 last-wins 而不报错。保留 upsert 才能让两条路径
  419. /// 在**异常数据**下也同构,不因换路径而改变 FULL 语义。</para>
  420. /// </summary>
  421. /// <param name="def">维度定义。</param>
  422. /// <param name="sourceColumns">源表全部列(决定 <c>raw_data</c> 的 JSON 形态)。</param>
  423. /// <param name="stagingColumns">staging 表实际存在的列(<c>factory_id</c> / <c>company_id</c> 按存在与否可选)。</param>
  424. /// <param name="rowIdCandidates">源表中实际存在的 <c>source_row_id</c> 候选列,按候选顺序。</param>
  425. public static string BuildSameDbStagingInsertSql(
  426. S0DimDefinition def,
  427. IReadOnlyList<S0DimSourceColumn> sourceColumns,
  428. IReadOnlySet<string> stagingColumns,
  429. IReadOnlyList<string> rowIdCandidates)
  430. {
  431. if (sourceColumns.Count == 0)
  432. throw new InvalidOperationException($"[{def.Key}] 源表无列,无法生成同库贴源语句");
  433. if (rowIdCandidates.Count == 0)
  434. throw new InvalidOperationException($"[{def.Key}] 源表无 source_row_id 候选列");
  435. var rowId = RowIdExpr(rowIdCandidates);
  436. var json = "JSON_OBJECT(" + string.Join(", ", sourceColumns.Select(
  437. c => $"'{c.Name}', {JsonValueExpr(c)}")) + ")";
  438. var cols = new List<string> { S0DimDefinition.TenantColumn };
  439. var vals = new List<string> { $"{SrcAlias}.{Q(S0DimDefinition.TenantColumn)}" };
  440. // 与旧 MdpStagingWriter 对齐:staging 有该列才写,且 factory 走 NormalizeFactoryId 口径
  441. // (>0 且 != tenant 才保留,否则 NULL)。ItemMaster 无 factory_id/FactoryId 列 →
  442. // 表达式恒 NULL,与旧路径实测的 16150/16150 全 NULL 一致。
  443. if (stagingColumns.Contains("factory_id"))
  444. {
  445. var f = SourceFactoryExpr(sourceColumns);
  446. cols.Add("factory_id");
  447. vals.Add(f is null
  448. ? "NULL"
  449. : $"IF({f} > 0 AND {f} <> {SrcAlias}.{Q(S0DimDefinition.TenantColumn)}, {f}, NULL)");
  450. }
  451. if (stagingColumns.Contains("company_id"))
  452. {
  453. var c = FindColumn(sourceColumns, "company_id");
  454. cols.Add("company_id");
  455. vals.Add(c is null ? "NULL" : $"{SrcAlias}.{Q(c)}");
  456. }
  457. cols.AddRange(["source_system", "source_table", "source_row_id", "source_biz_key",
  458. "raw_data", "sync_batch_id", "sync_time", "process_status", "create_time"]);
  459. vals.AddRange(["@SourceSystem", "@SourceTable", rowId, BizKeyExpr(def, rowId),
  460. json, "@BatchId", "@Now", "'PENDING'", "@Now"]);
  461. var updates = new List<string>
  462. {
  463. "`source_row_id`=VALUES(`source_row_id`)",
  464. "`raw_data`=VALUES(`raw_data`)",
  465. "`sync_batch_id`=VALUES(`sync_batch_id`)",
  466. "`sync_time`=VALUES(`sync_time`)",
  467. "`process_status`='PENDING'",
  468. $"`{S0DimDefinition.TenantColumn}`=VALUES(`{S0DimDefinition.TenantColumn}`)"
  469. };
  470. if (stagingColumns.Contains("factory_id")) updates.Add("`factory_id`=VALUES(`factory_id`)");
  471. if (stagingColumns.Contains("company_id")) updates.Add("`company_id`=VALUES(`company_id`)");
  472. updates.Add("`update_time`=@Now");
  473. return $"""
  474. INSERT INTO {Q(def.StagingTable)}
  475. ({string.Join(", ", cols.Select(Q))})
  476. SELECT
  477. {string.Join(",\n ", vals)}
  478. FROM {Q(def.SourceTable)} {SrcAlias}
  479. WHERE {SrcAlias}.{Q(S0DimDefinition.TenantColumn)} = @TenantId
  480. ON DUPLICATE KEY UPDATE
  481. {string.Join(",\n ", updates)}
  482. """;
  483. }
  484. /// <summary>源侧工厂列:对齐旧路径的 <c>factory_id ?? FactoryId</c> 顺序;都不存在返回 null。</summary>
  485. private static string? SourceFactoryExpr(IReadOnlyList<S0DimSourceColumn> cols)
  486. {
  487. var hit = FindColumn(cols, "factory_id") ?? FindColumn(cols, "FactoryId");
  488. return hit is null ? null : $"{SrcAlias}.{Q(hit)}";
  489. }
  490. private static string? FindColumn(IReadOnlyList<S0DimSourceColumn> cols, string name) =>
  491. cols.FirstOrDefault(c => string.Equals(c.Name, name, StringComparison.OrdinalIgnoreCase))?.Name;
  492. }