跳到主要内容

Workflow.foreach()

.foreach() 方法创建循环,为数组中的每个项目执行一个步骤。它始终返回包含每次迭代输出的数组,并保持原始顺序。

使用示例
使用示例的直接链接

workflow.foreach(step1, { concurrency: 2 })

参数
参数的直接链接

step:

Step
在循环中执行的步骤实例。前一个步骤必须返回数组类型。

opts?:

object
循环的可选配置。concurrency 选项控制可以并行运行的迭代次数(默认值:1)
number

返回值
返回值的直接链接

workflow:

Workflow
用于方法链式调用的 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:

string
foreach 步骤的步骤 ID

completedCount:

number
截至目前已完成的迭代次数

totalCount:

number
迭代总次数

currentIndex:

number
刚刚完成的迭代的索引

iterationStatus:

'success' | 'failed' | 'suspended'
刚刚完成的迭代的状态

iterationOutput?:

Record<string, any>
迭代的输出(当 iterationStatus 为 'success' 时存在)

恢复单次迭代
恢复单次迭代的直接链接

.foreach() 内的步骤暂停时,每次迭代会独立暂停。将 forEachIndex 传递给 run.resume(),以使用各自的 resumeData 每次恢复一个迭代。省略 forEachIndex 会使用相同的数据恢复每个暂停的迭代。

await run.resume({
step: 'approve',
resumeData: { ok: true },
forEachIndex: 1,
})