Conductor 事件处理实战:用 Event 任务与 Event Handler 实现工作流间事件驱动编排
发布时间:2026/9/11 13:23:30来源:尧图网络
Conductor 事件处理实战用 Event 任务与 Event Handler 实现工作流间事件驱动编排【免费下载链接】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 的官方实验教程docs/devguide/labs/eventhandlers.md完整演示如何在两个工作流之间通过事件完成发布事件 → 订阅事件 → 触发动作的闭环第一个工作流用EVENT系统任务发出事件并挂起一个WAIT任务第二个工作流由 Event Handler 自动启动再发出事件把WAIT任务标记为完成。读完你将掌握 Event 任务的sink语法、Event Handler 的完整定义结构、start_workflow与complete_task两类核心动作的用法以及事件驱动的底层处理链路并能在本地 Conductor 上原样复现整个实验。实验背景Conductor 的两套事件接口Conductor 为工作流提供事件能力共开放两个接口本次实验两者都会用到接口作用参考文档Event 任务EVENT系统任务在工作流内部发布事件到指定队列sinkEVENT 系统任务参考Event Handler事件处理器订阅事件队列并对事件执行一个或多个动作启动工作流、完成任务等Event Handler 配置参考简单来说EVENT任务负责发Event Handler 负责收并响应。只有两者配合事件驱动的编排才能闭环——只发事件、不定义 Event Handler事件会成功发出但不会触发任何后续动作这正是本实验第一步要验证的现象。创建两个工作流定义先向 Conductor 的/metadata/workflow端点发送POST请求注册两个工作流test_workflow_for_eventHandler包含一个EVENT任务用于启动另一个工作流和一个WAIT系统任务等待被事件完成。test_workflow_startedBy_eventHandler包含一个EVENT任务用于发布事件以完成上面那个WAIT任务。工作流一发出启动事件并挂起等待{ name: test_workflow_for_eventHandler, description: A test workflow to start another workflow with EventHandler, version: 1, tasks: [ { name: test_start_workflow_event, taskReferenceName: start_workflow_with_event, type: EVENT, sink: conductor }, { name: test_task_tobe_completed_by_eventHandler, taskReferenceName: test_task_tobe_completed_by_eventHandler, type: WAIT } ] }工作流二由事件启动并回发完成事件{ name: test_workflow_startedBy_eventHandler, description: A test workflow which is started by EventHandler, and then goes on to complete task in another workflow., version: 1, tasks: [ { name: test_complete_task_event, taskReferenceName: complete_task_with_event, inputParameters: { sourceWorkflowId: ${workflow.input.sourceWorkflowId} }, type: EVENT, sink: conductor } ] }注意第二份定义里EVENT任务通过inputParameters引用工作流输入中的sourceWorkflowId——这个值正是接下来由start_workflow动作注入的它是回发事件能够定位到第一个工作流实例的关键数据。Event 任务的工作方式EVENT是系统任务System Task与其他任务一样在工作流中直接声明但有一个重要区别EVENT任务以及WAIT任务无需预先注册 TaskDef 就能直接使用因此本实验不需要注册任何自定义任务。EVENT任务的核心参数是sink它决定了事件发布到哪里参数必填说明sink是格式为provider:provider-specific destination表达式在运行时解析inputParameters否用户自定义的事件负载字段asyncComplete否默认false设为true时发布成功后任务仍保持IN_PROGRESS需由外部更新或 Event Handler 的complete_task/fail_task动作来结束在本仓库的 OSS 实现中已注册的 provider 标识包括conductor、kafka、sqs、nats、jsm、nats_stream、amqp_queue、amqp_exchange需对应服务端模块启用。provider 拥有第一个冒号之后的目的地语法可能是 Kafka topic、SQS 队列 URL、NATS subject 或 AMQP 队列/交换机等。当sink只写conductor时运行时会被展开为conductor:workflowName:taskReferenceName即conductor:test_workflow_for_eventHandler:start_workflow_with_eventEvent Handler 必须订阅这个展开后的完整名称。事件发布时Conductor 会解析inputParameters并在负载中加入工作流元数据包括字段值workflowInstanceId父工作流执行 IDworkflowType父工作流名称workflowVersion父工作流版本correlationId父工作流关联 IDtaskToDomain父工作流任务域映射任务输出中还包含event_produced展开后的 sink但该字段不会随事件消息发送到 broker。EVENT任务以自身任务 ID 作为 broker 消息身份标识下游消费者可以用这个稳定值做去重检测。事件已发出但还没有被处理当你尝试启动test_workflow_for_eventHandler时会发现事件成功发出了但第二个工作流test_workflow_startedBy_eventHandler并没有被启动。原因是事件虽然已发送但我们还没有定义 Event Handler 来让 Conductor 针对该事件执行任何动作。这正是事件驱动编排的关键认知发布与消费是解耦的两件事缺了 Event Handler 这层订阅 动作的中间件事件就只是被丢弃在队列里的消息。创建 Event HandlerEvent Handler 的定义与 Task/Workflow 定义类似也是纯 JSON 元数据通过 REST API 注册。它从name开始逐步补齐事件源、条件和动作列表。第一步命名{ name: test_start_workflow }第二步指定要监听的事件队列Event Handler 需要知道自己要监听哪个队列这在event参数中定义。使用 Conductor 内置队列时event格式为conductor:{workflow_name}:{taskReferenceName}使用 SQS 时格式为sqs:{my_sqs_queue_name}{ name: test_start_workflow, event: conductor:test_workflow_for_eventHandler:start_workflow_with_event }从源码实现看EventHandler.javaevent字段带NotEmpty校验格式为provider:provider-specific queue URI运行时解析会在第一个冒号处拆分。除conductor和sqs外事件标识同样支持kafka、nats、jsm、nats_stream、amqp_queue、amqp_exchange等已注册 provider前提是对应模块启用。第三步定义动作列表Event Handler 通过actions数组声明对该事件队列要执行的一系列动作并用active标记是否启用{ name: test_start_workflow, event: conductor:test_workflow_for_eventHandler:start_workflow_with_event, actions: [ insert-actions-here ], active: true }关于active在 EventHandler.java 的注释中明确说明如果设为 false事件处理器将被停用只有active为 true 的处理器才会被订阅处理DefaultEventProcessor.java 通过getEventHandlersForEvent(event, true)只取启用状态的处理器。第四步定义 start_workflow 动作start_workflow动作用于启动一个工作流。其参数可以使用 Start Workflow Request 中的任意字段如name、version、correlationId、input、taskToDomain。这里我们在input中传入sourceWorkflowId值来自${workflowInstanceId}——即事件消息中的父工作流实例 ID这样第二个工作流就知道要回发的目标{ action: start_workflow, start_workflow: { name: test_workflow_startedBy_eventHandler, input: { sourceWorkflowId: ${workflowInstanceId} } } }动作表达式如${workflowInstanceId}从事件负载的根部解析。这里解析的负载来自 Event 任务发布的消息消息中已包含workflowInstanceId、workflowType、workflowVersion等元数据字段因此${workflowInstanceId}能直接取到第一个工作流的执行 ID。注册到 /event 端点向/api/event端点发送POST请求Controller 实现在 EventResource.java映射常量为/api/event见 RequestMappingConstants.java{ name: test_start_workflow, event: conductor:test_workflow_for_eventHandler:start_workflow_with_event, actions: [ { action: start_workflow, start_workflow: { name: test_workflow_startedBy_eventHandler, input: { sourceWorkflowId: ${workflowInstanceId} } } } ], active: true }第五步创建 complete_task 动作的 Event Handler类似地再创建一个 Event Handler 用于完成任务。complete_task动作需要定位到目标任务可以指定taskId也可以同时指定workflowId与taskRefName。这里用${sourceWorkflowId}即启动第二个工作流时注入的输入值引用第一个工作流并通过taskRefName指向挂起的WAIT任务{ name: test_complete_task_event, event: conductor:test_workflow_startedBy_eventHandler:complete_task_with_event, actions: [ { action: complete_task, complete_task: { workflowId: ${sourceWorkflowId}, taskRefName: test_task_tobe_completed_by_eventHandler } } ], active: true }从源码实现看SimpleActionProcessor.javacomplete_task动作会先把workflowId、taskId、taskRefName、reasonForIncompletion和output等字段放入输入映射再用事件负载做表达式替换随后通过workflowExecutor.updateTask将任务置为COMPLETED状态。若指定了workflowIdtaskRefName处理器会先取出工作流再按引用名查找任务对循环任务会取迭代值最高的实例。若按workflowIdtaskRefName都找不到任务输出中会带上error字段提示。Event Handler 的完整字段参考综合以上实验与 Event Handler 配置文档Event Handler 的完整字段如下字段必填行为name是非空且在一个 Conductor 实例内必须唯一源码中带NotEmpty校验event是provider:provider-specific queue URI运行时按第一个冒号拆分condition否针对负载根节点求值的条件表达式缺省视为 true如$.status READYactions是非空的动作列表动作之间并发执行、非原子active否默认false只有为true时才被订阅处理evaluatorType否指定注册的求值器否则使用默认脚本求值器Action.Type枚举EventHandler.java声明了 6 种动作但 OSS开源版的SimpleActionProcessor只实现了其中 3 种动作OSS Conductor行为start_workflow支持启动指定工作流并向其输入注入事件元数据complete_task支持完成任务指定taskId或workflowIdtaskRefNamefail_task支持使任务失败可设置reasonForIncompletionterminate_workflow不支持共享模型中存在但 OSS 动作处理器未实现update_workflow_variables不支持共享模型中存在但 OSS 动作处理器未实现在动作处理器中complete_task与fail_task走同一套completeTask逻辑只是最终状态分别为COMPLETED与FAILEDSimpleActionProcessor.java。start_workflow动作执行时会先对input和correlationId、taskToDomain等参数做表达式替换向工作流输入注入conductor.event.messageId与conductor.event.name元数据再调用workflowExecutor.startWorkflow启动SimpleActionProcessor.java。此外动作还支持expandInlineJSON布尔字段设为true时负载中的内联 JSON 字符串会被展开为完整 JSON 文档后再进行表达式解析。事件处理链路与可靠性机制理解底层处理链路有助于排查实验中出现的问题。核心处理器是 DefaultEventProcessor.java可通过conductor.default-event-processor.enabledfalse关闭其处理流程为队列监听器收到消息后按queue.getType() : queue.getName()还原事件标识调用metadataService.getEventHandlersForEvent(event, true)取出所有启用状态的 Event Handler对每个 Handler 先求值condition条件不满足则记录SKIPPED状态的事件执行并跳过动作条件通过后为actions中的每个动作创建独立的EventExecutionID 为消息ID_动作索引并提交到固定线程池并发执行执行成功则记录COMPLETED遇到瞬时错误则移除执行记录并重投递消息以便重试DefaultEventProcessor.java。这里有几个值得注意的可靠性语义至少一次投递at-least-oncebroker 投递不保证恰好一次消费端的动作及其副作用应设计为幂等去重依赖稳定消息 ID每个动作以broker 消息 ID 动作索引为标识持久化记录稳定的 broker 消息 IDEvent 任务即任务 ID可支撑持久化去重但下游的工作流启动、任务更新等外部副作用仍需幂等动作并发且非原子同一 Handler 的多个动作是并发执行的不构成原子事务任务定位是精确机制complete_task/fail_task通过taskId或workflowIdtaskRefName精确定位任务OSS 不会把业务关联键解析成等待中的任务。验证实验完整事件流将以上所有定义注册完成后启动test_workflow_for_eventHandler工作流整个事件闭环应当依次发生第一个工作流中的EVENT任务向conductor:test_workflow_for_eventHandler:start_workflow_with_event发布事件Event Handlertest_start_workflow监听到该事件执行start_workflow动作启动test_workflow_startedBy_eventHandler工作流并把sourceWorkflowId第一个工作流的实例 ID注入其输入同时第一个工作流继续执行到WAIT任务进入IN_PROGRESS状态挂起等待第二个工作流中的EVENT任务发布事件到conductor:test_workflow_startedBy_eventHandler:complete_task_with_eventEvent Handlertest_complete_task_event监听到该事件执行complete_task动作通过workflowIdtaskRefName找到第一个工作流中挂起的WAIT任务并置为COMPLETED两个工作流最终都进入COMPLETED状态。本仓库的docs/devguide/cookbook/examples/events/目录还提供了与本实验同构的完整示例可供对照学习start-workflow-handler.json含condition条件过滤与correlationId的start_workflow处理器、complete-wait-handler.json通过 Kafka 队列完成WAIT任务并回写输出的处理器以及配套的 publish-internal-event-workflow.json、fulfill-order-workflow.json 等工作流定义。进阶阅读Event Handler 完整配置参考事件标识格式、条件与负载表达式、动作能力矩阵、并发与去重语义Event Handlers API/api/event的增删改查端点与各字段说明EVENT 系统任务参考sink展开规则、发布负载与asyncComplete行为发布事件指南EVENT与KAFKA_PUBLISH的选型与生产实践事件总线编排provider 矩阵、路由、webhook、信号与投递可观测性。【免费下载链接】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),仅供参考
网站建设高端定制企业官网