Skip to content

BatchProcessor

Defined in: packages/runtime/src/batch-processor.ts:68

BatchProcessor Facade(不依赖 LokvisRuntimeImpl,通过构造注入接口)

new BatchProcessor(opts): BatchProcessor

Defined in: packages/runtime/src/batch-processor.ts:76

LokvisRuntime

EventBus

boolean

MemoryGuard

BatchProcessor

enqueue(options): BatchJob

Defined in: packages/runtime/src/batch-processor.ts:94

EnqueueOptions

BatchJob


list(): BatchJob[]

Defined in: packages/runtime/src/batch-processor.ts:130

BatchJob[]


get(jobId): BatchJob | undefined

Defined in: packages/runtime/src/batch-processor.ts:134

string

BatchJob | undefined


cancel(jobId): Promise<void>

Defined in: packages/runtime/src/batch-processor.ts:139

string

Promise<void>


dispose(): Promise<void>

Defined in: packages/runtime/src/batch-processor.ts:180

取消所有非终态 job + 清理 progress 订阅(W21.6 runtime.dispose 用)。

与单 job cancel 不同:不等待 runtime.cancel(wfId) 完成 —— runtime 自身的 dispose 会通过 executor.cancelAll 统一取消所有 workflow, 这里只做 batch 层面的状态标记 + 订阅清理,避免双 await 死锁。

Promise<void>


pause(jobId): Promise<void>

Defined in: packages/runtime/src/batch-processor.ts:196

string

Promise<void>


resume(jobId): Promise<void>

Defined in: packages/runtime/src/batch-processor.ts:205

string

Promise<void>


retryFailed(jobId): Promise<void>

Defined in: packages/runtime/src/batch-processor.ts:215

string

Promise<void>


onProgress(jobId, handler): () => void

Defined in: packages/runtime/src/batch-processor.ts:232

string

(p) => void

() => void


waitForItemsStarted(jobId, count, timeoutMs?): Promise<void>

Defined in: packages/runtime/src/batch-processor.ts:248

等待 job 中至少 count 项进入 processing 状态(TD-2.2 长期方案)。

替代测试中固定 setTimeout 赌注:基于“当前状态快照 + batch:item:started 事件订阅”双重判定,确定性等待而非时间赌注。

  • 先快照当前 processing 数,若已 >= count 立即 resolve(schedule 循环内 同步设置 item.status=‘processing’ + 发 batch:item:started,enqueue 返回时首批项通常已进入 processing)
  • 否则订阅 batch:item:started,累计到 count 时 resolve(兜底异步场景)
  • 超时(默认 5s)reject,避免坏 job 永久挂起

string

number

number = 5000

Promise<void>


waitForCompletion(jobId, timeoutMs?): Promise<BatchJob>

Defined in: packages/runtime/src/batch-processor.ts:286

string

number = ...

Promise<BatchJob>