using System.Text; using System.Text.RegularExpressions; using Admin.NET.Plugin.AiDOP.DataPlatform; using Microsoft.Extensions.Logging; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Schema; /// /// 中台表(mdp_std_* / mdp_stg_* / dwd_*)的增量对齐。 /// 表不存在则按声明建表;已存在则只 ADD COLUMN / ADD KEY。不 DROP、不改类型、不改已有索引。 /// public static class MdpSchemaAligner { private static readonly Regex CreateHead = new( @"^CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\s+`?(?[A-Za-z0-9_]+)`?(?:\s+LIKE\s+`?(?[A-Za-z0-9_]+)`?)?", RegexOptions.IgnoreCase | RegexOptions.Compiled | RegexOptions.Singleline); private static readonly Regex IndexHead = new( @"^(?UNIQUE\s+)?(?:KEY|INDEX)\s+`?(?[A-Za-z0-9_]+)`?\s*\((?[^)]+)\)", RegexOptions.IgnoreCase | RegexOptions.Compiled); public static bool IsNeutralTable(string name) => name.StartsWith("mdp_std_", StringComparison.OrdinalIgnoreCase) || name.StartsWith("mdp_stg_", StringComparison.OrdinalIgnoreCase) || name.StartsWith("dwd_", StringComparison.OrdinalIgnoreCase); /// /// 非中立建表语句原样交给数据库;中立建表改为增量对齐,避免 CREATE IF NOT EXISTS 在表已存在时不补列。 /// public static async Task ExecuteAsync(ISqlSugarClient db, string sql, params SugarParameter[] parameters) { if (!ContainsNeutralCreate(sql)) return await db.Ado.ExecuteCommandAsync(sql, parameters); var applied = 0; foreach (var statement in SplitStatements(sql)) { if (IsNeutralCreate(statement)) applied += await AlignOneAsync(db, statement); else applied += await db.Ado.ExecuteCommandAsync(statement, parameters); } return applied; } public static Task ExecuteAsync(ISqlSugarClient db, string sql, List parameters) => ExecuteAsync(db, sql, parameters?.ToArray() ?? []); public static async Task EnsureWrittenByColumnAsync(ISqlSugarClient db, string table) { if (!IsSafeName(table)) throw new InvalidOperationException($"非法表名:{table}"); var hasWritten = await db.Ado.GetIntAsync( """ SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='written_by' """, new SugarParameter("@t", table)); if (hasWritten > 0) return; var ddl = $"ALTER TABLE `{table}` ADD COLUMN written_by VARCHAR(32) NOT NULL DEFAULT 'LEGACY_UNVERIFIED' COMMENT 'DB_SYNC/API_INBOUND/PLATFORM_FORM/SEED/LEGACY_UNVERIFIED'"; await db.Ado.ExecuteCommandAsync(ddl); Log(ddl); } /// /// 给已有 source_system 的 mdp_std_* 补 written_by,并按贴源表名登记真实来源。 /// 确证不了的行不改 source_system,只保持 written_by=LEGACY_UNVERIFIED。 /// public static async Task EnsureSourceIdentityAsync(ISqlSugarClient db) { await db.Ado.ExecuteCommandAsync( """ CREATE TABLE IF NOT EXISTS mdp_source_table_registry ( id BIGINT NOT NULL AUTO_INCREMENT, source_table VARCHAR(64) NOT NULL COMMENT '贴源记录的原始表名', source_system VARCHAR(50) NOT NULL COMMENT '真实来源系统码', remark VARCHAR(255) NULL, PRIMARY KEY (id), UNIQUE KEY uk_source_table (source_table) ) COMMENT='原始表名 → 真实来源系统 的唯一对照' """); var stgTables = await db.Ado.SqlQueryAsync( """ SELECT TABLE_NAME AS Name FROM information_schema.TABLES WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME LIKE 'mdp_stg%' """); foreach (var table in stgTables) { if (!IsSafeName(table.Name)) continue; var hasSourceTable = await db.Ado.GetIntAsync( """ SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='source_table' """, new SugarParameter("@t", table.Name)); if (hasSourceTable == 0) continue; var names = await db.Ado.SqlQueryAsync( $"SELECT DISTINCT source_table AS Name FROM `{table.Name}` WHERE source_table IS NOT NULL AND TRIM(source_table)<>''"); foreach (var row in names) { var system = MdpSourceIdentity.ClassifySourceTable(row.Name); if (system == null) { Log($"来源表未登记(命名规则无法确证,不猜测): {row.Name}"); continue; } await db.Ado.ExecuteCommandAsync( """ INSERT INTO mdp_source_table_registry (source_table, source_system, remark) SELECT @t, @s, '1.0.565 按表名规则登记' FROM DUAL WHERE NOT EXISTS (SELECT 1 FROM mdp_source_table_registry WHERE source_table=@t) """, new SugarParameter("@t", row.Name.Trim()), new SugarParameter("@s", system)); } } var stdTables = await db.Ado.SqlQueryAsync( """ SELECT TABLE_NAME AS Name FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND COLUMN_NAME='source_system' AND TABLE_NAME LIKE 'mdp_std%' """); foreach (var table in stdTables.Select(t => t.Name).Distinct(StringComparer.OrdinalIgnoreCase)) { if (!IsSafeName(table)) continue; var hasWritten = await db.Ado.GetIntAsync( """ SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='written_by' """, new SugarParameter("@t", table)); if (hasWritten == 0) { var ddl = $"ALTER TABLE `{table}` ADD COLUMN written_by VARCHAR(32) NOT NULL DEFAULT 'LEGACY_UNVERIFIED' COMMENT 'DB_SYNC/API_INBOUND/PLATFORM_FORM/SEED/LEGACY_UNVERIFIED'"; await db.Ado.ExecuteCommandAsync(ddl); Log(ddl); } var defaults = await db.Ado.SqlQueryAsync( """ SELECT COLUMN_DEFAULT AS Name FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME='source_system' """, new SugarParameter("@t", table)); var currentDefault = defaults.FirstOrDefault()?.Name?.Trim().Trim('\'') ?? ""; if (currentDefault == "AIDOP") { var ddl = $"ALTER TABLE `{table}` ALTER COLUMN source_system SET DEFAULT ''"; await db.Ado.ExecuteCommandAsync(ddl); Log(ddl); } await db.Ado.ExecuteCommandAsync( $""" UPDATE `{table}` SET written_by='DB_SYNC' WHERE written_by='LEGACY_UNVERIFIED' AND source_system IN ('DOPDEMORQ_SQLSERVER','T8','AIDOPDEV_MYSQL') """); await db.Ado.ExecuteCommandAsync( $""" UPDATE `{table}` SET written_by='SEED' WHERE written_by='LEGACY_UNVERIFIED' AND source_system IN ('UAT_GENERATOR','UAT','DEMO') """); await db.Ado.ExecuteCommandAsync( $""" UPDATE `{table}` SET written_by='API_INBOUND' WHERE written_by='LEGACY_UNVERIFIED' AND source_system='API' """); await db.Ado.ExecuteCommandAsync( $""" UPDATE `{table}` SET written_by='PLATFORM_FORM' WHERE written_by='LEGACY_UNVERIFIED' AND source_system='AIDOP_NATIVE' """); } await BackfillProvenSourceAsync(db, stdTables.Select(t => t.Name).Distinct(StringComparer.OrdinalIgnoreCase)); } private static async Task BackfillProvenSourceAsync(ISqlSugarClient db, IEnumerable stdTables) { var stg = await db.Ado.SqlQueryAsync( """ SELECT c.TABLE_NAME AS Name FROM information_schema.COLUMNS c WHERE c.TABLE_SCHEMA=DATABASE() AND c.TABLE_NAME LIKE 'mdp_stg%' AND c.COLUMN_NAME='source_biz_key' AND EXISTS ( SELECT 1 FROM information_schema.COLUMNS t WHERE t.TABLE_SCHEMA=DATABASE() AND t.TABLE_NAME=c.TABLE_NAME AND t.COLUMN_NAME='source_table') AND EXISTS ( SELECT 1 FROM information_schema.COLUMNS t WHERE t.TABLE_SCHEMA=DATABASE() AND t.TABLE_NAME=c.TABLE_NAME AND t.COLUMN_NAME='tenant_id') """); var stgNames = stg.Select(x => x.Name).Where(IsSafeName).Distinct(StringComparer.OrdinalIgnoreCase).ToList(); if (stgNames.Count == 0) return; var union = string.Join(" UNION ALL ", stgNames.Select(n => $"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")); var proven = $""" SELECT tenant_id, source_biz_key, MIN(source_system) AS source_system FROM ({union}) u GROUP BY tenant_id, source_biz_key HAVING COUNT(DISTINCT source_system)=1 """; foreach (var table in stdTables) { if (!IsSafeName(table)) continue; var ready = await db.Ado.GetIntAsync( """ SELECT COUNT(*) FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t AND COLUMN_NAME IN ('tenant_id','source_biz_key','source_system','written_by','id') """, new SugarParameter("@t", table)); if (ready < 5) continue; try { await db.Ado.ExecuteCommandAsync( $""" UPDATE `{table}` s JOIN ({proven}) x ON x.tenant_id=s.tenant_id AND x.source_biz_key=s.source_biz_key LEFT JOIN `{table}` o ON o.tenant_id=s.tenant_id AND o.source_system=x.source_system AND o.source_biz_key=s.source_biz_key AND o.id<>s.id SET s.source_system=x.source_system, s.written_by='DB_SYNC' WHERE s.source_system IN ('AIDOP','','AIDOP_LEGACY') AND o.id IS NULL """); } catch (Exception ex) { Log($"来源回填跳过 {table}: {ex.Message}"); } } } private static bool IsSafeName(string? name) => !string.IsNullOrWhiteSpace(name) && Regex.IsMatch(name, "^[A-Za-z0-9_]+$"); public static IReadOnlyList Plan(string statement, SchemaSnapshot snapshot) { var text = statement.Trim().TrimEnd(';').Trim(); var head = CreateHead.Match(text); if (!head.Success || !IsNeutralTable(head.Groups["name"].Value)) return [text]; var table = head.Groups["name"].Value; if (head.Groups["like"].Success) return snapshot.Exists ? [] : [text]; if (!snapshot.Exists) return [text]; var alters = new List(); foreach (var part in SplitTopLevel(Body(text))) { var line = part.Trim().TrimEnd(',').Trim(); if (line.Length == 0) continue; if (line.StartsWith("PRIMARY", StringComparison.OrdinalIgnoreCase) || line.StartsWith("CONSTRAINT", StringComparison.OrdinalIgnoreCase) || line.StartsWith("FULLTEXT", StringComparison.OrdinalIgnoreCase) || line.StartsWith("SPATIAL", StringComparison.OrdinalIgnoreCase)) continue; var index = IndexHead.Match(line); if (index.Success) { var indexName = index.Groups["name"].Value; if (snapshot.Indexes.Contains(indexName)) continue; var unique = index.Groups["unique"].Success ? "UNIQUE KEY" : "KEY"; alters.Add( $"ALTER TABLE `{table}` ADD {unique} `{indexName}` ({index.Groups["cols"].Value.Trim()})"); continue; } if (line.StartsWith("UNIQUE", StringComparison.OrdinalIgnoreCase) || line.StartsWith("KEY", StringComparison.OrdinalIgnoreCase) || line.StartsWith("INDEX", StringComparison.OrdinalIgnoreCase)) continue; var column = ColumnName(line); if (column.Length == 0 || snapshot.Columns.Contains(column)) continue; alters.Add($"ALTER TABLE `{table}` ADD COLUMN {line}"); } return alters; } private static async Task AlignOneAsync(ISqlSugarClient db, string statement) { var head = CreateHead.Match(statement.Trim()); var table = head.Groups["name"].Value; var snapshot = await LoadAsync(db, table); var planned = Plan(statement, snapshot); var applied = 0; foreach (var ddl in planned) { if (ddl.Contains("DROP ", StringComparison.OrdinalIgnoreCase)) throw new InvalidOperationException("中台结构对齐器禁止执行 DROP"); await db.Ado.ExecuteCommandAsync(ddl); Log(ddl); applied++; } return applied; } private static async Task LoadAsync(ISqlSugarClient db, string table) { var exists = await db.Ado.GetIntAsync( """ SELECT COUNT(*) FROM information_schema.TABLES WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t """, new SugarParameter("@t", table)); if (exists == 0) return SchemaSnapshot.Missing; var columns = await db.Ado.SqlQueryAsync( """ SELECT COLUMN_NAME AS Name FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t """, new SugarParameter("@t", table)); var indexes = await db.Ado.SqlQueryAsync( """ SELECT DISTINCT INDEX_NAME AS Name FROM information_schema.STATISTICS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t """, new SugarParameter("@t", table)); return new SchemaSnapshot( true, columns.Select(x => x.Name), indexes.Select(x => x.Name)); } private static void Log(string ddl) { try { App.GetService>()?.LogInformation("中台结构对齐 {Ddl}", ddl); } catch { // 单测与未建容器时不记日志 } } internal static bool ContainsNeutralCreate(string sql) => SplitStatements(sql).Any(IsNeutralCreate); private static bool IsNeutralCreate(string statement) { var head = CreateHead.Match(statement.Trim()); return head.Success && IsNeutralTable(head.Groups["name"].Value); } internal static List SplitStatements(string sql) { var list = new List(); var sb = new StringBuilder(); var depth = 0; var quote = false; for (var i = 0; i < sql.Length; i++) { var c = sql[i]; if (c == '\'') { if (quote && i + 1 < sql.Length && sql[i + 1] == '\'') { sb.Append("''"); i++; continue; } quote = !quote; sb.Append(c); continue; } if (!quote && c == '(') depth++; if (!quote && c == ')' && depth > 0) depth--; if (!quote && depth == 0 && c == ';') { Add(list, sb); continue; } sb.Append(c); } Add(list, sb); return list; } private static void Add(List list, StringBuilder sb) { var text = sb.ToString().Trim(); if (text.Length > 0) list.Add(text); sb.Clear(); } private static string Body(string create) { var start = create.IndexOf('('); var end = create.LastIndexOf(')'); if (start < 0 || end <= start) return string.Empty; return create.Substring(start + 1, end - start - 1); } private static List SplitTopLevel(string body) { var list = new List(); var sb = new StringBuilder(); var depth = 0; var quote = false; foreach (var c in body) { if (c == '\'') quote = !quote; if (!quote && c == '(') depth++; if (!quote && c == ')' && depth > 0) depth--; if (!quote && depth == 0 && c == ',') { list.Add(sb.ToString()); sb.Clear(); continue; } sb.Append(c); } if (sb.Length > 0) list.Add(sb.ToString()); return list; } private static string ColumnName(string line) { var token = line.Split([' ', '\t', '\r', '\n'], 2, StringSplitOptions.RemoveEmptyEntries)[0]; return token.Trim('`'); } private sealed class NameRow { public string Name { get; set; } = string.Empty; } /// 仅用于取日志泛型,避免把静态类当作日志类别。 private sealed class MdpSchemaAlignerLog; } public sealed class SchemaSnapshot { public static readonly SchemaSnapshot Missing = new(false, [], []); public SchemaSnapshot(bool exists, IEnumerable columns, IEnumerable indexes) { Exists = exists; Columns = new HashSet(columns.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase); Indexes = new HashSet(indexes.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase); } public bool Exists { get; } public HashSet Columns { get; } public HashSet Indexes { get; } }