Conductor 事件处理器(Event Handler)参考:事件标识、条件表达式、动作能力矩阵与去重语义

发布时间:2026/9/10 1:11:04

Conductor 事件处理器(Event Handler)参考:事件标识、条件表达式、动作能力矩阵与去重语义 Conductor 事件处理器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/documentation/configuration/eventhandlers.md完整梳理事件处理器的数据模型、事件标识格式、条件与载荷表达式的作用域、动作能力矩阵含 OSS 与 Orkes 的差异以及并发执行与去重语义并结合 EventHandler 数据模型 与 DefaultEventProcessor 等源码说明每条规则在运行时是如何被解析和执行的。读完后你可以独立完成注册一个事件处理器 → 正确书写条件与占位符 → 选择合适的动作 → 评估幂等与去重边界的完整链路。事件处理器是什么一条消费、评估、分派的规则事件处理器event handler是注册在 Conductor 服务端的一条规则它消费某个 provider 事件对可选条件求值然后分派一个或多个动作actions。只有通过 Event Handlers API 注册且处于active状态的处理器才会被订阅参与处理——这是理解整篇参考的第一前提。一个最小可运行的示例来自 start-workflow-handler.json{ name: start_fulfillment_on_order_ready, event: conductor:publish_order_event:order-status, condition: $.status READY, actions: [ { action: start_workflow, start_workflow: { name: fulfill_order, version: 1, correlationId: ${orderId}, input: { orderId: ${orderId}, sourceEventId: ${workflowInstanceId} } } } ], active: true }对应源码中EventHandler顶层字段的定义EventHandler.java#L32-L55字段是否必填行为name是非空、唯一的处理器名称源码中以NotEmpty约束event是provider:provider 特定队列 URI运行时按第一个冒号切分condition否针对载荷根求值省略视为恒真actions是非空列表源码NotEmpty约束动作并发执行active否默认falseevaluatorType否指定已注册的求值器未指定时使用默认脚本求值器注意active的默认值是false注册了一个忘了置为true的处理器不会有任何流量被消费这是最容易被忽视的配置了却不生效场景。事件标识provider:queue URI与第一个冒号切分事件标识event字段的格式为provider:provider-specific queue URI运行时解析时在第一个冒号处切分——前半段必须是一个已注册且服务端已启用模块的 provider 键后半段整体作为该 provider 内部的队列标识因此队列名里可以再包含冒号比如示例中的publish_order_event:order-status。官方文档列出的有效 provider 键为conductor、kafka、sqs、nats、jsm、nats_stream、amqp_queue、amqp_exchange前提是服务端启用了相应模块。这一解析方式在源码中可以直接印证DefaultEventProcessor.java#L120 中事件名由queue.getType() : queue.getName()拼成随后与处理器注册的event字符串做匹配。也就是说第一个冒号切分不是文档约定而是队列类型与队列名的运行时拼接结果。在 Orkes 托管版上事件处理器的前置条件是完成托管 broker 集成的配置然后在该集成基础上使用OSS 示例中的 provider 键如conductor只是说明 API 用法不代表完成了集成配置。条件与载荷表达式作用域是载荷根这是参考文档中最容易出错的语义部分官方逐条明确了以下规则active默认为false不激活不消费见上文。缺省条件视为真不写condition的处理器对队列上的每条消息都执行动作。条件直接对载荷根求值例如$.status READY。这里的$指向投递进来的 payload 本身而不是任何包装层。evaluatorType决定求值器若其值对应一个已注册的 evaluator则使用该求值器否则回落到默认脚本求值器。动作占位符同样从载荷根解析例如${orderId}。expandInlineJSON: true会在表达式解析前把字符串化的 JSON 字段展开为真正的 JSON 结构再求值。源码对前三条的实现见 DefaultEventProcessor.java#L163-L179Object payloadObject getPayloadObject(msg.getPayload()); for (EventHandler eventHandler : eventHandlerList) { String condition eventHandler.getCondition(); String evaluatorType eventHandler.getEvaluatorType(); // 缺省条件 → success true直接落入动作执行 boolean success true; if (StringUtils.isNotEmpty(condition) evaluators.get(evaluatorType) ! null) { Object result evaluators.get(evaluatorType) .evaluate(condition, jsonUtils.expand(payloadObject)); success ScriptEvaluator.toBoolean(result); } else if (StringUtils.isNotEmpty(condition)) { success ScriptEvaluator.evalBool(condition, jsonUtils.expand(payloadObject)); } ... }两个值得注意的实现细节jsonUtils.expand(payloadObject)在进入条件求值前对载荷做了一次展开这与expandInlineJSON的语义相呼应——载荷中的内联 JSON 字符串会先被展开表达式再对展开后的结构求值。条件求值失败success false时处理器会写入一条状态为SKIPPED的事件执行记录DefaultEventProcessor.java#L181-L197并continue跳过该处理器的所有动作——即条件为假不执行任何动作但留有可审计的记录。动作能力矩阵OSS 与 Orkes 的边界官方能力矩阵原文如下这是 OSS 版本选型时的硬边界动作OSS ConductorOrkes行为start_workflowYesYes启动指定工作流并追加 Conductor 事件元数据到其输入complete_taskYesYes完成一个被精确定位的任务fail_taskYesYes使一个被精确定位的任务失败可设置reasonForIncompletionterminate_workflowNoYes终止目标工作流update_workflow_variablesNoYes更新目标工作流的变量也就是说complete_task与fail_task要求精确的任务定位——提供taskId或者同时提供workflowId与taskRefName。这是精确的任务寻址机制OSS 处理器不会把一个业务关联键correlation key自动解析成某个等待中的任务。terminate_workflow和update_workflow_variables存在于共享模型中但 OSS 动作处理器未实现请求可以反序列化通过却会在处理阶段以不支持失败。从源码结构看模型侧与文档描述一致EventHandler.java#L142-L153 中Action.Type枚举声明了start_workflow、complete_task、fail_task、terminate_workflow、update_workflow_variables此外源码还声明了一个start_agent取值但官方能力矩阵未将其列入——OSS 动作处理器实际支持的动作范围仍以文档矩阵表为准。各动作的参数结构同样定义在 EventHandler.java 的内嵌类中StartWorkflowL396-L415name、version、correlationId、input、taskToDomainTaskDetailsL295-L314workflowId、taskRefName、taskId、output、reasonForIncompletionTerminateWorkflowworkflowId、terminationReasonUpdateWorkflowVariablesworkflowId、variables、appendArray。实战示例用事件完成/失败一个被精确定位的任务当事件本身携带任务身份时可以用任务动作承接。官方 Consume and route events 给出的一对典型处理器审批事件到达时完成任务——{ name: complete_payment_wait, event: kafka:payment-events, condition: $.status APPROVED, actions: [ { action: complete_task, complete_task: { workflowId: ${workflowId}, taskRefName: wait_for_payment, output: { paymentId: ${paymentId}, approved: true } } } ], active: true }拒绝事件到达时失败任务建议注册为独立的处理器通过条件区分语义{ name: fail_payment_wait, event: kafka:payment-events, condition: $.status REJECTED, actions: [ { action: fail_task, fail_task: { taskId: ${rejectionTaskId}, reasonForIncompletion: ${reason}, output: { providerStatus: ${status} } } } ], active: true }注意两个示例分别演示了两种任务定位方式前者用workflowIdtaskRefName后者用taskId二者均从 broker 载荷根解析占位符。reasonForIncompletion仅在fail_task中有语义output中的字段同样是按载荷根做表达式解析的。并发执行与去重broker 消息 ID 动作索引参考文档对投递语义的表述非常克制这也是设计幂等方案时必须完整继承的部分动作并发执行且不是原子的每个动作单独记录记录 ID 由 broker 消息 ID 加上该动作在列表中的索引构成稳定的 broker 消息 ID可以在事件执行记录落库之后启用持久化去重但下游的工作流启动、任务更新以及外部副作用仍然要求幂等条件为假时记录一条 SKIPPED 事件执行不执行任何动作。源码层面这条链路可以逐一对应DefaultEventProcessor.java#L236-L262protected CompletableFutureListEventExecution executeActionsForEventHandler( EventHandler eventHandler, Message msg) { ListCompletableFutureEventExecution futuresList new ArrayList(); int i 0; for (Action action : eventHandler.getActions()) { String id msg.getId() _ i; // 消息 ID 动作索引 EventExecution eventExecution new EventExecution(id, msg.getId()); ... if (executionService.addEventExecution(eventExecution)) { futuresList.add(CompletableFuture.supplyAsync( () - execute(eventExecution, action, getPayloadObject(msg.getPayload())), eventActionExecutorService)); } else { LOGGER.warn(Duplicate delivery/execution of message: {}, msg.getId()); } } return CompletableFutures.allAsList(futuresList); }由此可以得到几个可直接落地的结论去重的粒度是消息 × 动作索引重复投递的同一条 broker 消息在同一个动作索引上会被addEventExecution的返回值false识别为重复并跳过这正是稳定消息 ID 才能去重的原因。并发的代价动作提交到独立的线程池event-action-executorService执行多个动作之间没有事务边界部分成功部分失败是可能状态。瞬时失败不落终态从源码结构看L272-L317动作执行经由RetryTemplate包裹若捕获到瞬时异常EventExecution保持IN_PROGRESS状态并走processTransientFailures移除已暂存的记录使消息可以被重试非瞬时异常则把状态置为FAILED并记录异常信息。消息 ack 策略仅当执行未失败且没有瞬时失败动作时才ack对无 unack 超时的队列瞬时失败会通过rePublish重新投回队列否则走nackDefaultEventProcessor.java#L127-L140。因此官方建议保持不变工作流启动、任务更新和任何外部副作用都应自行保证幂等把持久化去重理解为重复消息不会重复执行同一动作而不是端到端的 exactly-once。服务端开关与线程池参数事件处理管线整体由 DefaultEventProcessor 承载两个可直接使用的服务端配置配置项默认值作用conductor.default-event-processor.enabledtrue缺省即启用置为false完全禁用事件处理conductor.app.event-processor-thread-count2动作执行线程池大小必须 0第二项来自 ConductorProperties.java#L68 中eventProcessorThreadCount字段ConfigurationProperties(conductor.app)前缀构造器中若该值 ≤ 0 会直接抛出IllegalStateException提示要禁用事件处理请设置conductor.default-event-processor.enabledfalse而不是把线程数设成 0。此外executeEvent只拉取active的处理器metadataService.getEventHandlersForEvent(event, true)L157从代码上再次印证了active状态决定订阅的文档语义。参考与延伸完整的字段、端点与状态语义见 Event Handlers APIREST 挂载在/api/event成功变更返回空200 OK首次使用的手把手演练注册、条件匹配、任务动作、投递与幂等见 Consume and route events本文作为其动作与表达式的参考手册事件执行记录的数据模型见 EventExecution.java可与 EventHandler.java 对照阅读处理器动作的执行入口接口为 ActionProcessor实现类为 SimpleActionProcessor测试用例 TestSimpleActionProcessor 与 TestDefaultEventProcessor 可作为行为回归验证入口。【免费下载链接】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),仅供参考
延伸阅读

