MdpSchemaAligner.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465
  1. using System.Text;
  2. using System.Text.RegularExpressions;
  3. using Admin.NET.Plugin.AiDOP.DataPlatform;
  4. using Microsoft.Extensions.Logging;
  5. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
  6. /// <summary>
  7. /// 中台表(mdp_std_* / mdp_stg_* / dwd_*)的增量对齐。
  8. /// 表不存在则按声明建表;已存在则只 ADD COLUMN / ADD KEY。不 DROP、不改类型、不改已有索引。
  9. /// </summary>
  10. public static class MdpSchemaAligner
  11. {
  12. private static readonly Regex CreateHead = new(
  13. @"^CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\s+`?(?<name>[A-Za-z0-9_]+)`?(?:\s+LIKE\s+`?(?<like>[A-Za-z0-9_]+)`?)?",
  14. RegexOptions.IgnoreCase | RegexOptions.Compiled | RegexOptions.Singleline);
  15. private static readonly Regex IndexHead = new(
  16. @"^(?<unique>UNIQUE\s+)?(?:KEY|INDEX)\s+`?(?<name>[A-Za-z0-9_]+)`?\s*\((?<cols>[^)]+)\)",
  17. RegexOptions.IgnoreCase | RegexOptions.Compiled);
  18. public static bool IsNeutralTable(string name) =>
  19. name.StartsWith("mdp_std_", StringComparison.OrdinalIgnoreCase)
  20. || name.StartsWith("mdp_stg_", StringComparison.OrdinalIgnoreCase)
  21. || name.StartsWith("dwd_", StringComparison.OrdinalIgnoreCase);
  22. /// <summary>
  23. /// 非中立建表语句原样交给数据库;中立建表改为增量对齐,避免 CREATE IF NOT EXISTS 在表已存在时不补列。
  24. /// </summary>
  25. public static async Task<int> ExecuteAsync(ISqlSugarClient db, string sql, params SugarParameter[] parameters)
  26. {
  27. if (!ContainsNeutralCreate(sql))
  28. return await db.Ado.ExecuteCommandAsync(sql, parameters);
  29. var applied = 0;
  30. foreach (var statement in SplitStatements(sql))
  31. {
  32. if (IsNeutralCreate(statement))
  33. applied += await AlignOneAsync(db, statement);
  34. else
  35. applied += await db.Ado.ExecuteCommandAsync(statement, parameters);
  36. }
  37. return applied;
  38. }
  39. public static Task<int> ExecuteAsync(ISqlSugarClient db, string sql, List<SugarParameter> parameters) =>
  40. ExecuteAsync(db, sql, parameters?.ToArray() ?? []);
  41. public static async Task EnsureWrittenByColumnAsync(ISqlSugarClient db, string table)
  42. {
  43. if (!IsSafeName(table))
  44. throw new InvalidOperationException($"非法表名:{table}");
  45. var hasWritten = await db.Ado.GetIntAsync(
  46. """
  47. SELECT COUNT(*) FROM information_schema.COLUMNS
  48. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='written_by'
  49. """,
  50. new SugarParameter("@t", table));
  51. if (hasWritten > 0) return;
  52. var ddl =
  53. $"ALTER TABLE `{table}` ADD COLUMN written_by VARCHAR(32) NOT NULL DEFAULT 'LEGACY_UNVERIFIED' COMMENT 'DB_SYNC/API_INBOUND/PLATFORM_FORM/SEED/LEGACY_UNVERIFIED'";
  54. await db.Ado.ExecuteCommandAsync(ddl);
  55. Log(ddl);
  56. }
  57. /// <summary>
  58. /// 给已有 source_system 的 mdp_std_* 补 written_by,并按贴源表名登记真实来源。
  59. /// 确证不了的行不改 source_system,只保持 written_by=LEGACY_UNVERIFIED。
  60. /// </summary>
  61. public static async Task EnsureSourceIdentityAsync(ISqlSugarClient db)
  62. {
  63. await db.Ado.ExecuteCommandAsync(
  64. """
  65. CREATE TABLE IF NOT EXISTS mdp_source_table_registry (
  66. id BIGINT NOT NULL AUTO_INCREMENT,
  67. source_table VARCHAR(64) NOT NULL COMMENT '贴源记录的原始表名',
  68. source_system VARCHAR(50) NOT NULL COMMENT '真实来源系统码',
  69. remark VARCHAR(255) NULL,
  70. PRIMARY KEY (id),
  71. UNIQUE KEY uk_source_table (source_table)
  72. ) COMMENT='原始表名 → 真实来源系统 的唯一对照'
  73. """);
  74. var stgTables = await db.Ado.SqlQueryAsync<NameRow>(
  75. """
  76. SELECT TABLE_NAME AS Name FROM information_schema.TABLES
  77. WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME LIKE 'mdp_stg%'
  78. """);
  79. foreach (var table in stgTables)
  80. {
  81. if (!IsSafeName(table.Name)) continue;
  82. var hasSourceTable = await db.Ado.GetIntAsync(
  83. """
  84. SELECT COUNT(*) FROM information_schema.COLUMNS
  85. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='source_table'
  86. """,
  87. new SugarParameter("@t", table.Name));
  88. if (hasSourceTable == 0) continue;
  89. var names = await db.Ado.SqlQueryAsync<NameRow>(
  90. $"SELECT DISTINCT source_table AS Name FROM `{table.Name}` WHERE source_table IS NOT NULL AND TRIM(source_table)<>''");
  91. foreach (var row in names)
  92. {
  93. var system = MdpSourceIdentity.ClassifySourceTable(row.Name);
  94. if (system == null)
  95. {
  96. Log($"来源表未登记(命名规则无法确证,不猜测): {row.Name}");
  97. continue;
  98. }
  99. await db.Ado.ExecuteCommandAsync(
  100. """
  101. INSERT INTO mdp_source_table_registry (source_table, source_system, remark)
  102. SELECT @t, @s, '1.0.565 按表名规则登记'
  103. FROM DUAL
  104. WHERE NOT EXISTS (SELECT 1 FROM mdp_source_table_registry WHERE source_table=@t)
  105. """,
  106. new SugarParameter("@t", row.Name.Trim()),
  107. new SugarParameter("@s", system));
  108. }
  109. }
  110. var stdTables = await db.Ado.SqlQueryAsync<NameRow>(
  111. """
  112. SELECT TABLE_NAME AS Name FROM information_schema.COLUMNS
  113. WHERE TABLE_SCHEMA=DATABASE() AND COLUMN_NAME='source_system' AND TABLE_NAME LIKE 'mdp_std%'
  114. """);
  115. foreach (var table in stdTables.Select(t => t.Name).Distinct(StringComparer.OrdinalIgnoreCase))
  116. {
  117. if (!IsSafeName(table)) continue;
  118. var hasWritten = await db.Ado.GetIntAsync(
  119. """
  120. SELECT COUNT(*) FROM information_schema.COLUMNS
  121. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='written_by'
  122. """,
  123. new SugarParameter("@t", table));
  124. if (hasWritten == 0)
  125. {
  126. var ddl =
  127. $"ALTER TABLE `{table}` ADD COLUMN written_by VARCHAR(32) NOT NULL DEFAULT 'LEGACY_UNVERIFIED' COMMENT 'DB_SYNC/API_INBOUND/PLATFORM_FORM/SEED/LEGACY_UNVERIFIED'";
  128. await db.Ado.ExecuteCommandAsync(ddl);
  129. Log(ddl);
  130. }
  131. var defaults = await db.Ado.SqlQueryAsync<NameRow>(
  132. """
  133. SELECT COLUMN_DEFAULT AS Name FROM information_schema.COLUMNS
  134. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='source_system'
  135. """,
  136. new SugarParameter("@t", table));
  137. var currentDefault = defaults.FirstOrDefault()?.Name?.Trim().Trim('\'') ?? "";
  138. if (currentDefault == "AIDOP")
  139. {
  140. var ddl = $"ALTER TABLE `{table}` ALTER COLUMN source_system SET DEFAULT ''";
  141. await db.Ado.ExecuteCommandAsync(ddl);
  142. Log(ddl);
  143. }
  144. await db.Ado.ExecuteCommandAsync(
  145. $"""
  146. UPDATE `{table}` SET written_by='DB_SYNC'
  147. WHERE written_by='LEGACY_UNVERIFIED'
  148. AND source_system IN ('DOPDEMORQ_SQLSERVER','T8','AIDOPDEV_MYSQL')
  149. """);
  150. await db.Ado.ExecuteCommandAsync(
  151. $"""
  152. UPDATE `{table}` SET written_by='SEED'
  153. WHERE written_by='LEGACY_UNVERIFIED'
  154. AND source_system IN ('UAT_GENERATOR','UAT','DEMO')
  155. """);
  156. await db.Ado.ExecuteCommandAsync(
  157. $"""
  158. UPDATE `{table}` SET written_by='API_INBOUND'
  159. WHERE written_by='LEGACY_UNVERIFIED' AND source_system='API'
  160. """);
  161. await db.Ado.ExecuteCommandAsync(
  162. $"""
  163. UPDATE `{table}` SET written_by='PLATFORM_FORM'
  164. WHERE written_by='LEGACY_UNVERIFIED' AND source_system='AIDOP_NATIVE'
  165. """);
  166. }
  167. await BackfillProvenSourceAsync(db, stdTables.Select(t => t.Name).Distinct(StringComparer.OrdinalIgnoreCase));
  168. }
  169. private static async Task BackfillProvenSourceAsync(ISqlSugarClient db, IEnumerable<string> stdTables)
  170. {
  171. var stg = await db.Ado.SqlQueryAsync<NameRow>(
  172. """
  173. SELECT c.TABLE_NAME AS Name
  174. FROM information_schema.COLUMNS c
  175. WHERE c.TABLE_SCHEMA=DATABASE() AND c.TABLE_NAME LIKE 'mdp_stg%' AND c.COLUMN_NAME='source_biz_key'
  176. AND EXISTS (
  177. SELECT 1 FROM information_schema.COLUMNS t
  178. WHERE t.TABLE_SCHEMA=DATABASE() AND t.TABLE_NAME=c.TABLE_NAME AND t.COLUMN_NAME='source_table')
  179. AND EXISTS (
  180. SELECT 1 FROM information_schema.COLUMNS t
  181. WHERE t.TABLE_SCHEMA=DATABASE() AND t.TABLE_NAME=c.TABLE_NAME AND t.COLUMN_NAME='tenant_id')
  182. """);
  183. var stgNames = stg.Select(x => x.Name).Where(IsSafeName).Distinct(StringComparer.OrdinalIgnoreCase).ToList();
  184. if (stgNames.Count == 0) return;
  185. var union = string.Join(" UNION ALL ", stgNames.Select(n =>
  186. $"SELECT g.tenant_id, g.source_biz_key, r.source_system FROM `{n}` g JOIN mdp_source_table_registry r ON r.source_table=g.source_table"));
  187. var proven =
  188. $"""
  189. SELECT tenant_id, source_biz_key, MIN(source_system) AS source_system
  190. FROM ({union}) u
  191. GROUP BY tenant_id, source_biz_key
  192. HAVING COUNT(DISTINCT source_system)=1
  193. """;
  194. foreach (var table in stdTables)
  195. {
  196. if (!IsSafeName(table)) continue;
  197. var ready = await db.Ado.GetIntAsync(
  198. """
  199. SELECT COUNT(*) FROM information_schema.COLUMNS
  200. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
  201. AND COLUMN_NAME IN ('tenant_id','source_biz_key','source_system','written_by','id')
  202. """,
  203. new SugarParameter("@t", table));
  204. if (ready < 5) continue;
  205. try
  206. {
  207. await db.Ado.ExecuteCommandAsync(
  208. $"""
  209. UPDATE `{table}` s
  210. JOIN ({proven}) x ON x.tenant_id=s.tenant_id AND x.source_biz_key=s.source_biz_key
  211. LEFT JOIN `{table}` o
  212. ON o.tenant_id=s.tenant_id AND o.source_system=x.source_system
  213. AND o.source_biz_key=s.source_biz_key AND o.id<>s.id
  214. SET s.source_system=x.source_system, s.written_by='DB_SYNC'
  215. WHERE s.source_system IN ('AIDOP','','AIDOP_LEGACY')
  216. AND o.id IS NULL
  217. """);
  218. }
  219. catch (Exception ex)
  220. {
  221. Log($"来源回填跳过 {table}: {ex.Message}");
  222. }
  223. }
  224. }
  225. private static bool IsSafeName(string? name) =>
  226. !string.IsNullOrWhiteSpace(name) && Regex.IsMatch(name, "^[A-Za-z0-9_]+$");
  227. public static IReadOnlyList<string> Plan(string statement, SchemaSnapshot snapshot)
  228. {
  229. var text = statement.Trim().TrimEnd(';').Trim();
  230. var head = CreateHead.Match(text);
  231. if (!head.Success || !IsNeutralTable(head.Groups["name"].Value))
  232. return [text];
  233. var table = head.Groups["name"].Value;
  234. if (head.Groups["like"].Success)
  235. return snapshot.Exists ? [] : [text];
  236. if (!snapshot.Exists)
  237. return [text];
  238. var alters = new List<string>();
  239. foreach (var part in SplitTopLevel(Body(text)))
  240. {
  241. var line = part.Trim().TrimEnd(',').Trim();
  242. if (line.Length == 0) continue;
  243. if (line.StartsWith("PRIMARY", StringComparison.OrdinalIgnoreCase)
  244. || line.StartsWith("CONSTRAINT", StringComparison.OrdinalIgnoreCase)
  245. || line.StartsWith("FULLTEXT", StringComparison.OrdinalIgnoreCase)
  246. || line.StartsWith("SPATIAL", StringComparison.OrdinalIgnoreCase))
  247. continue;
  248. var index = IndexHead.Match(line);
  249. if (index.Success)
  250. {
  251. var indexName = index.Groups["name"].Value;
  252. if (snapshot.Indexes.Contains(indexName)) continue;
  253. var unique = index.Groups["unique"].Success ? "UNIQUE KEY" : "KEY";
  254. alters.Add(
  255. $"ALTER TABLE `{table}` ADD {unique} `{indexName}` ({index.Groups["cols"].Value.Trim()})");
  256. continue;
  257. }
  258. if (line.StartsWith("UNIQUE", StringComparison.OrdinalIgnoreCase)
  259. || line.StartsWith("KEY", StringComparison.OrdinalIgnoreCase)
  260. || line.StartsWith("INDEX", StringComparison.OrdinalIgnoreCase))
  261. continue;
  262. var column = ColumnName(line);
  263. if (column.Length == 0 || snapshot.Columns.Contains(column)) continue;
  264. alters.Add($"ALTER TABLE `{table}` ADD COLUMN {line}");
  265. }
  266. return alters;
  267. }
  268. private static async Task<int> AlignOneAsync(ISqlSugarClient db, string statement)
  269. {
  270. var head = CreateHead.Match(statement.Trim());
  271. var table = head.Groups["name"].Value;
  272. var snapshot = await LoadAsync(db, table);
  273. var planned = Plan(statement, snapshot);
  274. var applied = 0;
  275. foreach (var ddl in planned)
  276. {
  277. if (ddl.Contains("DROP ", StringComparison.OrdinalIgnoreCase))
  278. throw new InvalidOperationException("中台结构对齐器禁止执行 DROP");
  279. await db.Ado.ExecuteCommandAsync(ddl);
  280. Log(ddl);
  281. applied++;
  282. }
  283. return applied;
  284. }
  285. private static async Task<SchemaSnapshot> LoadAsync(ISqlSugarClient db, string table)
  286. {
  287. var exists = await db.Ado.GetIntAsync(
  288. """
  289. SELECT COUNT(*) FROM information_schema.TABLES
  290. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
  291. """,
  292. new SugarParameter("@t", table));
  293. if (exists == 0)
  294. return SchemaSnapshot.Missing;
  295. var columns = await db.Ado.SqlQueryAsync<NameRow>(
  296. """
  297. SELECT COLUMN_NAME AS Name FROM information_schema.COLUMNS
  298. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
  299. """,
  300. new SugarParameter("@t", table));
  301. var indexes = await db.Ado.SqlQueryAsync<NameRow>(
  302. """
  303. SELECT DISTINCT INDEX_NAME AS Name FROM information_schema.STATISTICS
  304. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
  305. """,
  306. new SugarParameter("@t", table));
  307. return new SchemaSnapshot(
  308. true,
  309. columns.Select(x => x.Name),
  310. indexes.Select(x => x.Name));
  311. }
  312. private static void Log(string ddl)
  313. {
  314. try
  315. {
  316. App.GetService<ILogger<MdpSchemaAlignerLog>>()?.LogInformation("中台结构对齐 {Ddl}", ddl);
  317. }
  318. catch
  319. {
  320. // 单测与未建容器时不记日志
  321. }
  322. }
  323. internal static bool ContainsNeutralCreate(string sql) =>
  324. SplitStatements(sql).Any(IsNeutralCreate);
  325. private static bool IsNeutralCreate(string statement)
  326. {
  327. var head = CreateHead.Match(statement.Trim());
  328. return head.Success && IsNeutralTable(head.Groups["name"].Value);
  329. }
  330. internal static List<string> SplitStatements(string sql)
  331. {
  332. var list = new List<string>();
  333. var sb = new StringBuilder();
  334. var depth = 0;
  335. var quote = false;
  336. for (var i = 0; i < sql.Length; i++)
  337. {
  338. var c = sql[i];
  339. if (c == '\'')
  340. {
  341. if (quote && i + 1 < sql.Length && sql[i + 1] == '\'')
  342. {
  343. sb.Append("''");
  344. i++;
  345. continue;
  346. }
  347. quote = !quote;
  348. sb.Append(c);
  349. continue;
  350. }
  351. if (!quote && c == '(') depth++;
  352. if (!quote && c == ')' && depth > 0) depth--;
  353. if (!quote && depth == 0 && c == ';')
  354. {
  355. Add(list, sb);
  356. continue;
  357. }
  358. sb.Append(c);
  359. }
  360. Add(list, sb);
  361. return list;
  362. }
  363. private static void Add(List<string> list, StringBuilder sb)
  364. {
  365. var text = sb.ToString().Trim();
  366. if (text.Length > 0) list.Add(text);
  367. sb.Clear();
  368. }
  369. private static string Body(string create)
  370. {
  371. var start = create.IndexOf('(');
  372. var end = create.LastIndexOf(')');
  373. if (start < 0 || end <= start) return string.Empty;
  374. return create.Substring(start + 1, end - start - 1);
  375. }
  376. private static List<string> SplitTopLevel(string body)
  377. {
  378. var list = new List<string>();
  379. var sb = new StringBuilder();
  380. var depth = 0;
  381. var quote = false;
  382. foreach (var c in body)
  383. {
  384. if (c == '\'') quote = !quote;
  385. if (!quote && c == '(') depth++;
  386. if (!quote && c == ')' && depth > 0) depth--;
  387. if (!quote && depth == 0 && c == ',')
  388. {
  389. list.Add(sb.ToString());
  390. sb.Clear();
  391. continue;
  392. }
  393. sb.Append(c);
  394. }
  395. if (sb.Length > 0) list.Add(sb.ToString());
  396. return list;
  397. }
  398. private static string ColumnName(string line)
  399. {
  400. var token = line.Split([' ', '\t', '\r', '\n'], 2, StringSplitOptions.RemoveEmptyEntries)[0];
  401. return token.Trim('`');
  402. }
  403. private sealed class NameRow
  404. {
  405. public string Name { get; set; } = string.Empty;
  406. }
  407. /// <summary>仅用于取日志泛型,避免把静态类当作日志类别。</summary>
  408. private sealed class MdpSchemaAlignerLog;
  409. }
  410. public sealed class SchemaSnapshot
  411. {
  412. public static readonly SchemaSnapshot Missing = new(false, [], []);
  413. public SchemaSnapshot(bool exists, IEnumerable<string> columns, IEnumerable<string> indexes)
  414. {
  415. Exists = exists;
  416. Columns = new HashSet<string>(columns.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase);
  417. Indexes = new HashSet<string>(indexes.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase);
  418. }
  419. public bool Exists { get; }
  420. public HashSet<string> Columns { get; }
  421. public HashSet<string> Indexes { get; }
  422. }