using Admin.NET.Plugin.AiDOP.Entity.DataPlatform; namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors; /// 统一入站抽数执行器(方式甲 DB / 方式乙 API)。 public interface IMdpSourcePullExecutor { /// DB_SYNC / API_PULL string SupportedType { get; } Task PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default); } public sealed class MdpPullContext { public long TenantId { get; set; } public string BatchId { get; set; } = ""; /// 强制全量(忽略水位与滚动下界)。 public bool FullRefresh { get; set; } /// 工厂作用域;>0 时入站写入与源过滤按工厂对齐。 public long FactoryId { get; set; } /// 为 true 时禁止把无租户源行归给当前上下文租户;源租户不匹配则跳过。 public bool RequireMatchingSourceTenant { get; set; } /// 关联 mdp_sync_task.task_code,用于读取调度上的同步窗口。 public string? TaskCode { get; set; } /// FULL / INCR / ROLLING(可由调度表覆盖)。 public string? SyncWindowType { get; set; } /// 如 7d / 24h / 2026-01-01。 public string? SyncWindowValue { get; set; } /// ROLLING 解析后的下界(含)。 public DateTime? WindowFrom { get; set; } /// 全量分页偏移(仅 FullRefresh / 无 incr 时由 PullAll 推进)。 public int Offset { get; set; } /// 启用复合 keyset 游标(库存冷链等场景);未设置时保持既有 OFFSET/单列游标行为。 public bool UseKeysetCursor { get; set; } /// 时间/主游标列,如 UpdateTime / CreateTime。 public string? CursorColumn { get; set; } /// 唯一 tie-breaker 列,如 RecID。 public string? TieBreakerColumn { get; set; } /// 本页起始游标时间(ISO/字面量)。 public string? CursorValue { get; set; } /// 本页起始 tie-breaker。 public string? TieBreakerValue { get; set; } /// 冻结上界时间。 public string? UpperCursorValue { get; set; } /// 冻结上界 RecID。 public string? UpperTieBreakerValue { get; set; } /// bootstrap 下界(含),如近 12 个月。 public DateTime? BootstrapFrom { get; set; } /// 为 true 时本页不写实体 LastCursor,由编排在整轮成功后持久化。 public bool DeferCursorPersist { get; set; } /// LocationDetail 等 NULL 时间段:仅拉 CursorColumn IS NULL 且按 TieBreaker 推进。 public bool NullTimePhase { get; set; } /// 为 true 时不从实体 LastCursor 回填起始 keyset(bootstrap/reconcile 首刷用)。 public bool SkipPersistedKeysetCursor { get; set; } } public sealed class MdpPullResult { public int RowsPulled { get; set; } public int RowsWritten { get; set; } public string? NewCursor { get; set; } public string? Message { get; set; } } /// 按实体/源类型选择 DB 或 API 执行器。 public sealed class MdpSourcePullDispatcher : ITransient { private readonly MdpDbPullExecutor _dbExecutor; private readonly MdpApiPullExecutor _apiExecutor; private readonly ISqlSugarClient _db; public MdpSourcePullDispatcher( MdpDbPullExecutor dbExecutor, MdpApiPullExecutor apiExecutor, ISqlSugarClient db) { _dbExecutor = dbExecutor; _apiExecutor = apiExecutor; _db = db; } public async Task PullByEntityCodeAsync(string entityCode, MdpPullContext ctx, CancellationToken cancellationToken = default) { var entity = await _db.Queryable() .Where(x => x.EntityCode == entityCode && x.Status == 1) .FirstAsync(cancellationToken) ?? throw new InvalidOperationException($"mdp_entity 未找到启用实体:{entityCode}"); var source = await _db.Queryable() .Where(x => x.Id == entity.SourceId && x.Status == 1) .FirstAsync(cancellationToken) ?? throw new InvalidOperationException($"mdp_source id={entity.SourceId} 未找到或未启用"); if (string.IsNullOrWhiteSpace(ctx.BatchId)) ctx.BatchId = $"MDP_PULL_{DateTime.Now:yyyyMMddHHmmss}"; await ApplyScheduleWindowAsync(ctx, cancellationToken); MdpSyncWindowResolver.Apply(ctx); if (string.Equals(source.SourceType, "FILE_EXCEL", StringComparison.OrdinalIgnoreCase)) throw new InvalidOperationException("文件源不支持 Pull,请使用文件导入接口"); IMdpSourcePullExecutor executor = !string.IsNullOrWhiteSpace(entity.SourceApiPath) || string.Equals(source.SourceType, "API", StringComparison.OrdinalIgnoreCase) ? _apiExecutor : _dbExecutor; return await executor.PullAsync(source, entity, ctx, cancellationToken); } /// /// 分页抽尽:FullRefresh 用 OFFSET;增量依赖实体 LastCursor 在页间推进; /// UseKeysetCursor 时用复合 keyset,达 maxPages 且末页仍满批则抛错。 /// public async Task PullAllByEntityCodeAsync( string entityCode, MdpPullContext ctx, CancellationToken cancellationToken = default, int maxPages = 500) { var totalPulled = 0; var totalWritten = 0; string? lastCursor = null; string? lastMessage = null; ctx.Offset = 0; if (ctx.UseKeysetCursor && !ctx.SkipPersistedKeysetCursor && string.IsNullOrWhiteSpace(ctx.CursorValue) && string.IsNullOrWhiteSpace(ctx.TieBreakerValue) && !ctx.NullTimePhase) { var entity0 = await _db.Queryable() .Where(x => x.EntityCode == entityCode && x.Status == 1) .FirstAsync(cancellationToken); if (entity0 != null && MdpDbPullExecutor.TryDecodeKeysetCursor(entity0.LastCursor, out var c, out var t)) { ctx.CursorValue = c; ctx.TieBreakerValue = t; } } for (var page = 0; page < maxPages; page++) { cancellationToken.ThrowIfCancellationRequested(); var pageResult = await PullByEntityCodeAsync(entityCode, ctx, cancellationToken); totalPulled += pageResult.RowsPulled; totalWritten += pageResult.RowsWritten; lastCursor = pageResult.NewCursor ?? lastCursor; lastMessage = pageResult.Message; var entity = await _db.Queryable() .Where(x => x.EntityCode == entityCode && x.Status == 1) .FirstAsync(cancellationToken); var batchSize = entity?.BatchSize > 0 ? entity.BatchSize : 1000; if (pageResult.RowsPulled <= 0 || pageResult.RowsPulled < batchSize) break; if (ctx.UseKeysetCursor) { if (page == maxPages - 1) throw new InvalidOperationException( $"实体 {entityCode} keyset 拉取达到 maxPages={maxPages} 且末页仍满批,禁止静默截断"); // 下一页游标已由执行器写回 ctx.CursorValue/TieBreakerValue continue; } if (ctx.FullRefresh) ctx.Offset += pageResult.RowsPulled; // 增量:实体 LastCursor 已在执行器内更新,下一页自动收窄 } if (ctx.UseKeysetCursor && ctx.DeferCursorPersist && !string.IsNullOrEmpty(lastCursor)) { var entity = await _db.Queryable() .Where(x => x.EntityCode == entityCode && x.Status == 1) .FirstAsync(cancellationToken); if (entity != null) { var now = DateTime.Now; await _db.Updateable() .SetColumns(x => new MdpEntity { LastCursor = lastCursor, LastSyncTo = now, UpdateTime = now }) .Where(x => x.Id == entity.Id) .ExecuteCommandAsync(cancellationToken); } } return new MdpPullResult { RowsPulled = totalPulled, RowsWritten = totalWritten, NewCursor = lastCursor, Message = $"OK pages pulled={totalPulled} written={totalWritten}; {lastMessage}" }; } /// 若上下文带 TaskCode 且未显式指定窗口,则从 mdp_sync_task_schedule 读取。 private async Task ApplyScheduleWindowAsync(MdpPullContext ctx, CancellationToken cancellationToken) { if (ctx.FullRefresh) return; if (string.IsNullOrWhiteSpace(ctx.TaskCode)) return; if (!string.IsNullOrWhiteSpace(ctx.SyncWindowType)) return; var schedules = await _db.Queryable() .Where(x => x.TaskCode == ctx.TaskCode && (x.TenantId == ctx.TenantId || x.TenantId == 0)) .OrderByDescending(x => x.TenantId) .Take(1) .ToListAsync(cancellationToken); var schedule = schedules.FirstOrDefault(); if (schedule == null) return; ctx.SyncWindowType = schedule.SyncWindowType; ctx.SyncWindowValue = schedule.SyncWindowValue; } }