Conductor 工作流编排详解:用 FORK_JOIN 任务实现任务序列并行执行
2026/9/10 5:32:06 网站建设 项目流程

Conductor 工作流编排详解:用 FORK_JOIN 任务实现任务序列并行执行

【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor

本篇围绕 Conductor 的 Fork 任务(FORK_JOIN)展开:讲解其forkTasks参数结构与 JSON 配置规范、与 Join 任务的强制配对机制、无输出的聚合语义,以及一个邮件/SMS/HTTP 三路并行通知的完整示例。读完并结合仓库源码后,你可以直接写出可运行的并行编排定义,并理解 Fork/Join 在 Conductor 执行引擎(ForkJoinTaskMapperJoinTaskMapperJoin系统任务)中的真实调度与汇聚逻辑。

1. Fork 任务定位:静态分叉与并行任务序列

Fork 任务的类型声明为:

"type" : "FORK_JOIN"

Fork 又称静态分叉(static fork),用于把多个任务序列同时并行执行,分支内部甚至允许嵌套 Sub Workflow 任务。

两条硬性约定:

  1. Fork 之后必须紧跟 Join 任务:Join 会等待被分叉的任务序列完成后再推进到下一个任务,并负责收集各分叉任务的输出。这一点不是文档层面的"建议",引擎在调度阶段就会强校验——若 Fork 的下一个任务不是JOIN,工作流会直接终止(见下文第 4 节源码分析)。
  2. 分叉数量在定义期固定:每个分支内包含哪些任务,在编写工作流定义时就已确定。这与运行时根据输入决定分支数的 Dynamic Fork 形成对照;本文只覆盖静态 Fork。

2. 任务参数:forkTasks 的嵌套列表结构

Fork 任务配置中使用以下顶层参数:

参数类型说明必填/可选
forkTasksList[List[Task]]要并行调度的任务列表的列表([[...], [...]])。外层列表的每一项代表一条将并行执行的分支;内层列表是该分支内的任务配置序列。每个子列表内的任务可以串行执行,也可以包含更深层的嵌套 Fork必填

结构要点:

  • 外层列表 = 并行度forkTasks有多少个元素,就有多少条分支被同时调度。
  • 内层列表 = 分支内串行链:同一条分支内的任务按顺序执行,即分支1 = taskA → taskB表示 taskA 完成后才执行 taskB,但整条分支与分支 2、分支 3 是并行的。
  • 可嵌套:内层任务可以是任意任务类型(SIMPLESTART_WORKFLOW、甚至另一个FORK_JOIN),从而构建出"并行分支中再并行"的多级 DAG 结构。

在 WorkflowTask 模型中,该字段声明为嵌套校验的双层列表:

private List<@Valid List<@Valid WorkflowTask>> forkTasks = new LinkedList<>();

forkTasks的默认值是空列表,元数据校验(@Valid)会递归校验每一条分支内的每个任务对象。

3. JSON 配置规范

一个最小的 Fork 任务配置:

