IMdpSourcePullExecutor.cs 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266
  1. using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
  2. using Admin.NET.Plugin.AiDOP.DataPlatform;
  3. namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
  4. /// <summary>统一入站抽数执行器(方式甲 DB / 方式乙 API)。</summary>
  5. public interface IMdpSourcePullExecutor
  6. {
  7. /// <summary>DB_SYNC / API_PULL</summary>
  8. string SupportedType { get; }
  9. Task<MdpPullResult> PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default);
  10. }
  11. public sealed class MdpPullContext
  12. {
  13. public long TenantId { get; set; }
  14. public string BatchId { get; set; } = "";
  15. /// <summary>强制全量(忽略水位与滚动下界)。</summary>
  16. public bool FullRefresh { get; set; }
  17. /// <summary>工厂作用域;&gt;0 时入站写入与源过滤按工厂对齐。</summary>
  18. public long FactoryId { get; set; }
  19. /// <summary>为 true 时禁止把无租户源行归给当前上下文租户;源租户不匹配则跳过。</summary>
  20. public bool RequireMatchingSourceTenant { get; set; }
  21. /// <summary>关联 mdp_sync_task.task_code,用于读取调度上的同步窗口。</summary>
  22. public string? TaskCode { get; set; }
  23. /// <summary>FULL / INCR / ROLLING(可由调度表覆盖)。</summary>
  24. public string? SyncWindowType { get; set; }
  25. /// <summary>如 7d / 24h / 2026-01-01。</summary>
  26. public string? SyncWindowValue { get; set; }
  27. /// <summary>ROLLING 解析后的下界(含)。</summary>
  28. public DateTime? WindowFrom { get; set; }
  29. /// <summary>全量分页偏移(仅 FullRefresh / 无 incr 时由 PullAll 推进)。</summary>
  30. public int Offset { get; set; }
  31. /// <summary>启用复合 keyset 游标(库存冷链等场景);未设置时保持既有 OFFSET/单列游标行为。</summary>
  32. public bool UseKeysetCursor { get; set; }
  33. /// <summary>时间/主游标列,如 UpdateTime / CreateTime。</summary>
  34. public string? CursorColumn { get; set; }
  35. /// <summary>唯一 tie-breaker 列,如 RecID。</summary>
  36. public string? TieBreakerColumn { get; set; }
  37. /// <summary>本页起始游标时间(ISO/字面量)。</summary>
  38. public string? CursorValue { get; set; }
  39. /// <summary>本页起始 tie-breaker。</summary>
  40. public string? TieBreakerValue { get; set; }
  41. /// <summary>冻结上界时间。</summary>
  42. public string? UpperCursorValue { get; set; }
  43. /// <summary>冻结上界 RecID。</summary>
  44. public string? UpperTieBreakerValue { get; set; }
  45. /// <summary>bootstrap 下界(含),如近 12 个月。</summary>
  46. public DateTime? BootstrapFrom { get; set; }
  47. /// <summary>为 true 时本页不写实体 LastCursor,由编排在整轮成功后持久化。</summary>
  48. public bool DeferCursorPersist { get; set; }
  49. /// <summary>LocationDetail 等 NULL 时间段:仅拉 CursorColumn IS NULL 且按 TieBreaker 推进。</summary>
  50. public bool NullTimePhase { get; set; }
  51. /// <summary>为 true 时不从实体 LastCursor 回填起始 keyset(bootstrap/reconcile 首刷用)。</summary>
  52. public bool SkipPersistedKeysetCursor { get; set; }
  53. /// <summary>
  54. /// 贴源归属租户与业务租户分离(165 库存一次贴源、再按库位投影)。
  55. /// 为 true 时不要求该租户已登记本实体的来源系统。
  56. /// </summary>
  57. public bool StagingOwnerPull { get; set; }
  58. }
  59. public sealed class MdpPullResult
  60. {
  61. public int RowsPulled { get; set; }
  62. public int RowsWritten { get; set; }
  63. public string? NewCursor { get; set; }
  64. public string? Message { get; set; }
  65. }
  66. /// <summary>按实体/源类型选择 DB 或 API 执行器。</summary>
  67. public sealed class MdpSourcePullDispatcher : ITransient
  68. {
  69. private readonly MdpDbPullExecutor _dbExecutor;
  70. private readonly MdpApiPullExecutor _apiExecutor;
  71. private readonly ISqlSugarClient _db;
  72. public MdpSourcePullDispatcher(
  73. MdpDbPullExecutor dbExecutor,
  74. MdpApiPullExecutor apiExecutor,
  75. ISqlSugarClient db)
  76. {
  77. _dbExecutor = dbExecutor;
  78. _apiExecutor = apiExecutor;
  79. _db = db;
  80. }
  81. public async Task<MdpPullResult> PullByEntityCodeAsync(string entityCode, MdpPullContext ctx, CancellationToken cancellationToken = default)
  82. {
  83. var entity = await _db.Queryable<MdpEntity>()
  84. .Where(x => x.EntityCode == entityCode && x.Status == 1)
  85. .FirstAsync(cancellationToken)
  86. ?? throw new InvalidOperationException($"mdp_entity 未找到启用实体:{entityCode}");
  87. var source = await _db.Queryable<MdpSource>()
  88. .Where(x => x.Id == entity.SourceId && x.Status == 1)
  89. .FirstAsync(cancellationToken)
  90. ?? throw new InvalidOperationException($"mdp_source id={entity.SourceId} 未找到或未启用");
  91. if (entity.TenantId == 0 && ctx.TenantId > 0 && !ctx.StagingOwnerPull)
  92. {
  93. var neutral = MdpSourceIdentity.NeutralCode(source.SystemCode, source.SourceCode);
  94. if (!string.Equals(neutral, MdpSourceIdentity.Native, StringComparison.OrdinalIgnoreCase)
  95. && !string.Equals(neutral, "AIDOPDEV_MYSQL", StringComparison.OrdinalIgnoreCase))
  96. {
  97. var allowed = await _db.Ado.GetIntAsync(
  98. """
  99. SELECT COUNT(*) FROM mdp_tenant_std_source
  100. WHERE tenant_id=@t AND source_system=@s
  101. """,
  102. new SugarParameter("@t", ctx.TenantId),
  103. new SugarParameter("@s", neutral));
  104. if (allowed == 0)
  105. throw new InvalidOperationException(
  106. $"租户 {ctx.TenantId} 未登记来源 {neutral},拒绝拉取全局实体 {entity.EntityCode}");
  107. }
  108. }
  109. if (string.IsNullOrWhiteSpace(ctx.BatchId))
  110. ctx.BatchId = $"MDP_PULL_{DateTime.Now:yyyyMMddHHmmss}";
  111. await ApplyScheduleWindowAsync(ctx, cancellationToken);
  112. MdpSyncWindowResolver.Apply(ctx);
  113. if (string.Equals(source.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase))
  114. throw new InvalidOperationException("文件源不支持 Pull,请使用文件导入接口");
  115. IMdpSourcePullExecutor executor =
  116. !string.IsNullOrWhiteSpace(entity.SourceApiPath)
  117. || string.Equals(source.SourceType, "API", StringComparison.OrdinalIgnoreCase)
  118. ? _apiExecutor
  119. : _dbExecutor;
  120. return await executor.PullAsync(source, entity, ctx, cancellationToken);
  121. }
  122. /// <summary>
  123. /// 游标没有前进时必须改走 OFFSET,否则空游标增量会反复读第一页。
  124. /// 全量本来就用 OFFSET。keyset 自己带翻页谓词,不叠加 OFFSET。
  125. /// </summary>
  126. internal static bool ShouldPageByOffset(bool fullRefresh, bool useKeyset, string? cursorAtPageStart, string? newCursor)
  127. {
  128. if (useKeyset)
  129. return false;
  130. if (fullRefresh)
  131. return true;
  132. return string.IsNullOrWhiteSpace(newCursor)
  133. || string.Equals(newCursor, cursorAtPageStart, StringComparison.Ordinal);
  134. }
  135. /// <summary>
  136. /// 分页抽尽:FullRefresh 用 OFFSET;增量依赖实体 LastCursor 在页间推进;
  137. /// 游标为空或本页没有把游标推前时退化为 OFFSET,避免空游标只抽到第一页或反复读同一页。
  138. /// UseKeysetCursor 时用复合 keyset,达 maxPages 且末页仍满批则抛错。
  139. /// </summary>
  140. public async Task<MdpPullResult> PullAllByEntityCodeAsync(
  141. string entityCode,
  142. MdpPullContext ctx,
  143. CancellationToken cancellationToken = default,
  144. int maxPages = 500)
  145. {
  146. var totalPulled = 0;
  147. var totalWritten = 0;
  148. string? lastCursor = null;
  149. string? lastMessage = null;
  150. ctx.Offset = 0;
  151. if (ctx.UseKeysetCursor
  152. && !ctx.SkipPersistedKeysetCursor
  153. && string.IsNullOrWhiteSpace(ctx.CursorValue)
  154. && string.IsNullOrWhiteSpace(ctx.TieBreakerValue)
  155. && !ctx.NullTimePhase)
  156. {
  157. var entity0 = await _db.Queryable<MdpEntity>()
  158. .Where(x => x.EntityCode == entityCode && x.Status == 1)
  159. .FirstAsync(cancellationToken);
  160. if (entity0 != null
  161. && MdpDbPullExecutor.TryDecodeKeysetCursor(entity0.LastCursor, out var c, out var t))
  162. {
  163. ctx.CursorValue = c;
  164. ctx.TieBreakerValue = t;
  165. }
  166. }
  167. for (var page = 0; page < maxPages; page++)
  168. {
  169. cancellationToken.ThrowIfCancellationRequested();
  170. var cursorAtPageStart = lastCursor;
  171. var pageResult = await PullByEntityCodeAsync(entityCode, ctx, cancellationToken);
  172. totalPulled += pageResult.RowsPulled;
  173. totalWritten += pageResult.RowsWritten;
  174. lastCursor = pageResult.NewCursor ?? lastCursor;
  175. lastMessage = pageResult.Message;
  176. var entity = await _db.Queryable<MdpEntity>()
  177. .Where(x => x.EntityCode == entityCode && x.Status == 1)
  178. .FirstAsync(cancellationToken);
  179. var batchSize = entity?.BatchSize > 0 ? entity.BatchSize : 1000;
  180. if (pageResult.RowsPulled <= 0 || pageResult.RowsPulled < batchSize)
  181. break;
  182. if (ctx.UseKeysetCursor)
  183. {
  184. if (page == maxPages - 1)
  185. throw new InvalidOperationException(
  186. $"实体 {entityCode} keyset 拉取达到 maxPages={maxPages} 且末页仍满批,禁止静默截断");
  187. // 下一页游标已由执行器写回 ctx.CursorValue/TieBreakerValue
  188. continue;
  189. }
  190. if (ShouldPageByOffset(ctx.FullRefresh, ctx.UseKeysetCursor, cursorAtPageStart, pageResult.NewCursor))
  191. ctx.Offset += pageResult.RowsPulled;
  192. // 游标已前进:实体 LastCursor 已在执行器内更新,下一页用 > cursor,Offset 保持 0
  193. }
  194. if (ctx.UseKeysetCursor && ctx.DeferCursorPersist && !string.IsNullOrEmpty(lastCursor))
  195. {
  196. var entity = await _db.Queryable<MdpEntity>()
  197. .Where(x => x.EntityCode == entityCode && x.Status == 1)
  198. .FirstAsync(cancellationToken);
  199. if (entity != null)
  200. {
  201. var now = DateTime.Now;
  202. await _db.Updateable<MdpEntity>()
  203. .SetColumns(x => new MdpEntity
  204. {
  205. LastCursor = lastCursor,
  206. LastSyncTo = now,
  207. UpdateTime = now
  208. })
  209. .Where(x => x.Id == entity.Id)
  210. .ExecuteCommandAsync(cancellationToken);
  211. }
  212. }
  213. return new MdpPullResult
  214. {
  215. RowsPulled = totalPulled,
  216. RowsWritten = totalWritten,
  217. NewCursor = lastCursor,
  218. Message = $"OK pages pulled={totalPulled} written={totalWritten}; {lastMessage}"
  219. };
  220. }
  221. /// <summary>若上下文带 TaskCode 且未显式指定窗口,则从 mdp_sync_task_schedule 读取。</summary>
  222. private async Task ApplyScheduleWindowAsync(MdpPullContext ctx, CancellationToken cancellationToken)
  223. {
  224. if (ctx.FullRefresh) return;
  225. if (string.IsNullOrWhiteSpace(ctx.TaskCode)) return;
  226. if (!string.IsNullOrWhiteSpace(ctx.SyncWindowType)) return;
  227. var schedules = await _db.Queryable<MdpSyncTaskSchedule>()
  228. .Where(x => x.TaskCode == ctx.TaskCode && (x.TenantId == ctx.TenantId || x.TenantId == 0))
  229. .OrderByDescending(x => x.TenantId)
  230. .Take(1)
  231. .ToListAsync(cancellationToken);
  232. var schedule = schedules.FirstOrDefault();
  233. if (schedule == null) return;
  234. ctx.SyncWindowType = schedule.SyncWindowType;
  235. ctx.SyncWindowValue = schedule.SyncWindowValue;
  236. }
  237. }