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));
}
}