| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465 |
- using System.Text;
- using System.Text.RegularExpressions;
- using Admin.NET.Plugin.AiDOP.DataPlatform;
- using Microsoft.Extensions.Logging;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.Schema;
- /// <summary>
- /// 中台表(mdp_std_* / mdp_stg_* / dwd_*)的增量对齐。
- /// 表不存在则按声明建表;已存在则只 ADD COLUMN / ADD KEY。不 DROP、不改类型、不改已有索引。
- /// </summary>
- public static class MdpSchemaAligner
- {
- private static readonly Regex CreateHead = new(
- @"^CREATE\s+TABLE\s+IF\s+NOT\s+EXISTS\s+`?(?<name>[A-Za-z0-9_]+)`?(?:\s+LIKE\s+`?(?<like>[A-Za-z0-9_]+)`?)?",
- RegexOptions.IgnoreCase | RegexOptions.Compiled | RegexOptions.Singleline);
- private static readonly Regex IndexHead = new(
- @"^(?<unique>UNIQUE\s+)?(?:KEY|INDEX)\s+`?(?<name>[A-Za-z0-9_]+)`?\s*\((?<cols>[^)]+)\)",
- 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);
- /// <summary>
- /// 非中立建表语句原样交给数据库;中立建表改为增量对齐,避免 CREATE IF NOT EXISTS 在表已存在时不补列。
- /// </summary>
- public static async Task<int> 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<int> ExecuteAsync(ISqlSugarClient db, string sql, List<SugarParameter> 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);
- }
- /// <summary>
- /// 给已有 source_system 的 mdp_std_* 补 written_by,并按贴源表名登记真实来源。
- /// 确证不了的行不改 source_system,只保持 written_by=LEGACY_UNVERIFIED。
- /// </summary>
- 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<NameRow>(
- """
- 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<NameRow>(
- $"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<NameRow>(
- """
- 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<NameRow>(
- """
- 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<string> stdTables)
- {
- var stg = await db.Ado.SqlQueryAsync<NameRow>(
- """
- 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<string> 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<string>();
- 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<int> 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<SchemaSnapshot> 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<NameRow>(
- """
- 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<NameRow>(
- """
- 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<ILogger<MdpSchemaAlignerLog>>()?.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<string> SplitStatements(string sql)
- {
- var list = new List<string>();
- 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<string> 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<string> SplitTopLevel(string body)
- {
- var list = new List<string>();
- 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;
- }
- /// <summary>仅用于取日志泛型,避免把静态类当作日志类别。</summary>
- private sealed class MdpSchemaAlignerLog;
- }
- public sealed class SchemaSnapshot
- {
- public static readonly SchemaSnapshot Missing = new(false, [], []);
- public SchemaSnapshot(bool exists, IEnumerable<string> columns, IEnumerable<string> indexes)
- {
- Exists = exists;
- Columns = new HashSet<string>(columns.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase);
- Indexes = new HashSet<string>(indexes.Where(x => !string.IsNullOrWhiteSpace(x)), StringComparer.OrdinalIgnoreCase);
- }
- public bool Exists { get; }
- public HashSet<string> Columns { get; }
- public HashSet<string> Indexes { get; }
- }
|