BatchProcessor
Defined in: packages/runtime/src/batch-processor.ts:68
BatchProcessor Facade(不依赖 LokvisRuntimeImpl,通过构造注入接口)
Constructors
Section titled “Constructors”Constructor
Section titled “Constructor”new BatchProcessor(
opts):BatchProcessor
Defined in: packages/runtime/src/batch-processor.ts:76
Parameters
Section titled “Parameters”runtime
Section titled “runtime”eventBus
Section titled “eventBus”boolean
memoryGuard?
Section titled “memoryGuard?”MemoryGuard
Returns
Section titled “Returns”BatchProcessor
Methods
Section titled “Methods”enqueue()
Section titled “enqueue()”enqueue(
options):BatchJob
Defined in: packages/runtime/src/batch-processor.ts:94
Parameters
Section titled “Parameters”options
Section titled “options”Returns
Section titled “Returns”list()
Section titled “list()”list():
BatchJob[]
Defined in: packages/runtime/src/batch-processor.ts:130
Returns
Section titled “Returns”BatchJob[]
get(
jobId):BatchJob|undefined
Defined in: packages/runtime/src/batch-processor.ts:134
Parameters
Section titled “Parameters”string
Returns
Section titled “Returns”BatchJob | undefined
cancel()
Section titled “cancel()”cancel(
jobId):Promise<void>
Defined in: packages/runtime/src/batch-processor.ts:139
Parameters
Section titled “Parameters”string
Returns
Section titled “Returns”Promise<void>
dispose()
Section titled “dispose()”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 死锁。
Returns
Section titled “Returns”Promise<void>
pause()
Section titled “pause()”pause(
jobId):Promise<void>
Defined in: packages/runtime/src/batch-processor.ts:196
Parameters
Section titled “Parameters”string
Returns
Section titled “Returns”Promise<void>
resume()
Section titled “resume()”resume(
jobId):Promise<void>
Defined in: packages/runtime/src/batch-processor.ts:205
Parameters
Section titled “Parameters”string
Returns
Section titled “Returns”Promise<void>
retryFailed()
Section titled “retryFailed()”retryFailed(
jobId):Promise<void>
Defined in: packages/runtime/src/batch-processor.ts:215
Parameters
Section titled “Parameters”string
Returns
Section titled “Returns”Promise<void>
onProgress()
Section titled “onProgress()”onProgress(
jobId,handler): () =>void
Defined in: packages/runtime/src/batch-processor.ts:232
Parameters
Section titled “Parameters”string
handler
Section titled “handler”(p) => void
Returns
Section titled “Returns”() => void
waitForItemsStarted()
Section titled “waitForItemsStarted()”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 永久挂起
Parameters
Section titled “Parameters”string
number
timeoutMs?
Section titled “timeoutMs?”number = 5000
Returns
Section titled “Returns”Promise<void>
waitForCompletion()
Section titled “waitForCompletion()”waitForCompletion(
jobId,timeoutMs?):Promise<BatchJob>
Defined in: packages/runtime/src/batch-processor.ts:286
Parameters
Section titled “Parameters”string
timeoutMs?
Section titled “timeoutMs?”number = ...
Returns
Section titled “Returns”Promise<BatchJob>