Workflow.foreach()
.foreach() 方法创建循环,为数组中的每个项目执行一个步骤。它始终返回包含每次迭代输出的数组,并保持原始顺序。
使用示例使用示例的直接链接
workflow.foreach(step1, { concurrency: 2 })
参数参数的直接链接
step:
opts?:
返回值返回值的直接链接
workflow:
行为行为的直接链接
执行和等待执行和等待的直接链接
.foreach() 方法会在执行下一个步骤前处理所有项目。.foreach() 后面的步骤仅在每次迭代都完成后运行,与并发设置无关。使用 concurrency: 1(默认值)时,项目按顺序处理。使用更高并发时,项目以并行批次处理,但下一个步骤仍会等待所有批次完成。
如果你需要为每个项目运行多个操作,请使用嵌套 workflow 作为步骤。这样可将每个项目的所有操作归在一起,比链接多个 .foreach() 调用更清晰。示例请参阅在 foreach 中嵌套 workflow。
输出结构输出结构的直接链接
.foreach() 始终输出一个数组。输出数组中的每个元素对应于处理输入数组中相同索引元素的结果。
// Input: [{ value: 1 }, { value: 2 }, { value: 3 }]
// Step adds 10 to each value
// Output: [{ value: 11 }, { value: 12 }, { value: 13 }]
在 .foreach() 后使用 .then()using-then-after-foreach的直接链接
在 .foreach() 后链接 .then() 时,下一个步骤会接收整个输出数组作为输入。你可以聚合所有结果,或一起处理它们。
workflow
.foreach(processItemStep) // Output: array of processed items
.then(aggregateStep) // Input: the entire array
.commit()
在 .foreach() 后使用 .map()using-map-after-foreach的直接链接
使用 .map() 转换数组输出,然后再将其传递给下一个步骤:
workflow
.foreach(processItemStep)
.map(async ({ inputData }) => ({
total: inputData.reduce((sum, item) => sum + item.value, 0),
count: inputData.length,
}))
.then(nextStep)
.commit()
链接多个 .foreach() 调用chaining-multiple-foreach-calls的直接链接
链接 .foreach() 调用时,每个调用都在前一个步骤的数组上操作:
workflow
.foreach(stepA) // If input is [a, b, c], output is [A, B, C]
.foreach(stepB) // Operates on [A, B, C], output is [A', B', C']
.commit()
如果 .foreach() 中的步骤返回数组,输出会成为数组的数组。使用带有 .flat() 的 .map() 将其扁平化:
workflow
.foreach(chunkStep) // Output: [[chunk1, chunk2], [chunk3, chunk4]]
.map(async ({ inputData }) => inputData.flat()) // Output: [chunk1, chunk2, chunk3, chunk4]
.foreach(embedStep)
.commit()
流式传输期间的进度事件流式传输期间的进度事件的直接链接
使用 run.stream() 时,foreach 步骤会在每次迭代完成后发出 workflow-step-progress 事件。这让你无需等待整个 foreach 完成,即可跟踪实时进度。
const run = await workflow.createRun()
const stream = run.stream({ inputData })
for await (const chunk of stream) {
if (chunk.type === 'workflow-step-progress') {
console.log(`${chunk.payload.completedCount}/${chunk.payload.totalCount}`)
// e.g. "1/3", "2/3", "3/3"
}
}
每个进度事件载荷包含:
id:
completedCount:
totalCount:
currentIndex:
iterationStatus:
iterationOutput?:
恢复单次迭代恢复单次迭代的直接链接
当 .foreach() 内的步骤暂停时,每次迭代会独立暂停。将 forEachIndex 传递给 run.resume(),以使用各自的 resumeData 每次恢复一个迭代。省略 forEachIndex 会使用相同的数据恢复每个暂停的迭代。
await run.resume({
step: 'approve',
resumeData: { ok: true },
forEachIndex: 1,
})