更多相关文章

2026/9/10 1:11:03

Locust高并发压测实战:从脚本设计到瓶颈定位

如果你也经历过那种凌晨两点被监控电话叫醒,打开面板看到CPU打满、数据库连接池耗尽、线上服务一个接一个雪崩的场景,你应该能理解“压力测试”这四个字的分量。系统跑得好不好,不是上线那一刻才确定的,而是取决于你有没有在上线之…

2026/9/10 1:06:03

计算机硬件基础知识全解析:从CPU到电源的选型与排障指南

开头 “计算机基础”这四个字,听起来像是一个应该早就解决了的问题——毕竟我们每天用电脑工作、打游戏、刷视频,似乎离“基础”二字也不远。可实际情况是,我接触过太多能熟练写代码、能把系统玩出花来的朋友,一旦问到“内存频率和…

2026/9/10 2:06:11

618复盘别再手工对表了!2026年5款大促数据工具排行榜测评

618大促落幕之后,真正拉开差距的往往不是当天的销售额,而是复盘的质量。一场大促沉淀下来的订单、流量、投放、库存、售后数据,构成了下一场大促的决策底牌。市面上可用于大促复盘的数据工具,按产品形态大致可分为四类&#xff1a…

2026/9/10 2:06:11

FPGA实现AES-128加密算法:Verilog硬件设计与工程实践详解

