Conductor 中的 Human Task:配置人工介入任务并实现 Human-in-the-Loop 审批流
【免费下载链接】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 系统任务中的HUMAN类型任务用于在工作流中插入一个"人工闸门":工作流运行到该任务时会被暂停,任务保持IN_PROGRESS状态,直到外部触发(Task Update API、事件处理器或队列更新 API)将其标记为COMPLETED或FAILED。本文基于 Human Task 配置文档 展开,完整讲解 Human Task 的 JSON 配置、两种完成方式、五种监控/回调模式以及一个可直接复用的完整审批流示例,并结合开源仓库源码说明其底层执行机制,帮助你在工作流中可靠地落地人工审批与外部信号等待场景。
核心概念:Human Task 是什么
Human Task("type": "HUMAN")的作用是暂停工作流并等待外部信号。它充当一个 gate:任务创建后立即进入IN_PROGRESS状态,此后一直挂起,直到被外部触发器标记为COMPLETED或FAILED。
典型使用场景包括:
- 工作流需要暂停并等待人工干预,例如手动审批(manual approval);
- 等待来自外部来源的事件,例如 Kafka、SQS,或 Conductor 内部的队列机制。
从源码结构看,Human Task 的极简实现印证了它"只做等待"的语义。任务实现类 Human 继承自WorkflowSystemTask,其start方法仅将任务状态置为IN_PROGRESS,cancel方法仅置为CANCELED——没有任何轮询或回调逻辑,真正的状态推进完全依赖外部更新:
@Override public void start(WorkflowModel workflow, TaskModel task, WorkflowExecutor workflowExecutor) { task.setStatus(IN_PROGRESS); }JSON 配置
Human Task 不需要任何专用参数即可配置。最小配置如下:
{ "name": "human", "taskReferenceName": "human_ref", "inputParameters": {}, "type": "HUMAN" }需要注意的是:inputParameters可以为空,但并非不能携带数据——任务映射器 HumanTaskMapper 在执行时会通过ParametersUtils.getTaskInputV2对inputParameters做变量替换并写入TaskModel,同时记录startTime并将状态直接置为IN_PROGRESS:
TaskModel humanTask = taskMapperContext.createTaskModel(); humanTask.setTaskType(TASK_TYPE_HUMAN); humanTask.setInputData(humanTaskInput); humanTask.setStartTime(System.currentTimeMillis()); humanTask.setStatus(TaskModel.Status.IN_PROGRESS);这意味着你在inputParameters里写入的${...}表达式会被正常求值,可用于把上下文信息(如申请单号、发起人)带给等待中的人工任务,供审批人参考。
完成 Human Task 的两种官方方式
方式一:Task Update API
通过 Task Update API(POST api/tasks)完成 Human Task,需要提供taskId、任务状态和期望的任务输出。使用 CLI 的等价命令:
conductor task update-execution --workflow-id {workflowId} --task-ref-name waiting_around_ref --status COMPLETED --output '{"data_key":"somedatatoWait1","data_key2":"somedatatoWAit2"}'在服务端 REST 层,该能力对应 TaskResource 中的多个端点:POST /api/tasks(提交TaskResult)、POST /api/tasks/update-v2,以及按 ref name 更新的POST /api/tasks/{workflowId}/{taskRefName}/{status}(可选workerid查询参数,请求体即任务输出)。此外还有同步变体{workflowId}/{taskRefName}/{status}/sync,会直接返回更新后的 workflow,适合脚本化审批流程。
方式二:事件处理器 / Update Queue APIs
若启用了 SQS 等事件队列集成,Human Task 还可以通过 Update Queue API 完成,对应实现位于 QueueAdminResource:
POST api/queue/update/{workflowId}/{taskRefName}/{status}POST api/queue/update/{workflowId}/task/{taskId}/{status}
POST 消息体中携带的任何参数都会作为该任务的输出重复。例如发送如下 COMPLETED 消息:
curl -X "POST" "{{ server_host }}{{ api_prefix }}/queue/update/{workflowId}/waiting_around_ref/COMPLETED" \ -H 'Content-Type: application/json' \ -d '{"data_key":"somedatatoWait1","data_key2":"somedatatoWAit2"}'则 Human Task 的输出为:
{ "data_key":"somedatatoWait1", "data_key2":"somedatatoWAit2" }从源码看,该端点将请求委托给DefaultEventQueueProcessor.updateByTaskRefName / updateByTaskId(位于 DefaultEventQueueProcessor),即通过在事件队列中发布消息来驱动任务状态变更——这正是它能与 SQS/Kafka 等外部事件源解耦对接的原因。
另外,也可以配置一个使用complete_taskaction 的事件处理器来完成 Human Task,这样外部系统只需发布事件,无需直接调用 HTTP API。
监控 Human Task:回调与通知模式
当工作流到达 Human Task 时,通常需要收到通知或回调以触发下一步动作(发邮件、通知 Slack 频道、触发外部系统等)。文档推荐了五种模式:
模式 1:轮询 Workflow Status API
最简单的方式是轮询工作流执行状态,查找处于IN_PROGRESS状态的 Human 任务:
# Get workflow execution status curl '{{ server_host }}/api/workflow/{workflowId}' \ -H 'accept: application/json'解析响应,找到taskType: "HUMAN"且status: "IN_PROGRESS"的任务。
- 优点:实现简单,无需额外配置
- 缺点:需要轮询,非实时
模式 2:基于 Conductor 内部事件的事件处理器
Conductor 在任务状态变化时可以发布内部事件,可以配置事件处理器监听:
{ "name": "human_task_notification_handler", "event": "conductor:TASK_STATUS_CHANGE", "condition": "$.taskType == 'HUMAN' && $.status == 'IN_PROGRESS'", "actions": [ { "action": "start_workflow", "start_workflow": { "name": "notification_workflow", "input": { "workflowId": "${workflowId}", "taskRefName": "${taskRefName}", "taskStatus": "${status}" } } } ] }每当 Human Task 进入IN_PROGRESS状态时,即触发一个通知工作流。该处理器可配合 事件处理器配置文档 中的完整字段说明使用。
模式 3:通过 EVENT Task 做 Webhook 集成
在 Human Task 之前加一个EVENT任务,用于发送 webhook 通知:
{ "name": "notify_human_task", "taskReferenceName": "notify_ref", "type": "EVENT", "sink": "kafka:human-task-notifications", "inputParameters": { "workflowId": "${workflow.input.workflowId}", "taskRefName": "human_ref", "eventType": "HUMAN_TASK_PENDING" } }, { "name": "human_approval", "taskReferenceName": "human_ref", "type": "HUMAN" }随后配置事件处理器或外部消费者处理这些通知。sink的前缀(如kafka:)指向对应的消息中间件,需要相应的事件队列模块已启用。
模式 4:完成时携带回调输出的 Task Update
完成 Human Task 时,把回调信息放进输出中:
curl -X POST "{{ server_host }}/api/tasks" \ -H 'Content-Type: application/json' \ -d '{ "taskId": "${taskId}", "status": "COMPLETED", "output": { "approvedBy": "user@example.com", "approvedAt": "2026-04-22T10:30:00Z", "comments": "Approved for production deployment" } }'输出会供下游任务使用,可用于审计追踪或进一步通知。
模式 5:外部系统集成
用于实时通知的外部集成途径包括:
- Slack/Teams:用事件处理器触发一个通知工作流,向其 webhook 发帖
- 邮件:通过 SMTP 或邮件服务 API 发送通知
- 短信/推送:集成 Twilio、Pushover 或类似服务
- 自定义 Webhook:当有人工任务待处理时 POST 到你的内部系统
最佳实践
- 使用关联 ID:所有通知中都包含
workflowId和taskRefName,便于追踪 - 设置超时:考虑为未审批的人工任务添加超时/升级逻辑
- 审计追踪:记录所有人工任务完成的 timestamp 与用户信息
- 幂等性:确保通知处理器幂等,以处理重复事件
完整示例:带通知的审批工作流
下面是一个端到端的审批工作流定义:Slack 发起审批请求 → 等待人工审批 → 发送审批结果。
{ "name": "approval_workflow", "version": 1, "tasks": [ { "name": "send_approval_request", "taskReferenceName": "send_request_ref", "type": "HTTP", "inputParameters": { "http_request": { "method": "POST", "url": "https://hooks.slack.com/services/xxx", "body": { "text": "Approval needed for workflow ${workflow.input.requestId}" } } } }, { "name": "wait_for_approval", "taskReferenceName": "approval_ref", "type": "HUMAN" }, { "name": "send_approval_result", "taskReferenceName": "send_result_ref", "type": "HTTP", "inputParameters": { "http_request": { "method": "POST", "url": "https://hooks.slack.com/services/xxx", "body": { "text": "Approval ${approval_ref.output.status} by ${approval_ref.output.approvedBy}" } } } } ] }该工作流的执行过程:
- 发送 Slack 通知,告知需要审批;
- 工作流停在
wait_for_approval(HUMAN任务)处,可通过前文任一方式(Task Update API、Update Queue API 或complete_task事件处理器)提交审批结果,例如携带approvedBy、approvedAt的输出; - 审批完成后,第三个 HTTP 任务引用
${approval_ref.output.status}与${approval_ref.output.approvedBy},发送包含审批结果与审批人的后续通知——这体现了 Human Task 输出可被下游任务变量引用的完整数据链路。
小结
Human Task 用零参数的极简配置,实现了 Conductor 中最关键的可恢复执行能力之一:工作流可以"暂停但不丢失",等待任意时长的人工决策或外部事件。结合 Human 与 HumanTaskMapper 的源码可以看出,挂起逻辑完全无状态,状态推进全部经由 TaskResource 的更新端点或 QueueAdminResource 的队列更新端点完成,因此可以灵活对接审批系统、IM 通知、消息队列等多种外部信号源。对于需要 Human-in-the-Loop 的 Agent 工作流(例如让 LLM 产出结果后必须经人工确认再执行副作用操作),Human Task 是最直接的实现手段。
【免费下载链接】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),仅供参考