资讯详情

资讯详情

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

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 执行引擎ForkJoinTaskMapper、JoinTaskMapper、Join系统任务中的真实调度与汇聚逻辑。1. Fork 任务定位静态分叉与并行任务序列Fork 任务的类型声明为type : FORK_JOINFork 又称静态分叉static fork用于把多个任务序列同时并行执行分支内部甚至允许嵌套 Sub Workflow 任务。两条硬性约定Fork 之后必须紧跟 Join 任务Join 会等待被分叉的任务序列完成后再推进到下一个任务并负责收集各分叉任务的输出。这一点不是文档层面的建议引擎在调度阶段就会强校验——若 Fork 的下一个任务不是JOIN工作流会直接终止见下文第 4 节源码分析。分叉数量在定义期固定每个分支内包含哪些任务在编写工作流定义时就已确定。这与运行时根据输入决定分支数的 Dynamic Fork 形成对照本文只覆盖静态 Fork。2. 任务参数forkTasks 的嵌套列表结构Fork 任务配置中使用以下顶层参数参数类型说明必填/可选forkTasksList[List[Task]]要并行调度的任务列表的列表[[...], [...]]。外层列表的每一项代表一条将并行执行的分支内层列表是该分支内的任务配置序列。每个子列表内的任务可以串行执行也可以包含更深层的嵌套 Fork必填结构要点外层列表 并行度forkTasks有多少个元素就有多少条分支被同时调度。内层列表 分支内串行链同一条分支内的任务按顺序执行即分支1 taskA → taskB表示 taskA 完成后才执行 taskB但整条分支与分支 2、分支 3 是并行的。可嵌套内层任务可以是任意任务类型SIMPLE、START_WORKFLOW、甚至另一个FORK_JOIN从而构建出并行分支中再并行的多级 DAG 结构。在 WorkflowTask 模型中该字段声明为嵌套校验的双层列表private ListValid ListValid 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 任务文档。对于静态 ForkJoin 的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、全部分支与 JoinForkJoinTaskMapper 负责把FORK_JOIN类型的WorkflowTask映射为一批待调度任务先创建一个立即完成的 FORK 标记任务forkTask类型设为TASK_TYPE_FORKstartTime与endTime都取当前时间状态直接置为COMPLETED。它不执行任何实际逻辑仅作为并行段开始的记录点也接收 Fork 任务的inputParameters。逐条调度分支的首任务对workflowTask.getForkTasks()中的每个内层列表取该分支的第一个任务wfts.get(0)递归调用getTasksToBeScheduled生成对应任务模型。分支内后续任务在前一任务完成后由执行器按正常顺序继续调度。强校验下一个任务必须是 JOINWorkflowTask 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列表写入该任务的inputDataMapString, 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 JoinModeSYNC 同步模式与指数退避在 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 任务type为FORK_JOIN且forkTasks非空、每个内层列表至少含一个合法任务Fork 在工作流定义中的下一个任务是JOIN类型否则引擎会直接终止工作流joinOn中列出的引用名与forkTasks中实际使用的taskReferenceName完全一致需要不阻塞主流程的分支已有意从joinOn中排除并确认这些分支的失败/延迟不影响业务后续任务如需引用分叉输出使用Join 任务的引用名取值对延迟敏感的短分支场景可考虑为 Join 配置SYNC的joinMode以避免退避延迟。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),仅供参考
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →