Workflow.foreach()
La méthode .foreach() crée une boucle qui exécute une étape pour chaque élément d’un tableau. Elle renvoie toujours un tableau contenant la sortie de chaque itération, en préservant l’ordre d’origine.
Exemple d’utilisationLien direct vers Exemple d’utilisation
workflow.foreach(step1, { concurrency: 2 })
ParamètresLien direct vers Paramètres
step:
opts?:
Valeur renvoyéeLien direct vers Valeur renvoyée
workflow:
ComportementLien direct vers Comportement
Exécution et attenteLien direct vers Exécution et attente
La méthode .foreach() traite tous les éléments avant l’exécution de l’étape suivante. L’étape qui suit .foreach() ne s’exécute qu’après la fin de chaque itération, quels que soient les paramètres de concurrence. Avec concurrency: 1 (valeur par défaut), les éléments sont traités séquentiellement. Avec une concurrence plus élevée, les éléments sont traités par lots en parallèle, mais l’étape suivante attend toujours la fin de tous les lots.
Si vous devez exécuter plusieurs opérations par élément, utilisez un Workflow imbriqué comme étape. Toutes les opérations d’un élément restent ainsi regroupées, ce qui est plus propre que d’enchaîner plusieurs appels .foreach(). Consultez Workflows imbriqués dans foreach pour des exemples.
Structure de sortieLien direct vers Structure de sortie
.foreach() produit toujours un tableau. Chaque élément du tableau de sortie correspond au résultat du traitement de l’élément situé au même index dans le tableau d’entrée.
// Input: [{ value: 1 }, { value: 2 }, { value: 3 }]
// Step adds 10 to each value
// Output: [{ value: 11 }, { value: 12 }, { value: 13 }]
Utiliser .then() après .foreach()Lien direct vers using-then-after-foreach
Lorsque vous enchaînez .then() après .foreach(), l’étape suivante reçoit l’intégralité du tableau de sortie comme entrée. Vous pouvez agréger ou traiter tous les résultats ensemble.
workflow
.foreach(processItemStep) // Output: array of processed items
.then(aggregateStep) // Input: the entire array
.commit()
Utiliser .map() après .foreach()Lien direct vers using-map-after-foreach
Utilisez .map() pour transformer le tableau de sortie avant de le transmettre à l’étape suivante :
workflow
.foreach(processItemStep)
.map(async ({ inputData }) => ({
total: inputData.reduce((sum, item) => sum + item.value, 0),
count: inputData.length,
}))
.then(nextStep)
.commit()
Enchaîner plusieurs appels .foreach()Lien direct vers chaining-multiple-foreach-calls
Lorsque vous enchaînez des appels .foreach(), chacun opère sur le tableau provenant de l’étape précédente :
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()
Si une étape dans .foreach() renvoie un tableau, la sortie devient un tableau de tableaux. Utilisez .map() avec .flat() pour l’aplatir :
workflow
.foreach(chunkStep) // Output: [[chunk1, chunk2], [chunk3, chunk4]]
.map(async ({ inputData }) => inputData.flat()) // Output: [chunk1, chunk2, chunk3, chunk4]
.foreach(embedStep)
.commit()
Événements de progression pendant le streamingLien direct vers Événements de progression pendant le streaming
Lorsque vous utilisez run.stream(), les étapes foreach émettent un événement workflow-step-progress après la fin de chaque itération. Vous pouvez ainsi suivre la progression en temps réel sans attendre la fin de l’ensemble du 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"
}
}
La charge utile de chaque événement de progression contient :
id:
completedCount:
totalCount:
currentIndex:
iterationStatus:
iterationOutput?:
Reprendre une seule itérationLien direct vers Reprendre une seule itération
Lorsqu’une étape dans .foreach() se suspend, chaque itération se suspend indépendamment. Transmettez forEachIndex à run.resume() pour reprendre une itération à la fois avec ses propres resumeData. L’omission de forEachIndex reprend chaque itération suspendue avec les mêmes données.
await run.resume({
step: 'approve',
resumeData: { ok: true },
forEachIndex: 1,
})