简介:面向FPGA与硬件安全开发者的AES-128完整Verilog实现,基于Rijndael算法,覆盖密钥扩展、字节替换、行移位、列混淆等核心模块,并附带VHDL对照代码,适合用于学习对称加密算法的硬件加速、安全模块设计与芯片验证流程…

2026/9/10 2:01:11

CANN/ge图编译缓存功能

图编译缓存 【免费下载链接】ge GE(Graph Engine)是面向昇腾的图编译器和执行器,提供了计算图优化、多流并行、内存复用和模型下沉等技术手段,加速模型执行效率,减少模型内存占用。 GE 提供对 PyTorch、TensorFlow 前端…

2026/9/9 13:11:35

超人会飞不算本事:系统稳定依赖清晰规则与边界设计

开头先不绕弯子。“#斯坦李吐槽dc 所以超人是无缘无故会飞的嘛哈哈哈哈哈哈哈锤哥真是技术人才啊!#雷神 #复联”这类调侃式短标题,第一波冲击力在于它把两个宇宙的角色塞进同一个吐槽箱里,但细想一下就能发现,它真正碰到的根本不是…

2026/9/8 7:15:15

超人VS蜘蛛侠:拆解超级IP的影响力与传播方法论

把“蜘蛛侠 vs 超人”放在 CSDN 上聊,可能很多人第一反应是走错片场了。但如果把这两个角色看成“两个持续运营了 80 多年的文化产品”,你会发现,这场比较本质上是两个不同 IP 策略的长期结果对比:超人赢在定义了整个超级英雄题材…

