ADK Python Workflow 人机协作实战:request_input_rerun 单节点审批模式深度解析
【免费下载链接】adk-pythonAn open-source, code-first Python toolkit for building, evaluating, and deploying sophisticated AI agents with flexibility and control.项目地址: https://gitcode.com/GitHub_Trending/ad/adk-python
本篇技术文章基于 ADK Python 官方示例 request_input_rerun,讲解如何在 ADK Workflows 中用RequestInput事件配合@node(rerun_on_resume=True)装饰器,把"向人类要输入"和"处理人类输入"合并进同一个节点。读完后,你将掌握两种 Human-in-the-Loop(人机协作)模式的区别、可直接运行的完整示例代码、事件日志级别的执行流程验证,以及该机制在框架源码中的底层实现原理。
两种模式对比:request_input 与 request_input_rerun
ADK 提供了两种在 Workflow 中暂停并等待人工输入的模式,二者模拟的都是同一个客服场景:AI 起草一封回复客户投诉的邮件,人工审核后决定放行、驳回或退回修改。核心差异在于恢复执行时人类输入是如何被消费的。
| 维度 | 标准模式request_input | 重执行模式request_input_rerun |
|---|---|---|
| 暂停方式 | 某节点 yieldRequestInput事件后暂停 | 同样 yieldRequestInput事件 |
| 恢复后输入去向 | 自动作为下一个节点的入参(node_input) | 通过执行Context的resume_inputs被同一个节点读取 |
| 节点数量 | 需要两个节点:一个请求输入,一个处理输入 | 单个节点自包含"请求 + 处理" |
| 边定义复杂度 | 链路多一跳 | 更简洁 |
从源码结构看,标准模式的实现位于 request_input 示例:request_human_review节点负责 yieldRequestInput,工作流恢复后由框架把用户输入作为参数喂给后续节点handle_human_review(node_input: str):
def request_human_review(draft: str): yield RequestInput( message=( "Please review the following draft email and provide 'approve'," f" 'reject', or feedback to revise.\n\n---\n{draft}\n---" ), ) def handle_human_review(node_input: str): if node_input == "reject": yield Event(route="rejected") elif node_input == "approve": yield Event(route="approved") else: yield Event(state={"feedback": node_input}, route="revise")而重执行模式把这两个职责合并进一个节点:该节点被@node(rerun_on_resume=True)装饰,恢复时框架重新运行这个节点本身,节点内部通过ctx.resume_inputs判断自己是"首次执行(应暂停)"还是"被恢复(应处理输入)"。这样请求与处理逻辑保持在同一处,状态更内聚,边定义也更短。
完整示例:客服邮件审批工作流
完整代码见 agent.py,整体结构如下(保留 Apache 2.0 版权头省略):
from google.adk import Agent from google.adk import Context from google.adk import Event from google.adk import Workflow from google.adk.events import RequestInput from google.adk.workflow import node def process_input(node_input: str): """Takes the initial customer complaint as input and sets it in the state.""" yield Event(state={"complaint": node_input, "feedback": ""}) draft_email = Agent( name="draft_email", instruction=""" Please write a polite, helpful response email to the following customer complaint: "{complaint}" If there is any feedback from the manager to revise the draft, please incorporate it: "{feedback?}" """, output_key="draft", ) @node(rerun_on_resume=True) def human_review(draft: str, ctx: Context): resume_input = ctx.resume_inputs.get("human_review") if not resume_input: yield RequestInput( interrupt_id="human_review", message=( "Please review the following draft email and provide 'approve'," f" 'reject', or feedback to revise.\n\n---\n{draft}\n---" ), ) return if resume_input == "reject": yield Event(route="rejected") elif resume_input == "approve": yield Event(route="approved") else: yield Event(state={"feedback": resume_input}, route="revise") def reject_email(): yield Event(message="Draft rejected.") def send_email(draft: str): yield Event(message="Draft approved and sent successfully.") root_agent = Workflow( name="request_input_rerun", edges=[ ("START", process_input, draft_email, human_review), ( human_review, { "revise": draft_email, "approved": send_email, "rejected": reject_email, }, ), ], )工作流图(继承自示例 README):
各节点职责:
process_input:把用户输入的投诉文本写入共享 state(complaint),并预置空feedback。draft_email:一个 LLMAgent节点。instruction 中的{complaint}从 state 插值,{feedback?}带?后缀表示可选——仅当 state 中存在feedback时才拼入提示词;output_key="draft"把模型输出写回 state 的draft键,供下游节点以参数形式引用。human_review:核心的"请求 + 处理"双职责节点,详见下文实现要点。reject_email/send_email:两条终态分支,分别输出"草稿被驳回"与"草稿已发送"。
官方 README 给出的典型测试输入包括:The delivery was a week late、I received the wrong item、My account was charged twice。
实现要点:四个关键步骤
以下四步完整继承自示例 README 的 How To 章节,是编写任何 rerun 型审批节点的标准流程。
1. 用@node(rerun_on_resume=True)装饰节点,并在签名中声明Context
from google.adk.workflow import node from google.adk import Context @node(rerun_on_resume=True) def human_review(draft: str, ctx: Context): # ...rerun_on_resume=True告诉框架:当该节点因RequestInput暂停、后续又带着用户输入恢复时,不要推进到下一个节点,而是重新执行本节点。函数签名必须包含 workflow 的Context,才能访问resume_inputs。
2. 通过ctx.resume_inputs判断是否处于恢复执行
resume_input = ctx.resume_inputs.get('human_review')resume_inputs是一个以interrupt_id为键的字典。示例中查找用的键"human_review"与 yieldRequestInput时声明的interrupt_id一致(恰好也与节点同名)。
3. 首次执行时 yieldRequestInput并立即 return 暂停
if not resume_input: yield RequestInput( interrupt_id="human_review", message="Please review the draft...", ) return # Important: Stop execution of this node for now注意return的作用:节点首次运行到这里后停止,工作流进入等待状态,事件流中会挂起一个待响应的中断。
4. 恢复执行时消费输入并产出路由事件
if resume_input == "reject": yield Event(route="rejected") elif resume_input == "approve": yield Event(route="approved") else: yield Event(state={"feedback": resume_input}, route="revise")任何非approve/reject的输入都被当作修改意见:写入 state 的feedback(供draft_email的{feedback?}插值),并沿revise路由回退给 LLM 节点重写,形成"草稿—审阅—修改"循环。
边定义:单节点带来的简化
由于human_review一个节点承担了请求与处理全部逻辑,边定义比标准模式更短:
Workflow( name="request_input_rerun", edges=[ ("START", process_input, draft_email, human_review), (human_review, {"revise": draft_email, "approved": send_email}), ], )对比标准模式的边定义:("START", process_input, draft_email, request_human_review, handle_human_review)——多出一个handle_human_review节点。注意两条边定义中{"revise": ..., "approved": ...}这个字典表示条件路由:目标节点取决于节点 yield 的Event(route=...)取值,示例还额外映射了"rejected": reject_email(见 agent.py)。
事件流验证:一次完整的"修改后批准"执行记录
test 会话记录 完整回放了一次phone broke投诉的处理过程,是验证该模式行为的最佳证据。11 个事件的演进如下:
e-1:用户输入phone broke。e-2:process_input节点产出 stateDelta{"complaint": "phone broke", "feedback": ""},nodeInfo.path为request_input_rerun@1/process_input@1。e-3:draft_email首次运行,写出第一版草稿(stateDelta.draft,路径后缀draft_email@1)。e-4:human_review首次执行,对外表现为一次名为adk_request_input的函数调用,参数含interruptId(fc-1)、message(附完整草稿正文)、payload、response_schema;事件携带longRunningToolIds: ["fc-1"],表示工作流在此挂起等待。e-5:用户以functionResponse返回shorter——既非 approve 也非 reject。e-6:human_review被重执行,产出route: "revise"及stateDelta: {"feedback": "shorter"}。e-7:draft_email@2(同一节点的第二次执行,路径中的@2即执行序号)基于反馈重写出一版更短的邮件。e-8:human_review@2再次 yieldRequestInput(fc-2),等待二次审阅——证明同一节点可在循环中反复中断/恢复。e-9:用户返回approve。e-10:节点第三次执行,产出route: "approved"。e-11:send_email输出Draft approved and sent successfully.。
最终 state 中保留了complaint、最新版draft和feedback: "shorter"。这条记录证实了三件事:RequestInput在事件层被序列化为adk_request_input函数调用(调用与响应的配对按interruptId匹配);rerun_on_resume节点确实被原样重执行(@2序号);循环中断复用同一interrupt_id是被支持的——RequestInput 模型注释明确说明"跨循环迭代复用同一interrupt_id是受支持的,框架按计数匹配函数调用与响应,但为了事件日志清晰仍建议每次迭代使用唯一 ID"。
源码级原理:resume_inputs 如何被注入与清理
结合框架源码(src/google/adk),可以还原该机制的调用链:
- 事件定义:
RequestInput是 Pydantic 模型,见 request_input.py。三个关键字段:interrupt_id(默认自动生成 UUID,显式指定可让恢复端确定性定位)、payload(可选的恢复载荷)、message(展示给用户的提示文案)。 - Context 侧:agents/context.py 中
Context.resume_inputs返回dict[str, Any],文档注释说明其键为 interrupt id;构造函数中初始化为空字典(self._resume_inputs = resume_inputs or {}),因此"查不到值 = 首次执行"这一判空逻辑是安全的。 - 工作流调度侧:_workflow.py 负责在恢复时把输入路由回原节点。从源码结构看,节点执行结果中的
resume_inputs会被存回节点状态(node_state.resume_inputs = result.resume_inputs or {},约 L664),恢复路径上再从节点状态取出并构造带resume_inputs的子 Context 重新执行该节点(约 L678-L691);重执行完成后会node_state.resume_inputs.clear()(约 L810-L811),保证输入只被消费一次、不会在下一次正常执行时"残留命中"。此外,若提供了resume_inputs却找不到可恢复的执行记录,框架会发出告警"resume_inputs provided but no recovered executions"(约 L254-L256),这是排查"输入没被接住"问题的第一信号。 - 装饰器侧:workflow/_node.py 的
node装饰器文档以@node(name='my_node', rerun_on_resume=True)为例,说明该参数会覆盖内部节点的rerun_on_resume属性(约 L113),使节点具备"恢复时原地重跑"的语义。
这套机制的语义是:中断状态只记录"哪个节点在等输入",恢复时把用户输入按interrupt_id塞回该节点的Context,节点自己决定下一步路由——这正是能把"请求"与"处理"合并进单节点的底层原因。
选型建议与适用前提
- 当"请求输入"与"解释输入"逻辑简单、只需一个后继节点做分发时,标准
request_input模式(两节点)更直白,可参考 request_input 示例。 - 当输入处理逻辑复杂(如示例中的三路分发 + 写 state + 回环修订),或你希望审阅上下文(草稿内容、判断逻辑)集中在一个函数内便于维护与测试时,
request_input_rerun的单节点模式更内聚。 - 适用前提:节点函数签名需声明
Context;RequestInput应显式指定interrupt_id以便resume_inputs精确命中;revise回环在语义上允许无限次循环,生产场景中可考虑在节点内对迭代次数加限。 - 所有代码基于当前仓库的 ADK Python 实现验证;示例可放入独立的 app 目录后用
adk run/adk web交互式调试(该方式适用于已按 ADK CLI 约定组织好agent.py的示例目录)。
延伸阅读
- 示例源码:contributing/samples/workflows/request_input_rerun/agent.py
- 事件回放数据:contributing/samples/workflows/request_input_rerun/tests/phone_broke.json
- 对照示例:contributing/samples/workflows/request_input/agent.py
- 核心实现:src/google/adk/events/request_input.py、src/google/adk/workflow/_node.py、src/google/adk/workflow/_workflow.py、src/google/adk/agents/context.py
【免费下载链接】adk-pythonAn open-source, code-first Python toolkit for building, evaluating, and deploying sophisticated AI agents with flexibility and control.项目地址: https://gitcode.com/GitHub_Trending/ad/adk-python
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考