Conductor PULL_WORKFLOW_MESSAGES 系统任务详解:基于工作流消息队列(WMQ)的推送式消息消费实战
Conductor PULL_WORKFLOW_MESSAGES 系统任务详解基于工作流消息队列WMQ的推送式消息消费实战【免费下载链接】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/conductorConductor 的PULL_WORKFLOW_MESSAGES是一个异步系统任务system task用于从工作流消息队列Workflow Message Queue简称 WMQ中拉取外部系统推送的消息批次让工作流能够像事件循环一样空闲等待、处理消息、再回到等待状态。本指南以 pull-workflow-messages-task.md 为骨架结合 WMQ 架构文档与源码实现全面讲解该任务的可用前提、配置方式、执行语义、数据模型与故障行为。读完本文你将掌握在 Conductor 中如何用PULL_WORKFLOW_MESSAGES构建事件驱动的 agent 循环、webhook 驱动工作流与人机协同流程。一、任务定位PULL_WORKFLOW_MESSAGES是什么type: PULL_WORKFLOW_MESSAGESPULL_WORKFLOW_MESSAGES等待工作流消息队列中可用的消息并把收到的消息批次暴露给工作流。它的设计目标是服务那些使用workflow-message-queueWMQ特性、而不是依赖 worker 轮询器worker poller的工作流。在 Conductor 的系统任务体系中它与其他系统任务并列是 系统任务索引 中明确列出的内置能力Pull a batch from workflow-message queues; requiresconductor.workflow-message-queue.enabledtrue。从架构角度看WMQ 引入了每个工作流一个的、由 Redis 支撑的消息缓冲区PULL_WORKFLOW_MESSAGES从缓冲区中取出消息。这与 Conductor 既有机制有本质区别详见 workflow-message-queue-architecture.md 中的对比既有机制为何不满足 WMQ 需求事件处理器Event handlers作用于工作流定义层面而非每个工作流实例无法按实例 ID 定向到某个正在运行的工作流WAIT 任务没有结构化消息载荷外部完成 API 解锁时无法把数据带进任务输出HTTP 任务需要工作流主动去请求外部端点WMQ 反转了该模型工作流接收外部推送的数据任务如何注册特性开关AvailabilityPULL_WORKFLOW_MESSAGES只有在conductor.workflow-message-queue.enabledtrue时才会注册。它同时依赖对应的 WMQ 基础设施与配置。若特性被禁用使用该类型的工作流将无法被映射cannot be mapped。源码层面有三处ConditionalOnProperty把关PullWorkflowMessages.java系统任务 Bean 只有在该属性为true时才创建PullWorkflowMessagesTaskMapper.java对应的任务映射器TaskMapper同样受开关控制特性关闭时任务类型对SystemTaskRegistry未知任何引用它的工作流定义都会在校验阶段失败WorkflowMessageQueueResource.java推送端点的 REST Controller 也不注册端点直接 404。特性关闭时所有相关 Bean系统任务、映射器、DAO、清理监听器都不创建运行时零占用已有部署不受影响。二、配置开启 WMQ 与设置任务输入2.1 服务端开关在 Conductor 服务器配置中如 config.properties 对应格式conductor.workflow-message-queue.enabledtrue conductor.workflow-message-queue.maxQueueSize1000 conductor.workflow-message-queue.ttlSeconds86400 conductor.workflow-message-queue.maxBatchSize100这些属性由 WorkflowMessageQueueProperties.java 这一ConfigurationProperties类承载前缀为conductor.workflow-message-queue属性类型默认值说明enabledbooleanfalse总开关。所有 WMQ Bean 都以此为门控maxQueueSizeint1000单个工作流队列同时允许的最大消息数超出时 push 报错背压ttlSecondslong86400Redis key 的 TTL默认 24 小时每次 push 都会重置maxBatchSizeint100单次PULL_WORKFLOW_MESSAGES执行的服务端batchSize上限防止失控的批量大小2.2 任务定义中的 inputParameters任务映射器会解析任务的inputParameters队列 worker 消费这些参数。当工作流需要限制单次拉取数量时提供batchSize此外还需要提供所配置消息队列实现要求的任何队列特定输入。{ name: pull_messages, taskReferenceName: pull_messages, type: PULL_WORKFLOW_MESSAGES, inputParameters: { batchSize: 10 } }任务在消息可用之前保持in progress状态。特性配置与投递语义参见 Workflow Message Queue。输入参数语义源码确认在 PullWorkflowMessages.java 中定义了完整的输入/输出契约参数类型默认说明batchSizeint1单次调用最多出队消息数服务端按maxBatchSize封顶Math.min(requested, properties.getMaxBatchSize())且小于 1 时按 1 处理blockingbooleantrue为false时无消息则立即以空列表完成而不是等待映射器 PullWorkflowMessagesTaskMapper.java 通过ParametersUtils.getTaskInputV2完成输入解析支持${workflow.input.xxx}等引用创建任务后直接置为IN_PROGRESS等待系统任务 worker 轮询。三、任务执行语义与消息数据模型3.1 生命周期从 SCHEDULED 到 COMPLETEDPULL_WORKFLOW_MESSAGES是异步系统任务isAsync() true接入SystemTaskWorker轮询循环。其生命周期start()——任务进入SCHEDULED时调用一次将状态置为IN_PROGRESS源码execute()——每个轮询周期由SystemTaskWorker调用检查 Redis 队列队列为空返回false任务保持IN_PROGRESSAsyncSystemTaskExecutor以较短的callbackAfterSeconds重新入队任务消息队列非空原子弹出最多batchSize条消息写入输出状态置COMPLETED返回true源码。完成时的输出字段字段类型说明messages消息对象数组每条含id、workflowId、payload、receivedAtcountint实际返回消息数恒 batchSize轮询节奏PULL_WORKFLOW_MESSAGES重写了getEvaluationOffset()返回Optional.of(1L)源码即等待消息期间每 1 秒被重新评估一次而非默认的systemTaskWorkerCallbackDuration30 秒从而显著降低唤醒延迟。超时行为任务遵循工作流任务定义或任务定义元数据中配置的timeoutSeconds。若等待消息超时Conductor 走标准超时机制将任务转为TIMED_OUTPULL_WORKFLOW_MESSAGES自身无需特殊处理。3.2 消息 Schema每条消息都是包含以下字段的 JSON 对象{ id: 3f2504e0-4f89-11d3-9a0c-0305e82c3301, workflowId: 8e2c14e1-99ab-4c10-b4a8-a7b0d2f0e123, payload: { decision: approved, approvedBy: userexample.com }, receivedAt: 2025-06-15T10:30:00Z }字段类型说明idUUID v4 字符串推送端点摄入时生成返回给调用方workflowId字符串拥有该消息的工作流实例 ID与队列 key 冗余但保留以便下游追踪payload任意 JSON 对象外部调用方提供的数据Conductor 不解释也不校验结构receivedAtISO-8601 UTC 时间戳摄入时记录POJO 定义见 WorkflowMessage.java。3.3 底层存储Redis List 与原子出队WMQ 存储层在core定义接口、在redis-persistence实现接口WorkflowMessageQueueDAO.java操作包括push(workflowId, message)、pop(workflowId, maxCount)原子出队空队列返回空列表而非 null、size(workflowId)、delete(workflowId)Redis 实现RedisWorkflowMessageQueueDAO.java每个工作流一个 Redis List。属性细节Key 模式wmq:{workflowId}入队RPUSH——追加到尾部保证 FIFO 顺序出队LRANGE读取 LTRIM移除因 Conductor 在 decide 周期持有每工作流执行锁无需原子 Lua 脚本也安全Redis 6.2 可用LPOP key count简化TTL可配置默认 24 小时86400 秒每次RPUSH重置为完整 TTL最大容量可配置上限默认 1000 条队列满载时push返回错误四、端到端实战推送消息、拉取与事件循环4.1 完整链路Push → Pull第 1 步服务器开启特性见 2.1。第 2 步注册包含PULL_WORKFLOW_MESSAGES的工作流定义。以 wmq_echo_loop.json 为参考这是仓库自带的可运行示例用DO_WHILE包裹pull_messagebatchSize: 1与echo_payloadINLINE 任务将$.messages[0]回显。第 3 步向运行中的工作流推送消息curl -X POST http://localhost:8080/api/workflow/{workflowId}/messages \ -H Content-Type: application/json \ -d {text: hello}第 4 步任务完成后工作流通过output.messages取到批次{ messages: [ { id: 3f2504e0-4f89-11d3-9a0c-0305e82c3301, workflowId: 8e2c14e1-..., payload: { text: hello }, receivedAt: 2025-06-15T10:30:00Z } ], count: 1 }用户数据通过output.messages[0].payload访问id与receivedAt是 Conductor 摄入时附加的字段。推送错误语义WorkflowMessageQueueResource.java 的实现细节404 Not Found——工作流 ID 不存在或 WMQ 特性被禁用409 Conflict——工作流不在RUNNING状态completed、failed、terminated 等。消息不存储。注意推送后还会二次校验状态以捕获 TOCTOU 竞态推送瞬间工作流转入终态会删除队列并返回 409429 Too Many Requests——队列已满达到maxQueueSize。调用方必须退避重试。推送成功的副作用存储消息后端点会调用workflowExecutor.decide(workflowId)立即触发一次工作流求值周期让正在等待的PULL_WORKFLOW_MESSAGES任务无需等待下一个 SystemTaskWorker 轮询间隔即可被唤醒接近实时。4.2 事件循环模式Event loop pattern对于处理无界消息流的工作流用DO_WHILE包裹任务{ name: message_loop, taskReferenceName: message_loop_ref, type: DO_WHILE, loopCondition: $.message_loop_ref[iteration] 100, loopOver: [ { name: pull_message, taskReferenceName: pull_message_ref, type: PULL_WORKFLOW_MESSAGES, inputParameters: { batchSize: 1 } }, { name: process_message, taskReferenceName: process_message_ref, type: INLINE, inputParameters: { evaluatorType: javascript, expression: function e() { return { payload: $.messages[0].payload }; } e();, messages: ${pull_message_ref.output.messages} } } ] }循环在PULL_WORKFLOW_MESSAGES上停靠直到下一条消息到达。4.3 与 WorkflowSweeper 的交互作为异步系统任务PULL_WORKFLOW_MESSAGES的执行链路为引擎调度任务时放入系统任务队列QueueDAO支撑、按任务类型分 key 的队列SystemTaskWorkerCoordinator启动时将其注册给SystemTaskWorker开始轮询队列每次轮询调用AsyncSystemTaskExecutor.execute()首次执行调start()后续调execute()execute()返回false队列空时AsyncSystemTaskExecutor以systemTaskCallbackTime秒延迟由conductor.app.systemTaskWorkerCallbackDuration配置重新入队任务 IDREST 推送消息后立即调用workflowExecutor.decide(workflowId)sweeper 重新评估工作流并触发任何IN_PROGRESS异步任务的执行把唤醒延迟从完整轮询间隔降到接近实时。五、典型应用场景与 Agent 集成架构文档 workflow-message-queue-architecture.md 列出的主要用例用例描述Agent 循环agentic/agent loopsAI agent 工作流循环等待外部调用方的工具结果或人工确认调用方在工具响应时推送消息循环解锁Webhook 驱动的工作流异步 HTTP 回调需要把数据喂给暂停的工作流回调目标是 WMQ 推送端点而非轮询机制通知管道工作流循环、按可配置批次读消息并 fork 扇出到多个渠道人机协同human-in-the-loop操作员工具或 UI 向运行中的工作流注入结构化的人工决策或审批载荷WMQ 是框架无关的。在 Conductor 图中使用PULL_WORKFLOW_MESSAGES停靠执行直到消息到达再把返回的 payload 传给下一个任务。SDK 编写的 agent 参见 Conductor Agents。Kafka 桥接示例Kafka 消费者可以把每条 record 翻译成一次POST /api/workflow/{workflowId}/messages请求payload 形状如上从而把外部事件流桥接到工作流内部该消费者实现应保留在其所属 SDK 或服务仓库中与工作流 agent 步骤使用的框架相互独立。六、故障模式与韧性场景行为推送时 Redis 不可用dao.push()抛异常REST 端点返回 HTTP 500消息未存储无丢失调用方需重试拉取时 Redis 不可用dao.pop()抛异常execute()传播错误AsyncSystemTaskExecutor按标准重试/超时配置处理任务级失败PULL_WORKFLOW_MESSAGESIN_PROGRESS 时工作流终止WorkflowSweeper 检测到终态并取消待处理任务WorkflowMessageQueueCleanupListener删除队列 keyIN_PROGRESS 时工作流暂停任务保持IN_PROGRESS暂停期间推送的消息在 Redis 排队受maxQueueSize约束恢复后decide()触发重新评估任务在下次轮询取走积压消息队列容量超限push返回错误REST 端点返回 HTTP 429 或 400调用方需处理背压超大消息载荷payload 内联存储在 Redis list 条目中对非常大的数据采用外部存储引用模式数据放对象存储S3、GCS 等WMQ payload 只放引用 URL 或 key与 Conductor 外部 payload 存储同一模式重复投递Lua 出队脚本在单 Redis 实例内原子同一条消息不会在一次execute()内投递两次但系统层面为至少一次语义调用方与工作流设计者应尽量将消费视为幂等Conductor 多节点网络分区Redis List 共享原子出队保证同一工作流的两次并发execute()不会返回重叠消息工作流级锁ExecutionLockService在 decide 路径上提供额外防护七、生命周期清理与安全注意事项清理队列清理依赖 Redis TTLconductor.workflow-message-queue.ttlSeconds默认 24 小时每次 push 重置 TTL活跃队列不会被提前过期。没有显式的WorkflowStatusListener清理实现——因为WorkflowStatusListener是单 Bean 接口叠加 WMQ 实现会与归档监听器等其他实现冲突。Redis TTL 足够兜底任何孤儿队列如服务器在工作流完成前崩溃都会自动过期。安全与访问控制推送端点接收任意 JSONConductor 不校验 payload 结构应在 API 网关层加认证/授权推送 URL 中的workflowId足以定向任意运行中的工作流调用方必须可信或端点必须受保护payload 数据按配置的 TTL 存于 Redis敏感数据受 Redis 访问控制约束必要时应在应用层对敏感字段加密。八、组件地图组件位置用途WorkflowMessageQueueDAOcore定义存储契约接口WorkflowMessagecommon单条消息 POJOWorkflowMessageQueuePropertiescore所有 WMQ 设置的ConfigurationPropertiesPullWorkflowMessagescore将消息出队到工作流输出的系统任务PullWorkflowMessagesTaskMappercore任务映射器解析输入并创建IN_PROGRESS任务RedisWorkflowMessageQueueDAOredis-persistenceRedis List 版 DAO 实现WorkflowMessageQueueResourcerest推送端点 REST Controller继续阅读特性级配置与投递语义见 Workflow Message Queue完整架构、数据流图与 WorkflowSweeper 交互见 Workflow Message Queue — Architecture可运行示例 wmq_echo_loop.json系统任务全览见 System Tasks推送端点的 API 文档见 workflow.md。【免费下载链接】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),仅供参考