| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610 |
- 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;
- /// <summary>
- /// 方式甲:从 mdp_source DB 源按 mdp_entity.source_table_name 增量抽数 → target_table_name(贴源)。
- /// </summary>
- public sealed class MdpDbPullExecutor : IMdpSourcePullExecutor, ITransient
- {
- public string SupportedType => "DB_SYNC";
- /// <summary>
- /// 贴源 raw_data 的时间统一按 MySQL 字面量格式输出;各域转换层用
- /// STR_TO_DATE(..., '%Y-%m-%d %H:%i:%s.%f') 解析,ISO 8601 的 T 分隔符会解析失败。
- /// </summary>
- 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<MdpPullResult> 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<SugarParameter>();
- if (NeedsTenantParameter(tenantCol, factoryCol))
- 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<string, object?>(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<MdpEntity>()
- .SetColumns(x => new MdpEntity
- {
- LastCursor = maxCursor,
- LastSyncTo = now,
- UpdateTime = now
- })
- .Where(x => x.Id == entity.Id)
- .ExecuteCommandAsync(cancellationToken);
- }
- // 零行假成功守卫:本页 0 行时,用同谓词、去作用域的探针确认源侧是否本就无数据。
- // 源侧有数据而作用域后 0 行 ⇒ 全量被租户/工厂作用域过滤掉,不得记为无保留的 SUCCESS。
- var scopeMismatch = IsScopeMismatch(
- table.Rows.Count, ctx, tenantCol, factoryCol,
- sourceHasRowsOutsideScope: ShouldProbeScope(table.Rows.Count, ctx, tenantCol, factoryCol)
- && await SourceHasRowsOutsideScopeAsync(scope, entity, ctx, isSqlServer, useKeyset, parameters, cancellationToken));
- var diagnostic = scopeMismatch
- ? $"{ScopeMismatchCode}: 源表 {entity.SourceTableName} 在 tenantId={ctx.TenantId}/factoryId={ctx.FactoryId} 作用域外仍有数据,"
- + "但作用域过滤后 0 行;请核对该实体的工厂口径(factory_id 是否被写成 tenant_id 或真实工厂号与作用域不一致)。"
- : null;
- await WriteSyncLogAsync(source, entity, ctx, startedAt, table.Rows.Count, written, diagnostic,
- scopeMismatch ? SyncStatusPartial : SyncStatusSuccess);
- 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 = (scopeMismatch ? $"{ScopeMismatchCode} window=" : "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}" : "")
- };
- }
- internal const string SyncStatusSuccess = "SUCCESS";
- /// <summary>mdp_sync_log.status 既有 enum 成员('RUNNING','SUCCESS','PARTIAL','FAILED'),无需改表。</summary>
- internal const string SyncStatusPartial = "PARTIAL";
- internal const string ScopeMismatchCode = "SCOPE_MISMATCH";
- /// <summary>
- /// 是否需要跑作用域探针。仅在 ①本页 0 行 ②确实发射了作用域谓词 ③非 OFFSET 续页 时才探。
- /// 第 ③ 条:ctx.Offset>0 说明本轮前一页已经拉满了作用域内的数据,作用域显然没有过滤掉全部,
- /// 此时的 0 行是正常翻页到底,探针(不带 OFFSET)会误报。
- /// </summary>
- internal static bool ShouldProbeScope(int rowsRead, MdpPullContext ctx, string? tenantCol, string? factoryCol)
- => rowsRead == 0 && ctx.Offset <= 0 && NeedsTenantParameter(tenantCol, factoryCol);
- /// <summary>
- /// 最终判定:本页 0 行且探针证明"去掉作用域后源侧仍有数据" ⇒ 全量被作用域过滤,记 PARTIAL。
- /// 源侧本就无数据(合法空源 / 增量无新数据)时 <paramref name="sourceHasRowsOutsideScope"/> 为 false ⇒ 正常 SUCCESS。
- /// </summary>
- internal static bool IsScopeMismatch(
- int rowsRead, MdpPullContext ctx, string? tenantCol, string? factoryCol, bool sourceHasRowsOutsideScope)
- => ShouldProbeScope(rowsRead, ctx, tenantCol, factoryCol) && sourceHasRowsOutsideScope;
- private static async Task<bool> SourceHasRowsOutsideScopeAsync(
- ISqlSugarClient scope, MdpEntity entity, MdpPullContext ctx,
- bool isSqlServer, bool useKeyset, List<SugarParameter> parameters,
- CancellationToken cancellationToken)
- {
- cancellationToken.ThrowIfCancellationRequested();
- try
- {
- var probeSql = BuildScopeProbeSql(entity, isSqlServer, ctx, useKeyset);
- var probe = await scope.Ado.GetDataTableAsync(probeSql, StripScopeParameters(parameters));
- return probe.Rows.Count > 0;
- }
- catch
- {
- // 探针纯诊断用途,任何失败都不得影响主流程,按"无异常"处理。
- return false;
- }
- }
- /// <summary>探针复用正式查询的参数,但必须剔除作用域参数(探针 SQL 不再引用它们)。</summary>
- internal static List<SugarParameter> StripScopeParameters(List<SugarParameter> parameters)
- => parameters
- .Where(p => !IsScopeParameter(p.ParameterName))
- .ToList();
- private static bool IsScopeParameter(string? name)
- {
- var n = (name ?? "").TrimStart('@', ':', '?');
- return string.Equals(n, "scopeTenantId", StringComparison.OrdinalIgnoreCase)
- || string.Equals(n, "scopeFactoryId", StringComparison.OrdinalIgnoreCase);
- }
- private static string RequireSourceTable(MdpEntity entity)
- {
- 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}");
- return table;
- }
- /// <summary>增量/窗口谓词(不含租户/工厂作用域),供正式查询与作用域探针复用。</summary>
- private static List<string> BuildIncrementalPredicates(MdpEntity entity, bool isSqlServer, MdpPullContext ctx, out string orderBy)
- {
- var predicates = new List<string>();
- 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";
- }
- return predicates;
- }
- private static string BuildSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null)
- {
- var table = RequireSourceTable(entity);
- var predicates = BuildIncrementalPredicates(entity, isSqlServer, ctx, out var orderBy);
- 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}";
- }
- /// <summary>keyset 谓词(不含租户/工厂作用域),供正式查询与作用域探针复用。</summary>
- private static List<string> BuildKeysetPredicates(bool isSqlServer, MdpPullContext ctx, out string orderBy)
- {
- 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<string>();
- 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)
- )
- """);
- }
- }
- orderBy = ctx.NullTimePhase ? tieExpr : $"{cursorExpr}, {tieExpr}";
- return predicates;
- }
- private static string BuildKeysetSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null)
- {
- var table = RequireSourceTable(entity);
- var predicates = BuildKeysetPredicates(isSqlServer, ctx, out var orderBy);
- AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol);
- var where = " WHERE " + string.Join(" AND ", predicates);
- 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}";
- }
- /// <summary>
- /// 作用域探针:与正式查询完全同一组增量/窗口/keyset 谓词,唯独**去掉**租户/工厂作用域谓词。
- /// 正式查询 0 行而探针有行 ⇒ 差异只可能来自作用域过滤,可据此判定"全量被作用域过滤掉"。
- /// 只做存在性判断(TOP 1 / LIMIT 1),不取数据。
- /// </summary>
- internal static string BuildScopeProbeSql(MdpEntity entity, bool isSqlServer, MdpPullContext ctx, bool useKeyset)
- {
- var table = RequireSourceTable(entity);
- var predicates = useKeyset
- ? BuildKeysetPredicates(isSqlServer, ctx, out _)
- : BuildIncrementalPredicates(entity, isSqlServer, ctx, out _);
- var where = predicates.Count > 0 ? " WHERE " + string.Join(" AND ", predicates) : "";
- return isSqlServer
- ? $"SELECT TOP 1 1 AS probe FROM {table}{where}"
- : $"SELECT 1 AS probe FROM {table}{where} LIMIT 1";
- }
- private static void AddKeysetParameters(List<SugarParameter> 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));
- }
- }
- /// <summary>
- /// 租户/工厂作用域谓词。工厂维度走 <see cref="MdpFactoryScope"/> 的归一语义
- /// (factory_id 为 NULL / <=0 / 等于 tenant_id 一律视作"无真实工厂"→默认工厂 1),
- /// 与 <see cref="MdpStagingWriter"/> 写入侧完全一致;真实工厂号仍逐值隔离。
- /// </summary>
- internal static void AppendScopePredicates(List<string> predicates, bool isSqlServer, string? tenantCol, string? factoryCol)
- {
- if (!string.IsNullOrWhiteSpace(tenantCol))
- predicates.Add($"{QuoteIdent(tenantCol, isSqlServer)} = @scopeTenantId");
- if (!string.IsNullOrWhiteSpace(factoryCol))
- {
- predicates.Add(MdpFactoryScope.SqlPredicate(
- QuoteIdent(factoryCol, isSqlServer), "@scopeTenantId", "@scopeFactoryId"));
- }
- }
- /// <summary>
- /// 工厂谓词引用 @scopeTenantId,故只要发射了任一作用域谓词就必须绑定该参数
- /// (源表只有 factory_id 而无 tenant_id 时,旧逻辑会漏绑导致 SQL 参数缺失)。
- /// </summary>
- internal static bool NeedsTenantParameter(string? tenantCol, string? factoryCol)
- => !string.IsNullOrWhiteSpace(tenantCol) || !string.IsNullOrWhiteSpace(factoryCol);
- private static string QuoteIdent(string name, bool isSqlServer) => isSqlServer ? name : $"`{name}`";
- private static string? FindColumn(IReadOnlyCollection<string> 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<IReadOnlyCollection<string>> 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<string>();
- if (isSqlServer)
- {
- var rows = await scope.Ado.SqlQueryAsync<string>(
- "SELECT name FROM sys.columns WHERE object_id = OBJECT_ID(@t)",
- new SugarParameter("@t", simple));
- return rows ?? new List<string>();
- }
- var mysql = await scope.Ado.SqlQueryAsync<string>(
- """
- SELECT COLUMN_NAME FROM information_schema.COLUMNS
- WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
- """,
- new SugarParameter("@t", simple));
- return mysql ?? new List<string>();
- }
- catch
- {
- return Array.Empty<string>();
- }
- }
- 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<string, object?> 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<KeysetCursorDto>(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<string, object?> 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, string? status = null)
- {
- 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", ResolveSyncStatus(status, error)),
- new SugarParameter("@error", error));
- }
- catch
- {
- // 日志表结构可能与种子不一致时不阻断主流程
- }
- }
- /// <summary>
- /// 显式 status 优先(用于 PARTIAL 等非成功但非失败的场景);未指定时沿用"有 error 即 FAILED"。
- /// 注意:SCOPE_MISMATCH 走 status=PARTIAL + error_message 诊断,而不是 FAILED——
- /// 它是"读到 0 行且可疑"的告警,不应把正常调度打成失败。
- /// </summary>
- internal static string ResolveSyncStatus(string? status, string? error)
- {
- if (!string.IsNullOrWhiteSpace(status)) return status!;
- return string.IsNullOrEmpty(error) ? SyncStatusSuccess : "FAILED";
- }
- private sealed class MysqlLiteralDateTimeConverter : JsonConverter<DateTime>
- {
- 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<DateTimeOffset>
- {
- 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));
- }
- }
|