MdpDbPullExecutor.cs 30 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640
  1. using System.Data;
  2. using System.Globalization;
  3. using System.Text.Json;
  4. using System.Text.Json.Serialization;
  5. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  6. using SqlSugar;
  7. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  8. /// <summary>
  9. /// 方式甲:从 mdp_source DB 源按 mdp_entity.source_table_name 增量抽数 → target_table_name(贴源)。
  10. /// </summary>
  11. public sealed class MdpDbPullExecutor : IMdpSourcePullExecutor, ITransient
  12. {
  13. public string SupportedType => "DB_SYNC";
  14. /// <summary>
  15. /// 贴源 raw_data 的时间统一按 MySQL 字面量格式输出;各域转换层用
  16. /// STR_TO_DATE(..., '%Y-%m-%d %H:%i:%s.%f') 解析,ISO 8601 的 T 分隔符会解析失败。
  17. /// </summary>
  18. private static readonly JsonSerializerOptions RawDataJsonOptions = new()
  19. {
  20. Converters =
  21. {
  22. new MysqlLiteralDateTimeConverter(),
  23. new MysqlLiteralDateTimeOffsetConverter()
  24. }
  25. };
  26. private readonly MdpSourceScopeFactory _scopeFactory;
  27. private readonly ISqlSugarClient _db;
  28. private readonly MdpStagingWriter _writer;
  29. public MdpDbPullExecutor(MdpSourceScopeFactory scopeFactory, ISqlSugarClient db, MdpStagingWriter writer)
  30. {
  31. _scopeFactory = scopeFactory;
  32. _db = db;
  33. _writer = writer;
  34. }
  35. public async Task<MdpPullResult> PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default)
  36. {
  37. if (string.IsNullOrWhiteSpace(entity.SourceTableName))
  38. throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 source_table_name");
  39. if (string.IsNullOrWhiteSpace(entity.TargetTableName))
  40. throw new InvalidOperationException($"实体 {entity.EntityCode} 未配置 target_table_name");
  41. var startedAt = DateTime.Now;
  42. var scope = await _scopeFactory.GetScopeAsync(source.SourceCode, cancellationToken);
  43. var batchSize = entity.BatchSize > 0 ? entity.BatchSize : 1000;
  44. var isSqlServer = string.Equals(source.DbType, "SQLSERVER", StringComparison.OrdinalIgnoreCase)
  45. || string.Equals(source.DbType, "MSSQL", StringComparison.OrdinalIgnoreCase);
  46. var useKeyset = ctx.UseKeysetCursor
  47. && !string.IsNullOrWhiteSpace(ctx.CursorColumn)
  48. && !string.IsNullOrWhiteSpace(ctx.TieBreakerColumn);
  49. string? tenantCol = null;
  50. string? factoryCol = null;
  51. if (ctx.TenantId > 0)
  52. {
  53. var cols = await TryGetSourceColumnsAsync(scope, entity.SourceTableName!, isSqlServer, cancellationToken);
  54. tenantCol = FindColumn(cols, "tenant_id", "TenantId", "TenantID");
  55. if (ctx.FactoryId > 0)
  56. factoryCol = FindColumn(cols, "factory_id", "FactoryId");
  57. }
  58. var sql = useKeyset
  59. ? BuildKeysetSelectSql(entity, isSqlServer, batchSize, ctx, tenantCol, factoryCol)
  60. : BuildSelectSql(entity, isSqlServer, batchSize, ctx, tenantCol, factoryCol);
  61. var parameters = new List<SugarParameter>();
  62. if (NeedsTenantParameter(tenantCol, factoryCol))
  63. parameters.Add(new SugarParameter("@scopeTenantId", ctx.TenantId));
  64. if (factoryCol != null)
  65. parameters.Add(new SugarParameter("@scopeFactoryId", ctx.FactoryId));
  66. if (useKeyset)
  67. {
  68. AddKeysetParameters(parameters, ctx);
  69. }
  70. else
  71. {
  72. if (ctx.WindowFrom.HasValue)
  73. parameters.Add(new SugarParameter("@windowFrom", ctx.WindowFrom.Value));
  74. if (!ctx.FullRefresh
  75. && !string.Equals(ctx.SyncWindowType, "FULL", StringComparison.OrdinalIgnoreCase)
  76. && !string.IsNullOrWhiteSpace(entity.IncrColumn)
  77. && !string.IsNullOrWhiteSpace(entity.LastCursor)
  78. && !string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase))
  79. parameters.Add(new SugarParameter("@cursor", entity.LastCursor));
  80. // ROLLING:用窗口下界;若已有水位且水位晚于下界,则进一步用水位收窄
  81. if (!ctx.FullRefresh
  82. && string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase)
  83. && !string.IsNullOrWhiteSpace(entity.IncrColumn)
  84. && !string.IsNullOrWhiteSpace(entity.LastCursor)
  85. && ctx.WindowFrom.HasValue
  86. && DateTime.TryParse(entity.LastCursor, out var cursorDt)
  87. && cursorDt > ctx.WindowFrom.Value)
  88. parameters.Add(new SugarParameter("@cursor", entity.LastCursor));
  89. }
  90. var table = await scope.Ado.GetDataTableAsync(sql, parameters);
  91. cancellationToken.ThrowIfCancellationRequested();
  92. var written = 0;
  93. string? maxCursor = entity.LastCursor;
  94. var now = DateTime.Now;
  95. string? lastKeysetCursor = null;
  96. string? lastKeysetTie = null;
  97. var pending = new List<(IDictionary<string, object?> Row, string RawJson, string SourceRowId)>();
  98. foreach (DataRow row in table.Rows)
  99. {
  100. cancellationToken.ThrowIfCancellationRequested();
  101. var dict = new Dictionary<string, object?>(StringComparer.OrdinalIgnoreCase);
  102. foreach (DataColumn col in table.Columns)
  103. dict[col.ColumnName] = row[col] == DBNull.Value ? null : row[col];
  104. var sourceRowId = ResolveSourceRowId(dict);
  105. var rawJson = JsonSerializer.Serialize(dict, RawDataJsonOptions);
  106. // 游标推进必须留在逐行循环里:它依赖遍历顺序,挪到批量写入之后会算错。
  107. if (useKeyset)
  108. {
  109. lastKeysetCursor = FormatCursorValue(dict, ctx.CursorColumn!);
  110. lastKeysetTie = FormatCursorValue(dict, ctx.TieBreakerColumn!);
  111. if (!string.IsNullOrEmpty(lastKeysetTie))
  112. maxCursor = EncodeKeysetCursor(lastKeysetCursor, lastKeysetTie);
  113. }
  114. else if (!string.IsNullOrWhiteSpace(entity.IncrColumn) && dict.TryGetValue(entity.IncrColumn, out var incrVal) && incrVal != null)
  115. {
  116. var cursor = incrVal is DateTime dt ? dt.ToString("yyyy-MM-dd HH:mm:ss.fff") : incrVal.ToString();
  117. if (!string.IsNullOrEmpty(cursor) && (maxCursor == null || string.CompareOrdinal(cursor, maxCursor) > 0))
  118. maxCursor = cursor;
  119. }
  120. pending.Add((dict, rawJson, sourceRowId));
  121. }
  122. if (pending.Count > 0)
  123. {
  124. var results = await _writer.UpsertBatchAsync(source, entity, entity.SourceTableName!, pending, ctx);
  125. written += results.Sum();
  126. }
  127. if (useKeyset && !string.IsNullOrEmpty(lastKeysetTie))
  128. {
  129. ctx.CursorValue = lastKeysetCursor;
  130. ctx.TieBreakerValue = lastKeysetTie;
  131. }
  132. if (!ctx.DeferCursorPersist
  133. && !string.IsNullOrEmpty(maxCursor)
  134. && maxCursor != entity.LastCursor)
  135. {
  136. await _db.Updateable<MdpEntity>()
  137. .SetColumns(x => new MdpEntity
  138. {
  139. LastCursor = maxCursor,
  140. LastSyncTo = now,
  141. UpdateTime = now
  142. })
  143. .Where(x => x.Id == entity.Id)
  144. .ExecuteCommandAsync(cancellationToken);
  145. }
  146. // 零行假成功守卫:本页 0 行时,用同谓词、去作用域的探针确认源侧是否本就无数据。
  147. // 源侧有数据而作用域后 0 行 ⇒ 全量被租户/工厂作用域过滤掉,不得记为无保留的 SUCCESS。
  148. var scopeMismatch = IsScopeMismatch(
  149. table.Rows.Count, ctx, tenantCol, factoryCol,
  150. sourceHasRowsOutsideScope: ShouldProbeScope(table.Rows.Count, ctx, tenantCol, factoryCol)
  151. && await SourceHasRowsOutsideScopeAsync(scope, entity, ctx, isSqlServer, useKeyset, parameters, tenantCol, cancellationToken));
  152. var diagnostic = scopeMismatch
  153. ? $"{ScopeMismatchCode}: 源表 {entity.SourceTableName} 在 tenantId={ctx.TenantId}/factoryId={ctx.FactoryId} 作用域外仍有数据,"
  154. + "但作用域过滤后 0 行;请核对该实体的工厂口径(factory_id 是否被写成 tenant_id 或真实工厂号与作用域不一致)。"
  155. : null;
  156. await WriteSyncLogAsync(source, entity, ctx, startedAt, table.Rows.Count, written, diagnostic,
  157. scopeMismatch ? SyncStatusPartial : SyncStatusSuccess);
  158. var windowHint = useKeyset
  159. ? (ctx.NullTimePhase ? "KEYSET_NULL" : "KEYSET")
  160. : (ctx.SyncWindowType ?? (ctx.FullRefresh ? "FULL" : "INCR"));
  161. return new MdpPullResult
  162. {
  163. RowsPulled = table.Rows.Count,
  164. RowsWritten = written,
  165. NewCursor = maxCursor,
  166. Message = (scopeMismatch ? $"{ScopeMismatchCode} window=" : "OK window=") + windowHint
  167. + (ctx.WindowFrom.HasValue ? $" from={ctx.WindowFrom:yyyy-MM-dd HH:mm:ss}" : "")
  168. + (ctx.BootstrapFrom.HasValue ? $" bootstrapFrom={ctx.BootstrapFrom:yyyy-MM-dd HH:mm:ss}" : "")
  169. };
  170. }
  171. internal const string SyncStatusSuccess = "SUCCESS";
  172. /// <summary>mdp_sync_log.status 既有 enum 成员('RUNNING','SUCCESS','PARTIAL','FAILED'),无需改表。</summary>
  173. internal const string SyncStatusPartial = "PARTIAL";
  174. internal const string ScopeMismatchCode = "SCOPE_MISMATCH";
  175. /// <summary>
  176. /// 是否需要跑作用域探针。仅在 ①本页 0 行 ②确实发射了作用域谓词 ③非 OFFSET 续页 时才探。
  177. /// 第 ③ 条:ctx.Offset&gt;0 说明本轮前一页已经拉满了作用域内的数据,作用域显然没有过滤掉全部,
  178. /// 此时的 0 行是正常翻页到底,探针(不带 OFFSET)会误报。
  179. /// </summary>
  180. internal static bool ShouldProbeScope(int rowsRead, MdpPullContext ctx, string? tenantCol, string? factoryCol)
  181. => rowsRead == 0 && ctx.Offset <= 0 && NeedsTenantParameter(tenantCol, factoryCol);
  182. /// <summary>
  183. /// 最终判定:本页 0 行且探针证明"去掉作用域后源侧仍有数据" ⇒ 全量被作用域过滤,记 PARTIAL。
  184. /// 源侧本就无数据(合法空源 / 增量无新数据)时 <paramref name="sourceHasRowsOutsideScope"/> 为 false ⇒ 正常 SUCCESS。
  185. /// </summary>
  186. internal static bool IsScopeMismatch(
  187. int rowsRead, MdpPullContext ctx, string? tenantCol, string? factoryCol, bool sourceHasRowsOutsideScope)
  188. => ShouldProbeScope(rowsRead, ctx, tenantCol, factoryCol) && sourceHasRowsOutsideScope;
  189. private static async Task<bool> SourceHasRowsOutsideScopeAsync(
  190. ISqlSugarClient scope, MdpEntity entity, MdpPullContext ctx,
  191. bool isSqlServer, bool useKeyset, List<SugarParameter> parameters,
  192. string? tenantCol, CancellationToken cancellationToken)
  193. {
  194. cancellationToken.ThrowIfCancellationRequested();
  195. try
  196. {
  197. var keepTenant = !string.IsNullOrWhiteSpace(tenantCol);
  198. var probeSql = BuildScopeProbeSql(entity, isSqlServer, ctx, useKeyset, tenantCol);
  199. var probe = await scope.Ado.GetDataTableAsync(probeSql, StripScopeParameters(parameters, keepTenant));
  200. return probe.Rows.Count > 0;
  201. }
  202. catch
  203. {
  204. // 探针纯诊断用途,任何失败都不得影响主流程,按"无异常"处理。
  205. return false;
  206. }
  207. }
  208. /// <summary>
  209. /// 探针复用正式查询的参数,剔除探针 SQL 不再引用的作用域参数。
  210. /// <para><paramref name="keepTenantParameter"/> 必须与
  211. /// <see cref="BuildScopeProbeSql"/> 是否发射了租户谓词保持一致:
  212. /// 探针带租户谓词时要留下 <c>@scopeTenantId</c>,否则参数会缺失;
  213. /// 不带时要剔除,否则会多绑一个 SQL 里不存在的参数。</para>
  214. /// <para><c>@scopeFactoryId</c> 恒剔除 —— 探针的全部意义就是"去掉工厂口径再看一眼"。</para>
  215. /// </summary>
  216. internal static List<SugarParameter> StripScopeParameters(
  217. List<SugarParameter> parameters, bool keepTenantParameter = false)
  218. => parameters
  219. .Where(p => !IsStrippedScopeParameter(p.ParameterName, keepTenantParameter))
  220. .ToList();
  221. private static bool IsStrippedScopeParameter(string? name, bool keepTenantParameter)
  222. {
  223. var n = (name ?? "").TrimStart('@', ':', '?');
  224. if (string.Equals(n, "scopeFactoryId", StringComparison.OrdinalIgnoreCase)) return true;
  225. return !keepTenantParameter
  226. && string.Equals(n, "scopeTenantId", StringComparison.OrdinalIgnoreCase);
  227. }
  228. private static string RequireSourceTable(MdpEntity entity)
  229. {
  230. var table = entity.SourceTableName!.Trim();
  231. // 仅允许简单标识符/schema.table,防注入
  232. if (!System.Text.RegularExpressions.Regex.IsMatch(table, @"^[A-Za-z0-9_\.\[\]]+$"))
  233. throw new InvalidOperationException($"非法 source_table_name:{table}");
  234. return table;
  235. }
  236. /// <summary>增量/窗口谓词(不含租户/工厂作用域),供正式查询与作用域探针复用。</summary>
  237. private static List<string> BuildIncrementalPredicates(MdpEntity entity, bool isSqlServer, MdpPullContext ctx, out string orderBy)
  238. {
  239. var predicates = new List<string>();
  240. if (!string.IsNullOrWhiteSpace(entity.IncrColumn))
  241. {
  242. var incr = entity.IncrColumn.Trim();
  243. if (!System.Text.RegularExpressions.Regex.IsMatch(incr, @"^[A-Za-z0-9_]+$"))
  244. throw new InvalidOperationException($"非法 incr_column:{incr}");
  245. var incrExpr = isSqlServer ? incr : $"`{incr}`";
  246. orderBy = incrExpr;
  247. var isRolling = string.Equals(ctx.SyncWindowType, "ROLLING", StringComparison.OrdinalIgnoreCase);
  248. if (!ctx.FullRefresh && isRolling && ctx.WindowFrom.HasValue)
  249. predicates.Add($"{incrExpr} >= @windowFrom");
  250. var useCursor = !ctx.FullRefresh
  251. && !string.IsNullOrWhiteSpace(entity.LastCursor)
  252. && (!isRolling
  253. || (ctx.WindowFrom.HasValue
  254. && DateTime.TryParse(entity.LastCursor, out var cursorDt)
  255. && cursorDt > ctx.WindowFrom.Value));
  256. if (useCursor)
  257. predicates.Add($"{incrExpr} > @cursor");
  258. }
  259. else
  260. {
  261. orderBy = isSqlServer ? "(SELECT NULL)" : "1";
  262. }
  263. return predicates;
  264. }
  265. private static string BuildSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null)
  266. {
  267. var table = RequireSourceTable(entity);
  268. var predicates = BuildIncrementalPredicates(entity, isSqlServer, ctx, out var orderBy);
  269. AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol);
  270. var where = predicates.Count > 0 ? " WHERE " + string.Join(" AND ", predicates) : "";
  271. var offset = ctx.Offset > 0 ? ctx.Offset : 0;
  272. if (isSqlServer)
  273. {
  274. // SQL Server:OFFSET 需 ORDER BY;无 incr 时用稳定键兜底
  275. if (string.IsNullOrWhiteSpace(entity.IncrColumn))
  276. orderBy = "(SELECT NULL)";
  277. return $"SELECT * FROM {table}{where} ORDER BY {orderBy} OFFSET {offset} ROWS FETCH NEXT {batchSize} ROWS ONLY";
  278. }
  279. return offset > 0
  280. ? $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize} OFFSET {offset}"
  281. : $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize}";
  282. }
  283. /// <summary>keyset 谓词(不含租户/工厂作用域),供正式查询与作用域探针复用。</summary>
  284. private static List<string> BuildKeysetPredicates(bool isSqlServer, MdpPullContext ctx, out string orderBy)
  285. {
  286. var cursorCol = RequireIdent(ctx.CursorColumn, "CursorColumn");
  287. var tieCol = RequireIdent(ctx.TieBreakerColumn, "TieBreakerColumn");
  288. var cursorExpr = isSqlServer ? cursorCol : $"`{cursorCol}`";
  289. var tieExpr = isSqlServer ? tieCol : $"`{tieCol}`";
  290. var predicates = new List<string>();
  291. if (ctx.NullTimePhase)
  292. {
  293. predicates.Add($"{cursorExpr} IS NULL");
  294. if (!string.IsNullOrWhiteSpace(ctx.TieBreakerValue))
  295. predicates.Add($"{tieExpr} > @tieBreaker");
  296. }
  297. else
  298. {
  299. predicates.Add($"{cursorExpr} IS NOT NULL");
  300. if (ctx.BootstrapFrom.HasValue)
  301. predicates.Add($"{cursorExpr} >= @bootstrapFrom");
  302. // keyset lower bound
  303. predicates.Add($"""
  304. (
  305. @cursorTime IS NULL
  306. OR {cursorExpr} > @cursorTime
  307. OR ({cursorExpr} = @cursorTime AND {tieExpr} > @tieBreaker)
  308. )
  309. """);
  310. if (!string.IsNullOrWhiteSpace(ctx.UpperCursorValue))
  311. {
  312. predicates.Add($"""
  313. (
  314. {cursorExpr} < @upperTime
  315. OR ({cursorExpr} = @upperTime AND {tieExpr} <= @upperTie)
  316. )
  317. """);
  318. }
  319. }
  320. orderBy = ctx.NullTimePhase ? tieExpr : $"{cursorExpr}, {tieExpr}";
  321. return predicates;
  322. }
  323. private static string BuildKeysetSelectSql(MdpEntity entity, bool isSqlServer, int batchSize, MdpPullContext ctx, string? tenantCol = null, string? factoryCol = null)
  324. {
  325. var table = RequireSourceTable(entity);
  326. var predicates = BuildKeysetPredicates(isSqlServer, ctx, out var orderBy);
  327. AppendScopePredicates(predicates, isSqlServer, tenantCol, factoryCol);
  328. var where = " WHERE " + string.Join(" AND ", predicates);
  329. if (isSqlServer)
  330. return $"SELECT * FROM {table}{where} ORDER BY {orderBy} OFFSET 0 ROWS FETCH NEXT {batchSize} ROWS ONLY";
  331. return $"SELECT * FROM {table}{where} ORDER BY {orderBy} LIMIT {batchSize}";
  332. }
  333. /// <summary>
  334. /// 作用域探针:与正式查询同一组增量/窗口/keyset 谓词,**保留租户谓词、只去掉工厂谓词**。
  335. /// 正式查询 0 行而探针有行 ⇒ 差异只可能来自工厂口径,可据此判定"本租户的数据被工厂过滤掉了"。
  336. /// 只做存在性判断(TOP 1 / LIMIT 1),不取数据。
  337. ///
  338. /// <para><b>租户谓词必须保留</b>。此前两个作用域谓词一起剥掉,探针退化成
  339. /// <c>SELECT 1 FROM 源表 LIMIT 1</c> —— 在多租户共享同一张源表的模型下,
  340. /// 它对**任何**租户都会命中别的租户的行,于是每一个合法空租户都被误报成 SCOPE_MISMATCH。
  341. /// 实测该误报在共享库上累计 291 条,且全部落在源表本就 0 行的租户上,
  342. /// 而真正被工厂谓词滤空的那两个租户反而一条都没报出来 —— 信号方向是反的。</para>
  343. ///
  344. /// <para>探针要回答的是「同一租户内,是否因工厂口径而取空」,
  345. /// 不是「整张表有没有数据」。丢掉租户维度,这个问题就问不出来了。</para>
  346. /// </summary>
  347. internal static string BuildScopeProbeSql(
  348. MdpEntity entity, bool isSqlServer, MdpPullContext ctx, bool useKeyset, string? tenantCol = null)
  349. {
  350. var table = RequireSourceTable(entity);
  351. var predicates = useKeyset
  352. ? BuildKeysetPredicates(isSqlServer, ctx, out _)
  353. : BuildIncrementalPredicates(entity, isSqlServer, ctx, out _);
  354. if (!string.IsNullOrWhiteSpace(tenantCol))
  355. predicates.Add($"{QuoteIdent(tenantCol, isSqlServer)} = @scopeTenantId");
  356. var where = predicates.Count > 0 ? " WHERE " + string.Join(" AND ", predicates) : "";
  357. return isSqlServer
  358. ? $"SELECT TOP 1 1 AS probe FROM {table}{where}"
  359. : $"SELECT 1 AS probe FROM {table}{where} LIMIT 1";
  360. }
  361. private static void AddKeysetParameters(List<SugarParameter> parameters, MdpPullContext ctx)
  362. {
  363. if (ctx.BootstrapFrom.HasValue)
  364. parameters.Add(new SugarParameter("@bootstrapFrom", ctx.BootstrapFrom.Value));
  365. if (ctx.NullTimePhase)
  366. {
  367. parameters.Add(new SugarParameter("@tieBreaker",
  368. string.IsNullOrWhiteSpace(ctx.TieBreakerValue) ? "0" : ctx.TieBreakerValue));
  369. return;
  370. }
  371. object cursorTime = string.IsNullOrWhiteSpace(ctx.CursorValue)
  372. ? DBNull.Value
  373. : (object)ctx.CursorValue!;
  374. parameters.Add(new SugarParameter("@cursorTime", cursorTime));
  375. parameters.Add(new SugarParameter("@tieBreaker",
  376. string.IsNullOrWhiteSpace(ctx.TieBreakerValue) ? "0" : ctx.TieBreakerValue));
  377. if (!string.IsNullOrWhiteSpace(ctx.UpperCursorValue))
  378. {
  379. parameters.Add(new SugarParameter("@upperTime", ctx.UpperCursorValue));
  380. parameters.Add(new SugarParameter("@upperTie",
  381. string.IsNullOrWhiteSpace(ctx.UpperTieBreakerValue) ? "0" : ctx.UpperTieBreakerValue));
  382. }
  383. }
  384. /// <summary>
  385. /// 租户/工厂作用域谓词。工厂维度走 <see cref="MdpFactoryScope"/> 的归一语义
  386. /// (factory_id 为 NULL / &lt;=0 / 等于 tenant_id 一律视作"无真实工厂"→默认工厂 1),
  387. /// 与 <see cref="MdpStagingWriter"/> 写入侧完全一致;真实工厂号仍逐值隔离。
  388. /// </summary>
  389. internal static void AppendScopePredicates(List<string> predicates, bool isSqlServer, string? tenantCol, string? factoryCol)
  390. {
  391. if (!string.IsNullOrWhiteSpace(tenantCol))
  392. predicates.Add($"{QuoteIdent(tenantCol, isSqlServer)} = @scopeTenantId");
  393. if (!string.IsNullOrWhiteSpace(factoryCol))
  394. {
  395. predicates.Add(MdpFactoryScope.SqlPredicate(
  396. QuoteIdent(factoryCol, isSqlServer), "@scopeTenantId", "@scopeFactoryId"));
  397. }
  398. }
  399. /// <summary>
  400. /// 工厂谓词引用 @scopeTenantId,故只要发射了任一作用域谓词就必须绑定该参数
  401. /// (源表只有 factory_id 而无 tenant_id 时,旧逻辑会漏绑导致 SQL 参数缺失)。
  402. /// </summary>
  403. internal static bool NeedsTenantParameter(string? tenantCol, string? factoryCol)
  404. => !string.IsNullOrWhiteSpace(tenantCol) || !string.IsNullOrWhiteSpace(factoryCol);
  405. private static string QuoteIdent(string name, bool isSqlServer) => isSqlServer ? name : $"`{name}`";
  406. private static string? FindColumn(IReadOnlyCollection<string> columns, params string[] names)
  407. {
  408. foreach (var n in names)
  409. {
  410. var hit = columns.FirstOrDefault(c => string.Equals(c, n, StringComparison.OrdinalIgnoreCase));
  411. if (hit != null) return hit;
  412. }
  413. return null;
  414. }
  415. private static async Task<IReadOnlyCollection<string>> TryGetSourceColumnsAsync(
  416. ISqlSugarClient scope, string table, bool isSqlServer, CancellationToken cancellationToken)
  417. {
  418. try
  419. {
  420. var simple = table.Contains('.') ? table.Split('.').Last().Trim('[', ']') : table.Trim('[', ']');
  421. if (!System.Text.RegularExpressions.Regex.IsMatch(simple, @"^[A-Za-z0-9_]+$"))
  422. return Array.Empty<string>();
  423. if (isSqlServer)
  424. {
  425. var rows = await scope.Ado.SqlQueryAsync<string>(
  426. "SELECT name FROM sys.columns WHERE object_id = OBJECT_ID(@t)",
  427. new SugarParameter("@t", simple));
  428. return rows ?? new List<string>();
  429. }
  430. var mysql = await scope.Ado.SqlQueryAsync<string>(
  431. """
  432. SELECT COLUMN_NAME FROM information_schema.COLUMNS
  433. WHERE TABLE_SCHEMA=DATABASE() AND TABLE_NAME=@t
  434. """,
  435. new SugarParameter("@t", simple));
  436. return mysql ?? new List<string>();
  437. }
  438. catch
  439. {
  440. return Array.Empty<string>();
  441. }
  442. }
  443. private static string RequireIdent(string? name, string label)
  444. {
  445. var v = (name ?? "").Trim();
  446. if (!System.Text.RegularExpressions.Regex.IsMatch(v, @"^[A-Za-z0-9_]+$"))
  447. throw new InvalidOperationException($"非法 {label}:{name}");
  448. return v;
  449. }
  450. private static string? FormatCursorValue(Dictionary<string, object?> dict, string column)
  451. {
  452. if (!dict.TryGetValue(column, out var raw) || raw == null || raw is DBNull)
  453. return null;
  454. return raw switch
  455. {
  456. DateTime dt => dt.ToString("yyyy-MM-dd HH:mm:ss.fff", CultureInfo.InvariantCulture),
  457. DateTimeOffset dto => dto.ToString("yyyy-MM-dd HH:mm:ss.fff", CultureInfo.InvariantCulture),
  458. _ => raw.ToString()
  459. };
  460. }
  461. internal static string EncodeKeysetCursor(string? cursor, string tieBreaker)
  462. => System.Text.Json.JsonSerializer.Serialize(new KeysetCursorDto
  463. {
  464. Cursor = cursor,
  465. TieBreaker = tieBreaker
  466. });
  467. internal static bool TryDecodeKeysetCursor(string? raw, out string? cursor, out string? tieBreaker)
  468. {
  469. cursor = null;
  470. tieBreaker = null;
  471. if (string.IsNullOrWhiteSpace(raw))
  472. return false;
  473. try
  474. {
  475. var dto = System.Text.Json.JsonSerializer.Deserialize<KeysetCursorDto>(raw);
  476. if (dto == null || string.IsNullOrWhiteSpace(dto.TieBreaker))
  477. return false;
  478. cursor = dto.Cursor;
  479. tieBreaker = dto.TieBreaker;
  480. return true;
  481. }
  482. catch
  483. {
  484. return false;
  485. }
  486. }
  487. private sealed class KeysetCursorDto
  488. {
  489. public string? Cursor { get; set; }
  490. public string? TieBreaker { get; set; }
  491. }
  492. private static string ResolveSourceRowId(Dictionary<string, object?> dict)
  493. {
  494. // InvTransHist 等表常有空字符串 ID,若优先取 ID 会导致 source_row_id="" 唯一键互相覆盖。
  495. // 跳过 null/空白后,再落到 RecID 等真实键。
  496. foreach (var key in new[]
  497. {
  498. "RecID", "RecId", "recid", "Recid",
  499. "id", "Id", "ID",
  500. "noid", "billno", "BillNo", "djbh", "Djbh"
  501. })
  502. {
  503. if (!dict.TryGetValue(key, out var v) || v is null || v is DBNull)
  504. continue;
  505. var s = v.ToString();
  506. if (!string.IsNullOrWhiteSpace(s))
  507. return s;
  508. }
  509. return Guid.NewGuid().ToString("N");
  510. }
  511. private async Task WriteSyncLogAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, DateTime startedAt, int pulled, int written, string? error, string? status = null)
  512. {
  513. try
  514. {
  515. var endedAt = DateTime.Now;
  516. var syncType = ctx.FullRefresh || string.Equals(ctx.SyncWindowType, "FULL", StringComparison.OrdinalIgnoreCase)
  517. ? "FULL"
  518. : "INCR";
  519. await _db.Ado.ExecuteCommandAsync(@"
  520. INSERT INTO mdp_sync_log
  521. (tenant_id, entity_id, source_code, entity_name, sync_batch_id, sync_type, trigger_type,
  522. sync_start, sync_end, duration_ms, rows_read, rows_written, status, error_message)
  523. VALUES
  524. (@tenantId, @entityId, @sourceCode, @entityName, @batchId, @syncType, 'AUTO',
  525. @syncStart, @syncEnd, @durationMs, @rowsRead, @rowsWritten, @status, @error)",
  526. new SugarParameter("@tenantId", ctx.TenantId),
  527. new SugarParameter("@entityId", entity.Id),
  528. new SugarParameter("@sourceCode", source.SourceCode),
  529. new SugarParameter("@entityName", string.IsNullOrWhiteSpace(entity.EntityName) ? entity.EntityCode : entity.EntityName),
  530. new SugarParameter("@batchId", ctx.BatchId),
  531. new SugarParameter("@syncType", syncType),
  532. new SugarParameter("@syncStart", startedAt),
  533. new SugarParameter("@syncEnd", endedAt),
  534. new SugarParameter("@durationMs", (int)(endedAt - startedAt).TotalMilliseconds),
  535. new SugarParameter("@rowsRead", pulled),
  536. new SugarParameter("@rowsWritten", written),
  537. new SugarParameter("@status", ResolveSyncStatus(status, error)),
  538. new SugarParameter("@error", error));
  539. }
  540. catch
  541. {
  542. // 日志表结构可能与种子不一致时不阻断主流程
  543. }
  544. }
  545. /// <summary>
  546. /// 显式 status 优先(用于 PARTIAL 等非成功但非失败的场景);未指定时沿用"有 error 即 FAILED"。
  547. /// 注意:SCOPE_MISMATCH 走 status=PARTIAL + error_message 诊断,而不是 FAILED——
  548. /// 它是"读到 0 行且可疑"的告警,不应把正常调度打成失败。
  549. /// </summary>
  550. internal static string ResolveSyncStatus(string? status, string? error)
  551. {
  552. if (!string.IsNullOrWhiteSpace(status)) return status!;
  553. return string.IsNullOrEmpty(error) ? SyncStatusSuccess : "FAILED";
  554. }
  555. private sealed class MysqlLiteralDateTimeConverter : JsonConverter<DateTime>
  556. {
  557. public override DateTime Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
  558. => reader.GetDateTime();
  559. public override void Write(Utf8JsonWriter writer, DateTime value, JsonSerializerOptions options)
  560. => writer.WriteStringValue(value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture));
  561. }
  562. private sealed class MysqlLiteralDateTimeOffsetConverter : JsonConverter<DateTimeOffset>
  563. {
  564. public override DateTimeOffset Read(ref Utf8JsonReader reader, Type typeToConvert, JsonSerializerOptions options)
  565. => reader.GetDateTimeOffset();
  566. public override void Write(Utf8JsonWriter writer, DateTimeOffset value, JsonSerializerOptions options)
  567. => writer.WriteStringValue(value.ToString("yyyy-MM-dd HH:mm:ss.ffffff", CultureInfo.InvariantCulture));
  568. }
  569. }