using System.Data; using System.Globalization; using System.Text.Json; using System.Text.Json.Serialization; using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; using SqlSugar; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors; /// /// 方式甲:从 mdp_source DB 源按 mdp_entity.source_table_name 增量抽数 → target_table_name(贴源)。 /// public sealed class MdpDbPullExecutor : IMdpSourcePullExecutor, ITransient { public string SupportedType => "DB_SYNC"; /// /// 贴源 raw_data 的时间统一按 MySQL 字面量格式输出;各域转换层用 /// STR_TO_DATE(..., '%Y-%m-%d %H:%i:%s.%f') 解析,ISO 8601 的 T 分隔符会解析失败。 /// private static readonly JsonSerializerOptions RawDataJsonOptions = new() { Converters = { new MysqlLiteralDateTimeConverter(), new MysqlLiteralDateTimeOffsetConverter() } }; private readonly MdpSourceScopeFactory _scopeFactory; private readonly ISqlSugarClient _db; private readonly MdpStagingWriter _writer; public MdpDbPullExecutor(MdpSourceScopeFactory scopeFactory, ISqlSugarClient db, MdpStagingWriter writer) { _scopeFactory = scopeFactory; _db = db; _writer = writer; } public async Task PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default) { if (string.IsNullOrWhiteSpace(entity.SourceTableName)) throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 source_table_name"); if (string.IsNullOrWhiteSpace(entity.TargetTableName)) throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 target_table_name"); var startedAt = DateTime.Now; var scope = await _scopeFactory.GetScopeAsync(source.SourceCode, cancellationToken); var batchSize = entity.BatchSize > 0 ? entity.BatchSize : 1000; var isSqlServer = string.Equals(source.DbType, "SQLSERVER", StringComparison.OrdinalIgnoreCase) || string.Equals(source.DbType, "MSSQL", StringComparison.OrdinalIgnoreCase); var useKeyset = ctx.UseKeysetCursor && !string.IsNullOrWhiteSpace(ctx.CursorColumn) && !string.IsNullOrWhiteSpace(ctx.TieBreakerColumn); string? tenantCol = null; string? factoryCol = null; if (ctx.TenantId > 0) { var cols = await TryGetSourceColumnsAsync(scope, entity.SourceTableName!, isSqlServer, cancellationToken); tenantCol = FindColumn(cols, "tenant_id", "TenantId", "TenantID"); if (ctx.FactoryId > 0) factoryCol = FindColumn(cols, "factory_id", "FactoryId"); } var sql = useKeyset ? BuildKeysetSelectSql(entity, isSqlServer, batchSize, ctx, tenantCol, factoryCol) : BuildSelectSql(entity, isSqlServer, batchSize, ctx, tenantCol, factoryCol); var parameters = new List(); if (tenantCol != null) parameters.Add(new SugarParameter("@scopeTenantId", ctx.TenantId)); if (factoryCol != null) parameters.Add(new SugarParameter("@scopeFactoryId", ctx.FactoryId)); if (useKeyset) { AddKeysetParameters(parameters, ctx); } else { if (ctx.WindowFrom.HasValue) parameters.Add(new SugarParameter("@windowFrom", ctx.WindowFrom.Value)); if (!ctx.FullRefresh && !string.Equals(ctx.SyncWindowType, "FULL", StringComparison.OrdinalIgnoreCase) && !string.IsNullOrWhiteSpace(entity.IncrColumn) && !string.IsNullOrWhiteSpace(entity.LastCursor) && !string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase)) parameters.Add(new SugarParameter("@cursor", entity.LastCursor)); // ROLLING:用窗口下界;若已有水位且水位晚于下界,则进一步用水位收窄 if (!ctx.FullRefresh && string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase) && !string.IsNullOrWhiteSpace(entity.IncrColumn) && !string.IsNullOrWhiteSpace(entity.LastCursor) && ctx.WindowFrom.HasValue && DateTime.TryParse(entity.LastCursor, out var cursorDt) && cursorDt > ctx.WindowFrom.Value) parameters.Add(new SugarParameter("@cursor", entity.LastCursor)); } var table = await scope.Ado.GetDataTableAsync(sql, parameters); cancellationToken.ThrowIfCancellationRequested(); var written = 0; string? maxCursor = entity.LastCursor; var now = DateTime.Now; string? lastKeysetCursor = null; string? lastKeysetTie = null; foreach (DataRow row in table.Rows) { cancellationToken.ThrowIfCancellationRequested(); var dict = new Dictionary(StringComparer.OrdinalIgnoreCase); foreach (DataColumn col in table.Columns) dict[col.ColumnName] = row[col] == DBNull.Value ? null : row[col]; var sourceRowId = ResolveSourceRowId(dict); var rawJson = JsonSerializer.Serialize(dict, RawDataJsonOptions); if (useKeyset) { lastKeysetCursor = FormatCursorValue(dict, ctx.CursorColumn!); lastKeysetTie = FormatCursorValue(dict, ctx.TieBreakerColumn!); if (!string.IsNullOrEmpty(lastKeysetTie)) maxCursor = EncodeKeysetCursor(lastKeysetCursor, lastKeysetTie); } else if (!string.IsNullOrWhiteSpace(entity.IncrColumn) && dict.TryGetValue(entity.IncrColumn, out var incrVal) && incrVal != null) { var cursor = incrVal is DateTime dt ? dt.ToString("yyyy-MM-dd HH:mm:ss.fff") : incrVal.ToString(); if (!string.IsNullOrEmpty(cursor) && (maxCursor == null || string.CompareOrdinal(cursor, maxCursor) > 0)) maxCursor = cursor; } written += await _writer.UpsertAsync( source, entity, entity.SourceTableName!, dict, rawJson, sourceRowId, ctx); } if (useKeyset && !string.IsNullOrEmpty(lastKeysetTie)) { ctx.CursorValue = lastKeysetCursor; ctx.TieBreakerValue = lastKeysetTie; } if (!ctx.DeferCursorPersist && !string.IsNullOrEmpty(maxCursor) && maxCursor != entity.LastCursor) { await _db.Updateable() .SetColumns(x => new MdpEntity { LastCursor = maxCursor, LastSyncTo = now, UpdateTime = now }) .Where(x => x.Id == entity.Id) .ExecuteCommandAsync(cancellationToken); } await WriteSyncLogAsync(source, entity, ctx, startedAt, table.Rows.Count, written, null); var windowHint = useKeyset ? (ctx.NullTimePhase ? "KEYSET_NULL" : "KEYSET") : (ctx.SyncWindowType ?? (ctx.FullRefresh ? "FULL" : "INCR")); return new MdpPullResult { RowsPulled = table.Rows.Count, RowsWritten = written, NewCursor = maxCursor, Message = $"OK window={windowHint}" + (ctx.WindowFrom.HasValue ? $" from={ctx.WindowFrom:yyyy-MM-dd HH:mm:ss}" : "") + (ctx.BootstrapFrom.HasValue ? $" bootstrapFrom={ctx.BootstrapFrom:yyyy-MM-dd HH:mm:ss}" : "") }; } private static string BuildSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null) { var table = entity.SourceTableName!.Trim(); // 仅允许简单标识符/schema.table,防注入 if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_\.\[\]]+$")) throw new InvalidOperationException($"非法 source_table_name:{table}"); var predicates = new List(); string orderBy; if (!string.IsNullOrWhiteSpace(entity.IncrColumn)) { var incr = entity.IncrColumn.Trim(); if (!System.Text.RegularExpressions.Regex.IsMatch(incr, @"^[A-Za-z0-9_]+$")) throw new InvalidOperationException($"非法 incr_column:{incr}"); var incrExpr = isSqlServer ? incr : $"`{incr}`"; orderBy = incrExpr; var isRolling = string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase); if (!ctx.FullRefresh && isRolling && ctx.WindowFrom.HasValue) predicates.Add($"{incrExpr} >= @windowFrom"); var useCursor = !ctx.FullRefresh && !string.IsNullOrWhiteSpace(entity.LastCursor) && (!isRolling || (ctx.WindowFrom.HasValue && DateTime.TryParse(entity.LastCursor, out var cursorDt) && cursorDt > ctx.WindowFrom.Value)); if (useCursor) predicates.Add($"{incrExpr} > @cursor"); } else { orderBy = isSqlServer ? "(SELECT NULL)" : "1"; } AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol); var where = predicates.Count > 0 ? " WHERE " + string.Join(" AND ", predicates) : ""; var offset = ctx.Offset > 0 ? ctx.Offset : 0; if (isSqlServer) { // SQL Server:OFFSET 需 ORDER BY;无 incr 时用稳定键兜底 if (string.IsNullOrWhiteSpace(entity.IncrColumn)) orderBy = "(SELECT NULL)"; return $"SELECT * FROM {table}{where} ORDER BY {orderBy} OFFSET {offset} ROWS FETCH NEXT {batchSize} ROWS ONLY"; } return offset > 0 ? $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize} OFFSET {offset}" : $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize}"; } private static string BuildKeysetSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null) { var table = entity.SourceTableName!.Trim(); if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_\.\[\]]+$")) throw new InvalidOperationException($"非法 source_table_name:{table}"); var cursorCol = RequireIdent(ctx.CursorColumn, "CursorColumn"); var tieCol = RequireIdent(ctx.TieBreakerColumn, "TieBreakerColumn"); var cursorExpr = isSqlServer ? cursorCol : $"`{cursorCol}`"; var tieExpr = isSqlServer ? tieCol : $"`{tieCol}`"; var predicates = new List(); if (ctx.NullTimePhase) { predicates.Add($"{cursorExpr} IS NULL"); if (!string.IsNullOrWhiteSpace(ctx.TieBreakerValue)) predicates.Add($"{tieExpr} > @tieBreaker"); } else { predicates.Add($"{cursorExpr} IS NOT NULL"); if (ctx.BootstrapFrom.HasValue) predicates.Add($"{cursorExpr} >= @bootstrapFrom"); // keyset lower bound predicates.Add($""" ( @cursorTime IS NULL OR {cursorExpr} > @cursorTime OR ({cursorExpr} = @cursorTime AND {tieExpr} > @tieBreaker) ) """); if (!string.IsNullOrWhiteSpace(ctx.UpperCursorValue)) { predicates.Add($""" ( {cursorExpr} < @upperTime OR ({cursorExpr} = @upperTime AND {tieExpr} <= @upperTie) ) """); } } AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol); var where = " WHERE " + string.Join(" AND ", predicates); var orderBy = ctx.NullTimePhase ? tieExpr : $"{cursorExpr}, {tieExpr}"; if (isSqlServer) return $"SELECT * FROM {table}{where} ORDER BY {orderBy} OFFSET 0 ROWS FETCH NEXT {batchSize} ROWS ONLY"; return $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize}"; } private static void AddKeysetParameters(List parameters, MdpPullContext ctx) { if (ctx.BootstrapFrom.HasValue) parameters.Add(new SugarParameter("@bootstrapFrom", ctx.BootstrapFrom.Value)); if (ctx.NullTimePhase) { parameters.Add(new SugarParameter("@tieBreaker", string.IsNullOrWhiteSpace(ctx.TieBreakerValue) ? "0" : ctx.TieBreakerValue)); return; } object cursorTime = string.IsNullOrWhiteSpace(ctx.CursorValue) ? DBNull.Value : (object)ctx.CursorValue!; parameters.Add(new SugarParameter("@cursorTime", cursorTime)); parameters.Add(new SugarParameter("@tieBreaker", string.IsNullOrWhiteSpace(ctx.TieBreakerValue) ? "0" : ctx.TieBreakerValue)); if (!string.IsNullOrWhiteSpace(ctx.UpperCursorValue)) { parameters.Add(new SugarParameter("@upperTime", ctx.UpperCursorValue)); parameters.Add(new SugarParameter("@upperTie", string.IsNullOrWhiteSpace(ctx.UpperTieBreakerValue) ? "0" : ctx.UpperTieBreakerValue)); } } private static void AppendScopePredicates(List predicates, bool isSqlServer, string? tenantCol, string? factoryCol) { if (!string.IsNullOrWhiteSpace(tenantCol)) predicates.Add($"{QuoteIdent(tenantCol, isSqlServer)} = @scopeTenantId"); if (!string.IsNullOrWhiteSpace(factoryCol)) { var f = QuoteIdent(factoryCol, isSqlServer); predicates.Add($"COALESCE(NULLIF({f}, 0), 1) = @scopeFactoryId"); } } private static string QuoteIdent(string name, bool isSqlServer) => isSqlServer ? name : $"`{name}`"; private static string? FindColumn(IReadOnlyCollection columns, params string[] names) { foreach (var n in names) { var hit = columns.FirstOrDefault(c => string.Equals(c, n, StringComparison.OrdinalIgnoreCase)); if (hit != null) return hit; } return null; } private static async Task> TryGetSourceColumnsAsync( ISqlSugarClient scope, string table, bool isSqlServer, CancellationToken cancellationToken) { try { var simple = table.Contains('.') ? table.Split('.').Last().Trim('[', ']') : table.Trim('[', ']'); if (!System.Text.RegularExpressions.Regex.IsMatch(simple, @"^[A-Za-z0-9_]+$")) return Array.Empty(); if (isSqlServer) { var rows = await scope.Ado.SqlQueryAsync( "SELECT name FROM sys.columns WHERE object_id = OBJECT_ID(@t)", new SugarParameter("@t", simple)); return rows ?? new List(); } var mysql = await scope.Ado.SqlQueryAsync( """ SELECT COLUMN_NAME FROM information_schema.COLUMNS WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t """, new SugarParameter("@t", simple)); return mysql ?? new List(); } catch { return Array.Empty(); } } private static string RequireIdent(string? name, string label) { var v = (name ?? "").Trim(); if (!System.Text.RegularExpressions.Regex.IsMatch(v, @"^[A-Za-z0-9_]+$")) throw new InvalidOperationException($"非法 {label}:{name}"); return v; } private static string? FormatCursorValue(Dictionary dict, string column) { if (!dict.TryGetValue(column, out var raw) || raw == null || raw is DBNull) return null; return raw switch { DateTime dt => dt.ToString("yyyy-MM-dd HH:mm:ss.fff", CultureInfo.InvariantCulture), DateTimeOffset dto => dto.ToString("yyyy-MM-dd HH:mm:ss.fff", CultureInfo.InvariantCulture), _ => raw.ToString() }; } internal static string EncodeKeysetCursor(string? cursor, string tieBreaker) => System.Text.Json.JsonSerializer.Serialize(new KeysetCursorDto { Cursor = cursor, TieBreaker = tieBreaker }); internal static bool TryDecodeKeysetCursor(string? raw, out string? cursor, out string? tieBreaker) { cursor = null; tieBreaker = null; if (string.IsNullOrWhiteSpace(raw)) return false; try { var dto = System.Text.Json.JsonSerializer.Deserialize(raw); if (dto == null || string.IsNullOrWhiteSpace(dto.TieBreaker)) return false; cursor = dto.Cursor; tieBreaker = dto.TieBreaker; return true; } catch { return false; } } private sealed class KeysetCursorDto { public string? Cursor { get; set; } public string? TieBreaker { get; set; } } private static string ResolveSourceRowId(Dictionary dict) { // InvTransHist 等表常有空字符串 ID,若优先取 ID 会导致 source_row_id="" 唯一键互相覆盖。 // 跳过 null/空白后,再落到 RecID 等真实键。 foreach (var key in new[] { "RecID", "RecId", "recid", "Recid", "id", "Id", "ID", "noid", "billno", "BillNo", "djbh", "Djbh" }) { if (!dict.TryGetValue(key, out var v) || v is null || v is DBNull) continue; var s = v.ToString(); if (!string.IsNullOrWhiteSpace(s)) return s; } return Guid.NewGuid().ToString("N"); } private async Task WriteSyncLogAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, DateTime startedAt, int pulled, int written, string? error) { try { var endedAt = DateTime.Now; var syncType = ctx.FullRefresh || string.Equals(ctx.SyncWindowType, "FULL", StringComparison.OrdinalIgnoreCase) ? "FULL" : "INCR"; await _db.Ado.ExecuteCommandAsync(@" INSERT INTO mdp_sync_log (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type, sync_start, sync_end, duration_ms, rows_read, rows_written, status, error_message) VALUES (@tenantId, @entityId, @sourceCode, @entityName, @batchId, @syncType, 'AUTO', @syncStart, @syncEnd, @durationMs, @rowsRead, @rowsWritten, @status, @error)", new SugarParameter("@tenantId", ctx.TenantId), new SugarParameter("@entityId", entity.Id), new SugarParameter("@sourceCode", source.SourceCode), new SugarParameter("@entityName", string.IsNullOrWhiteSpace(entity.EntityName) ? entity.EntityCode : entity.EntityName), new SugarParameter("@batchId", ctx.BatchId), new SugarParameter("@syncType", syncType), new SugarParameter("@syncStart", startedAt), new SugarParameter("@syncEnd", endedAt), new SugarParameter("@durationMs", (int)(endedAt - startedAt).TotalMilliseconds), new SugarParameter("@rowsRead", pulled), new SugarParameter("@rowsWritten", written), new SugarParameter("@status", string.IsNullOrEmpty(error) ? "SUCCESS" : "FAILED"), new SugarParameter("@error", error)); } catch { // 日志表结构可能与种子不一致时不阻断主流程 } } private sealed class MysqlLiteralDateTimeConverter : JsonConverter { public override DateTime Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options) => reader.GetDateTime(); public override void Write(Utf8JsonWriter writer, DateTime value, JsonSerializerOptions options) => writer.WriteStringValue(value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture)); } private sealed class MysqlLiteralDateTimeOffsetConverter : JsonConverter { public override DateTimeOffset Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options) => reader.GetDateTimeOffset(); public override void Write(Utf8JsonWriter writer, DateTimeOffset value, JsonSerializerOptions options) => writer.WriteStringValue(value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture)); } }