{ "name": "fork", "taskReferenceName": "fork_ref", "inputParameters": {}, "type": "FORK_JOIN", "forkTasks": [ [ // fork branch { // task configuration }, { // task configuration } ], [ // another fork branch { // task configuration }, { // task configuration } ] ] }

Fork 必须与 Join 配对使用,完整的一对配置见 Join 任务文档。对于静态 Fork,Join 的joinOn参数(可选)指定需要等待完成的任务引用名列表;若不指定,Join 将不再等待任何分叉任务即进入下一步。

4. 输出语义:Fork 无输出,由 Join 聚合

Fork 任务本身没有输出。它只是并行结构的"发射点";输出聚合由随后的JOIN任务完成——Join 的输出是一个 Map,键为被 join 的任务引用名(taskReferenceName),值为对应任务的输出:

{ "taskReferenceName": { "outputKey": "outputValue" }, "anotherTaskReferenceName": { "outputKey": "outputValue" } }

因此后续任务引用分叉结果时,应通过 Join 任务引用名取数,而不是 Fork 任务引用名。

5. 完整示例:邮件、SMS、HTTP 三路并行通知

场景:工作流需要同时发出三种通知——email、SMS 和 HTTP。三者互相不依赖,天然适合 Fork 并行。执行拓扑如下:

对应的 Fork + Join JSON 配置(三条分支各自包含"生成通知负载 → 发送通知"两个串行任务):

[ { "name": "fork_join", "taskReferenceName": "my_fork_join_ref", "type": "FORK_JOIN", "forkTasks": [ [ { "name": "process_notification_payload", "taskReferenceName": "process_notification_payload_email", "type": "SIMPLE" }, { "name": "email_notification", "taskReferenceName": "email_notification_ref", "type": "SIMPLE" } ], [ { "name": "process_notification_payload", "taskReferenceName": "process_notification_payload_sms", "type": "SIMPLE" }, { "name": "sms_notification", "taskReferenceName": "sms_notification_ref", "type": "SIMPLE" } ], [ { "name": "process_notification_payload", "taskReferenceName": "process_notification_payload_http", "type": "SIMPLE" }, { "name": "http_notification", "taskReferenceName": "http_notification_ref", "type": "SIMPLE" } ] ] }, { "name": "notification_join", "taskReferenceName": "notification_join_ref", "type": "JOIN", "joinOn": [ "email_notification_ref", "sms_notification_ref" ] } ]

注意示例中joinOn只列出了 email 与 SMS 两条分支的末端任务:这意味着 Join 在等待这两条分支完成后就推进工作流,而http_notification_ref分支可以"不阻塞主流程"地继续执行。这是一种典型的"尽力而为分支"设计——把不可靠/允许延迟完成的通道排除在joinOn之外。Join 侧的更多细节(等待语义、输出聚合)见 Join 任务文档。

6. 源码剖析:Fork/Join 在引擎中如何被调度

以下结合当前仓库源码说明 Fork 配置背后的执行机制,便于理解参数为何如此设计。

6.1 ForkJoinTaskMapper:一次调度同时生成 Fork、全部分支与 Join

ForkJoinTaskMapper 负责把FORK_JOIN类型的WorkflowTask映射为一批待调度任务:

  1. 先创建一个立即完成的 FORK 标记任务forkTask类型设为TASK_TYPE_FORKstartTimeendTime都取当前时间,状态直接置为COMPLETED。它不执行任何实际逻辑,仅作为并行段开始的记录点(也接收 Fork 任务的inputParameters)。
  2. 逐条调度分支的首任务:对workflowTask.getForkTasks()中的每个内层列表,取该分支的第一个任务(wfts.get(0)),递归调用getTasksToBeScheduled生成对应任务模型。分支内后续任务在前一任务完成后由执行器按正常顺序继续调度。
  3. 强校验下一个任务必须是 JOIN
WorkflowTask joinWorkflowTask = workflowModel.getWorkflowDefinition().getNextTask(workflowTask.getTaskReferenceName()); if (joinWorkflowTask == null || !joinWorkflowTask.getType().equals(TaskType.JOIN.name())) { throw new TerminateWorkflowException( "Fork task definition is not followed by a join task. Check the blueprint"); }

若 Fork 后面不是JOIN任务(或根本没有下一个任务),直接抛出TerminateWorkflowException终止工作流。这解释了文档中"Fork 后必须跟 Join"是引擎级硬约束。 4.Join 任务一并调度:校验通过后,Join 任务随 Fork 与所有分支首任务一起被创建,初始状态为IN_PROGRESS(见 JoinTaskMapper),并把定义中的joinOn列表写入该任务的inputData

Map<String, Object> joinInput = new HashMap<>(); joinInput.put("joinOn", workflowTask.getJoinOn());

也就是说,Join 的等待清单在任务创建时就固化进了任务输入数据,后续执行只读这份快照。

6.2 Join 系统任务:轮询等待 + 输出聚合

Join 是一个异步系统任务(isAsync()返回true),其execute方法在每个评估周期做三件事:

  • 等待判定:遍历joinOn中每个引用名,用workflow.getTaskByRefName(ref)取任务;若某个引用尚未调度(null)则跳过,等下一轮评估。只有当所有joinOn任务都到达终态(isTerminal())时,Join 才以COMPLETED(或带错误状态)结束并返回true,否则返回false表示继续等待。
  • 输出聚合:对每个已完成且输出非空的分叉任务,执行task.addOutput(joinOnRef, forkOutput)——这正是第 4 节所述"以任务引用名为键的输出 Map"的实现位置。
  • 失败传播:若某个被等待任务失败、且该任务既非optional也不满足permissive条件(或所有任务已终态),Join 直接置为FAILED;若失败的分支只是被取消(如手动终止的子工作流),则 Join 置为CANCELED以区分"取消"与"真失败"。optional分支失败不会阻断 Join,只会让 Join 以COMPLETED_WITH_ERRORS结束。

WorkflowExecutorOps 中还实现了permissive语义(isJoinOnFailedPermissive):标记为 permissive 的分叉任务允许其失败时继续等待其余joinOn任务全部到达终态,而不是立即失败,进一步细化了并行分支的失败策略。

6.3 JoinMode:SYNC 同步模式与指数退避

在 WorkflowTask 上还有一个joinMode字段(枚举JoinMode),可取SYNC等值。它影响 Join 的评估频率而非等待语义:

if (workflowTask != null && WorkflowTask.JoinMode.SYNC == workflowTask.getJoinMode()) { // Synchronous mode: evaluate immediately every time (no backoff) return Optional.of(0L); }
  • SYNC模式:Join 每个评估周期都立即检查,分叉/汇聚延迟最小,适合分支执行很快、需要尽快收敛的场景。
  • 异步模式(默认):前几次轮询(不超过systemTaskPostponeThreshold)立即评估,之后按底数 1.2 的指数退避拉长评估间隔,避免长耗时分叉持续产生无意义的轮询开销。

7. 配置检查清单

  • Fork 任务typeFORK_JOIN,且forkTasks非空、每个内层列表至少含一个合法任务;
  • Fork 在工作流定义中的下一个任务JOIN类型(否则引擎会直接终止工作流);
  • joinOn中列出的引用名与forkTasks中实际使用的taskReferenceName完全一致;
  • 需要"不阻塞主流程"的分支已有意从joinOn中排除,并确认这些分支的失败/延迟不影响业务;
  • 后续任务如需引用分叉输出,使用Join 任务的引用名取值;
  • 对延迟敏感的短分支场景,可考虑为 Join 配置SYNCjoinMode以避免退避延迟。

Fork/Join 是 Conductor 中构建"并行扇出、同步收敛"结构的基本单元;当分支数量需要在运行时由输入决定时,可改用 Dynamic Fork(其 Join 会隐式等待全部分支完成,无需joinOn)。

【免费下载链接】conductorConductor is an event driven agentic workflow engine providing durable and highly resilient execution engine for applications and AI Agents项目地址: https://gitcode.com/GitHub_Trending/co/conductor

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询