IMdpSourcePullExecutor.cs 9.7 KB

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