2026/9/9 16:31:09

基于CNN的调制信号识别:MATLAB实现时频图分类实战

简介:本资源是一套面向通信工程与信号处理方向学习者、研究者的深度学习实践方案,聚焦调制信号自动检测与识别这一典型无线通信任务,解决传统方法依赖人工特征、低信噪比下性能下降等痛点。压缩包共12个文件(10.73MB)&…

2026/9/10 0:00:55

目录对比去重实战:用哈希算法精准清理重复文件

我电脑里现在还有一块换了三次机的“数据墓地”硬盘,里面存着2016年以前所有旧笔记本的完整备份。平时不觉得有什么,直到前阵子想把它整理归档,发现同一个安装包、同一批照片、同一份论文草稿,在几个不同的备份目录里反复出现。更…

2026/9/10 0:00:55

Leaflet离线地图完整Demo合集:内网部署与坐标纠偏实战

简介:这是一份面向Web GIS开发者的LeafLet离线地图示例合集,帮助开发者快速掌握离线地图从搭建到交互的完整流程。压缩包共723个文件,大小14.06MB,以319个js脚本、175个html页面和29个css样式文件为主体,配合png/svg图…

2026/9/10 0:00:55

MATLAB读取Rinex 3.02观测文件:多系统GNSS数据解析实战

简介:基于MATLAB开发的Rinex3.02版观测文件(o文件)读取代码包,面向卫星定位导航方向的学习者与研究人员,用于解决新版观测文件的数据解析、历元提取与时间转换问题。压缩包共4个文件,包含两个m脚本、一个19…

2026/9/7 16:23:03

USB Type-C PCB布局分区设计:电源、高速信号与PD协议全攻略

做硬件这行,Type-C接口算是典型的“看着简单,做起来全坑”的东西。光引脚就24个,高低速信号、电源、控制线全部塞在一个小小的连接器里,如果PCB布局不做规划,打样回来基本就是“插上没反应”、“高速掉线”、“静电一打…

2026/9/7 22:46:00

系统编程学习原型如何补齐稳定性边界

系统编程学习原型如何补齐稳定性边界预算有限时&#xff0c;我先优化明显多余的复制&#xff0c;而不是猜测性地换容器。用借用传递只读数据通常就能减少分配&#xff1a; fn parse(line: &str) -> Result<Item, Error> { /* ... */ }用基准确认热点确实在分配&am…

2026/9/9 10:21:54

雨花区哪家财务公司代理记账比较好?

在雨花区&#xff0c;企业处理财税事务常常面临诸多挑战&#xff0c;选择一家靠谱的财务公司至关重要。湖南巨勤财务管理咨询有限公司就是本地正规实体财税服务机构&#xff0c;深耕本地工商财税行业多年&#xff0c;熟悉当地工商局、税务局最新政策与申报流程。主营公司注册、…

还想了解更多?直接咨询顾问

免费诊断 + 免费方案 + 透明报价。

全国咨询热线400-8866-253
免费获取方案
咨询二维码