| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224 |
- using Admin.NET.Plugin.AiDOP.Entity.DataPlatform;
- namespace Admin.NET.Plugin.AiDOP.DataPlatform.Executors;
- /// <summary>统一入站抽数执行器(方式甲 DB / 方式乙 API)。</summary>
- public interface IMdpSourcePullExecutor
- {
- /// <summary>DB_SYNC / API_PULL</summary>
- string SupportedType { get; }
- Task<MdpPullResult> PullAsync(MdpSource source, MdpEntity entity, MdpPullContext ctx, CancellationToken cancellationToken = default);
- }
- public sealed class MdpPullContext
- {
- public long TenantId { get; set; }
- public string BatchId { get; set; } = "";
- /// <summary>强制全量(忽略水位与滚动下界)。</summary>
- public bool FullRefresh { get; set; }
- /// <summary>工厂作用域;>0 时入站写入与源过滤按工厂对齐。</summary>
- public long FactoryId { get; set; }
- /// <summary>为 true 时禁止把无租户源行归给当前上下文租户;源租户不匹配则跳过。</summary>
- public bool RequireMatchingSourceTenant { get; set; }
- /// <summary>关联 mdp_sync_task.task_code,用于读取调度上的同步窗口。</summary>
- public string? TaskCode { get; set; }
- /// <summary>FULL / INCR / ROLLING(可由调度表覆盖)。</summary>
- public string? SyncWindowType { get; set; }
- /// <summary>如 7d / 24h / 2026-01-01。</summary>
- public string? SyncWindowValue { get; set; }
- /// <summary>ROLLING 解析后的下界(含)。</summary>
- public DateTime? WindowFrom { get; set; }
- /// <summary>全量分页偏移(仅 FullRefresh / 无 incr 时由 PullAll 推进)。</summary>
- public int Offset { get; set; }
- /// <summary>启用复合 keyset 游标(库存冷链等场景);未设置时保持既有 OFFSET/单列游标行为。</summary>
- public bool UseKeysetCursor { get; set; }
- /// <summary>时间/主游标列,如 UpdateTime / CreateTime。</summary>
- public string? CursorColumn { get; set; }
- /// <summary>唯一 tie-breaker 列,如 RecID。</summary>
- public string? TieBreakerColumn { get; set; }
- /// <summary>本页起始游标时间(ISO/字面量)。</summary>
- public string? CursorValue { get; set; }
- /// <summary>本页起始 tie-breaker。</summary>
- public string? TieBreakerValue { get; set; }
- /// <summary>冻结上界时间。</summary>
- public string? UpperCursorValue { get; set; }
- /// <summary>冻结上界 RecID。</summary>
- public string? UpperTieBreakerValue { get; set; }
- /// <summary>bootstrap 下界(含),如近 12 个月。</summary>
- public DateTime? BootstrapFrom { get; set; }
- /// <summary>为 true 时本页不写实体 LastCursor,由编排在整轮成功后持久化。</summary>
- public bool DeferCursorPersist { get; set; }
- /// <summary>LocationDetail 等 NULL 时间段:仅拉 CursorColumn IS NULL 且按 TieBreaker 推进。</summary>
- public bool NullTimePhase { get; set; }
- /// <summary>为 true 时不从实体 LastCursor 回填起始 keyset(bootstrap/reconcile 首刷用)。</summary>
- 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; }
- }
- /// <summary>按实体/源类型选择 DB 或 API 执行器。</summary>
- 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<MdpPullResult> PullByEntityCodeAsync(string entityCode, MdpPullContext ctx, CancellationToken cancellationToken = default)
- {
- var entity = await _db.Queryable<MdpEntity>()
- .Where(x => x.EntityCode == entityCode && x.Status == 1)
- .FirstAsync(cancellationToken)
- ?? throw new InvalidOperationException($"mdp_entity 未找到启用实体:{entityCode}");
- var source = await _db.Queryable<MdpSource>()
- .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);
- }
- /// <summary>
- /// 分页抽尽:FullRefresh 用 OFFSET;增量依赖实体 LastCursor 在页间推进;
- /// UseKeysetCursor 时用复合 keyset,达 maxPages 且末页仍满批则抛错。
- /// </summary>
- public async Task<MdpPullResult> 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<MdpEntity>()
- .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<MdpEntity>()
- .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<MdpEntity>()
- .Where(x => x.EntityCode == entityCode && x.Status == 1)
- .FirstAsync(cancellationToken);
- if (entity != null)
- {
- var now = DateTime.Now;
- await _db.Updateable<MdpEntity>()
- .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}"
- };
- }
- /// <summary>若上下文带 TaskCode 且未显式指定窗口,则从 mdp_sync_task_schedule 读取。</summary>
- 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<MdpSyncTaskSchedule>()
- .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;
- }
- }
|