并行步骤可以通过 同时执行多个阻塞调用 来缩短工作流的总执行时间。
sleep、HTTP 调用和回调等阻塞调用可能需要花费时间,从毫秒到天不等。并行步骤旨在帮助处理此类并发长时间运行的操作。如果工作流必须执行多个彼此独立的阻塞调用,则使用并行分支可以同时启动调用并等待所有调用完成,从而缩短总执行时间。
例如,如果工作流必须先从多个独立系统中检索客户数据,然后才能继续,则并行分支允许并发 API 请求。如果有五个系统,每个系统都需要两秒才能响应,那么在工作流中按顺序执行这些步骤可能至少需要 10 秒;而并行执行这些步骤可能只需要两秒。
创建并行步骤
创建一个 parallel 步骤,以定义工作流中可以同时执行两个或多个步骤的部分。
YAML
- PARALLEL_STEP_NAME: parallel: exception_policy: POLICY shared: [VARIABLE_A, VARIABLE_B, ...] concurrency_limit: CONCURRENCY_LIMIT BRANCHES_OR_FOR: ...
JSON
[ { "PARALLEL_STEP_NAME": { "parallel": { "exception_policy": "POLICY", "shared": [ "VARIABLE_A", "VARIABLE_B", ... ], "concurrency_limit": "CONCURRENCY_LIMIT", "BRANCHES_OR_FOR": ... } } } ]
替换以下内容:
PARALLEL_STEP_NAME:并行步骤的名称。POLICY(可选):确定在发生未处理的异常时其他分支将采取的操作。默认政策continueAll不会导致任何进一步的操作,所有其他分支都将尝试运行。请注意,目前仅支持continueAll政策。VARIABLE_A、VARIABLE_B等:具有父级范围的可写变量列表,允许在并行步骤中进行赋值。如需了解详情,请参阅 共享变量。CONCURRENCY_LIMIT(可选):在将更多分支和迭代排队等待之前,单个工作流执行中可以同时执行的分支和迭代次数上限。这仅适用于单个parallel步骤,不会级联。必须是正整数,可以是字面量值,也可以是表达式。如需了解 详情,请参阅 并发限制。BRANCHES_OR_FOR:使用branches或for来表示以下其中一项:- 可以同时运行的分支。
- 迭代可以同时运行的循环。
请注意以下几点:
将实验性函数替换为并行步骤
如果您使用 experimental.executions.map 来支持并行工作,则可以迁移工作流以改用并行步骤,并行执行普通 for 循环。如需查看示例,请参阅
将实验性函数替换为并行步骤。
示例
这些示例演示了语法。
并行执行操作(使用分支)
如果工作流具有多个不同的步骤集,并且可以同时执行这些步骤集,那么将它们放在并行分支中可以缩短完成这些步骤所需的总时间。
在以下示例中,用户 ID 作为参数传递给工作流,并从两项不同的服务并行检索数据。 共享变量 允许在分支中写入值,并在分支 完成后读取值:
YAML
JSON
并行处理项(使用并行循环)
如果您需要对列表中的每个项执行相同的操作,则可以使用并行循环更快地完成执行。并行循环允许并行执行多个循环迭代。请注意,与 常规 for 循环不同,迭代可以 按任意顺序执行。
在以下示例中,一组用户通知在并行 for 循环中进行处理:
YAML
JSON
汇总数据(使用并行循环)
您可以处理一组项,同时收集对每个项执行的操作中的数据。例如,您可能需要跟踪已创建项的 ID,或维护包含错误的项的列表。
在以下示例中,对公共 BigQuery 数据集的 10 个单独查询各自返回文档或一组文档中的字词数。 共享变量 允许在所有迭代完成后累积和读取字词数。在计算所有文档中的字词数后,工作流会返回总数。