第 6 章 动态工作流与人机协作
第 6 章 动态工作流与人机协作
第 5 章的退货审批流程,是用静态图画的:一条条边、一个个分支,清清楚楚。但真实业务不是总这么"听话"——有时候需要"循环直到成功",有时候需要"根据实时情况临时决定下一步",还有时候需要"暂停下来,等一个人拍板"。这一章,我们给图工作流装上三样东西:动态能力、重试能力、和人。
6.1 静态图的局限
第 5 章的 Workflow(edges=[...]) 适合"流程结构固定"的场景。但有几个场景,静态图会力不从心:
场景一:循环。 "客服自动重试发送验证码,最多 3 次"——这是一个循环,静态图要表达循环很笨重(得画回边,而且循环次数还是写死的)。
场景二:运行时才确定的分支。 "根据用户输入的商品列表,每个商品分别走一遍质检"——商品数量用户说了算,图的结构在写代码时根本不知道有多少个分支。
场景三:暂停等人。 "退款金额超过 5000 元,必须人工审批"——工作流要暂停,等一个真人确认,然后再继续。
这些场景的共同点是:流程的控制流不是静态的,而是需要代码逻辑动态决定的。这就是动态工作流(Dynamic Workflow) 的用武之地。
6.2 动态工作流:用代码表达控制流
6.2.1 核心思想
动态工作流是"基于图的工作流的更灵活、更强大的替代方案"。它的核心洞察是:
与其用图语言(edges)笨拙地表达循环和复杂分支,不如直接在普通 Python 代码里写控制流。
ADK 的动态工作流允许你:
- 用一个**编排器节点(orchestrator node)**承载主要逻辑
- 在编排器里用普通的
while、if/else、for、asyncio.gather表达控制流 - 用
ctx.run_node()调用其他子节点
这意味着什么?你可以用你早就熟悉的 Python 语法,写出任何复杂的流程——循环、条件、递归、并行,都不需要发明新的图语法。
6.2.2 @node 装饰器:把函数变成节点
动态工作流的核心 API 是 @node 装饰器,它把普通函数包装成工作流节点:
from google.adk import Context, Workflow
from google.adk.workflow import node
@node(name="check_quality")
def quality_check(node_input: str) -> str:
"""模拟质检:代码太短则返回问题描述,否则返回空字符串。"""
if len(node_input) < 20:
return "code is too short"
return ""
@node(name="fix_code")
def fix_code(node_input: str) -> str:
"""模拟修复:给代码追加注释。"""
return node_input + "\n# improved"
@node(name="summarize")
def summarize(node_input: str) -> str:
"""最终汇总。"""
return f"Final result:\n{node_input}"
# 编排器节点:用 Python 控制流编排子节点
@node(rerun_on_resume=True)
async def code_pipeline(ctx: Context, user_request: str) -> str:
code = user_request
# 循环:直到质量检查通过
finding = await ctx.run_node(quality_check, code)
while finding:
code = await ctx.run_node(fix_code, code)
finding = await ctx.run_node(quality_check, code)
# 条件分支:长度足够时执行汇总
if len(code) > 30:
code = await ctx.run_node(summarize, code)
return code
root_agent = Workflow(
name="root_agent",
edges=[("START", code_pipeline)],
)这个例子展示了动态工作流的三件事:
第一,@node 把函数变成节点。 quality_check、fix_code、summarize 都是普通 Python 函数,用 @node(name=...) 装饰后就变成了图上的节点。
第二,ctx.run_node() 是调用子节点的方式。 编排器节点接收 ctx: Context 参数,通过 await ctx.run_node(node, input) 执行某个子节点,返回值就是子节点的输出。
第三,控制流就是普通 Python。 while finding: 实现"循环直到质检通过",if len(code) > 30: 实现条件分支。没有任何图语法——就是 Python。
6.2.3 并行:用 asyncio.gather
动态工作流的并行极其优雅——因为你可以直接用 Python 的 asyncio.gather:
@node(rerun_on_resume=True)
async def parallel_check(ctx: Context, user_request: str) -> str:
# 同时执行三个独立的质检节点
results = await asyncio.gather(
ctx.run_node(check_quality, user_request),
ctx.run_node(check_security, user_request),
ctx.run_node(check_performance, user_request),
)
return "、".join(f"完成: {r}" for r in results)不需要并行的图语法,asyncio.gather 天然表达"这些任务同时跑"。这就是"用熟悉的编程结构,而不是图路由"的威力。
6.2.4 动态工作流 vs 静态图:怎么选
| 场景 | 用哪个 |
|---|---|
| 流程结构固定、步骤明确 | 静态图(第 5 章) |
| 有循环、复杂分支、运行时才知道结构 | 动态工作流 |
| 需要并行执行多个独立任务 | 动态工作流(asyncio.gather) |
| 简单流程、想直接看完整路径 | 静态图 |
记住:静态图适合"画出来的流程",动态工作流适合"算出来的流程"。两者都是 ADK 2.0 一等公民,可以组合使用。
6.3 人机协作(HITL):让流程暂停,等真人
6.3.1 为什么需要人机协作
有些决策,不应该让模型自动做。典型场景:
- 退款金额超过阈值:几千块的退款,该让真人经理拍板
- 高危操作:删除数据、发送对外通知,需要人确认
- 信息不完整:流程需要用户提供关键信息(比如收货地址),Agent 问不清
人机协作(Human-in-the-loop,HITL) 就是:工作流在某个节点暂停,等待一个真人输入,收到后再继续。
6.3.2 RequestInput:让节点请求人工输入
ADK 用 RequestInput 实现 HITL。核心机制:
- 节点
yield RequestInput(...)→ 工作流暂停 - 系统收到人工输入 → 工作流恢复
- 恢复后,用户的输入传递给后继节点
看官方的最小示例:
from google.adk.events import RequestInput
from google.adk import Workflow
def step1(): # 人工输入节点
yield RequestInput(message="Enter a number:")
def step2(node_input):
return node_input * 2
root_agent = Workflow(
name="root_agent",
edges=[('START', step1, step2)],
)注意 yield——人工输入节点是生成器函数,yield RequestInput(...) 会暂停执行并发出一个"需要人工输入"的事件。收到用户输入后,这个输入会被传给下一个节点(step2 收到用户输入的数字)。
RequestInput 有三个关键参数:
yield RequestInput(
message="请确认退款金额:¥5,000", # 展示给用户的提示
payload={"order_id": "ORD-20260901-001", "amount": 5000}, # 结构化数据
response_schema=ApprovalDecision, # 要求人工回复必须符合的结构
)message:展示给用户的提示文本payload:随请求发送的结构化数据,客户端可以渲染成表单response_schema:要求人工回复必须符合的数据结构(pydantic 模型)
6.3.3 rerunOnResume:恢复时节点怎么处理
恢复工作流时,被中断的节点怎么处理?由 rerun_on_resume 决定:
false(默认):被中断的节点不会重新执行。用户的回复直接作为该节点的输出,传递给后继节点。true:节点体从头重新执行。
文档特别强调:任何调用 ctx.run_node() 的节点,必须设置 rerun_on_resume=True——这样恢复时才能重新传递已缓存的子节点结果。
6.3.4 工具级的人工确认
除了图节点级的 HITL,ADK 还有一个工具级的确认机制:在工具配置上设置 require_confirmation: true,Agent 在执行该工具前会暂停,等待用户 yes/no 批准。
这个机制适合"单个高危工具"的审批,比如"删除用户"、"发送对外邮件"。它和图节点的 HITL 不同——图节点的 HITL 是流程级的(暂停整个工作流),工具确认是单工具级的。
6.4 自动重试与遥测:让流程更健壮
6.4.1 问题:LLM 和外部调用会失败
生产环境里,LLM 会超时、外部 API 会 5xx、工具会抛异常。一个健壮的流程,必须能自动重试。
6.4.2 RetryConfig:声明式重试
ADK 的 RetryConfig 让你声明式地配置重试策略:
from google.adk import Agent
from google.adk.config.retry_config import RetryConfig
agent = Agent(
name="robust_agent",
model="gemini-flash-latest",
retry_config=RetryConfig(
max_retries=3, # 最多重试 3 次
retry_on=("TimeoutError", "APIError"), # 这些异常才重试
backoff_factor=2.0, # 退避因子(第 1 次等 1s,第 2 次等 2s...)
initial_delay=1.0, # 首次重试前等待秒数
),
)RetryConfig 的关键参数:
max_retries:最大重试次数retry_on:哪些异常触发重试backoff_factor/initial_delay:指数退避策略,避免重试风暴
配合遥测(telemetry),ADK 会原生捕获异常并记录重试过程——这让我们能观测到"哪个环节经常失败、重试了几次"。(遥测的完整讨论见第 9 章。)
6.4.3 重试在动态工作流中的位置
在动态工作流里,重试可以组合进控制流。比如"重试直到成功"用 while 表达:
@node(rerun_on_resume=True)
async def robust_pipeline(ctx: Context, request: str) -> str:
attempts = 0
while True:
try:
result = await ctx.run_node(call_payment_api, request)
return result
except Exception as e:
attempts += 1
if attempts >= 3:
raise # 重试 3 次仍失败,向上抛
print(f"Attempt {attempts} failed: {e}, retrying...")这是动态工作流的又一大优势:重试逻辑就是普通 Python 的 try/while,不需要框架发明特殊语法。
6.5 主线项目:退货流程加入人工审核退款节点
6.5.1 业务需求
云销的退款规则里有一条:退款金额超过 ¥1,000 需要人工审批。这是典型的合规要求——大额退款不能由模型自动放行,必须经理确认。
我们在第 5 章的退货审批流程基础上,用动态工作流改造:
- 资格检查通过 → 计算退款金额
- 金额 ≤ ¥1,000 → 自动放行
- 金额 > ¥1,000 → 暂停,请求人工审批 → 审批通过后放行,拒绝则告知用户
6.5.2 完整代码
from google.adk import Context, Workflow, Event
from google.adk.events import RequestInput
from google.adk.workflow import node
from pydantic import BaseModel
# ---------- 数据契约 ----------
class RefundRequest(BaseModel):
order_id: str
category: str
amount: float # 退款金额
class ApprovalDecision(BaseModel):
approved: bool
approver: str
# ---------- 节点 1:检查资格并计算金额(确定性) ----------
@node(name="check_and_calc")
def check_and_calc(node_input: RefundRequest) -> dict:
# 模拟:根据订单计算退款金额
mock_amount = {"ORD-20260901-001": 599.0, "ORD-20260901-002": 2599.0}
amount = mock_amount.get(node_input.order_id, node_input.amount)
eligible = amount > 0
return {"eligible": eligible, "amount": amount, "order_id": node_input.order_id}
# ---------- 节点 2:人工审批节点(HITL) ----------
def human_approval(node_input: dict):
"""大额退款需要人工审批。"""
yield RequestInput(
message=(
f"订单 {node_input['order_id']} 申请退款 ¥{node_input['amount']},"
f"超过 ¥1,000 需要人工审批。请确认:"
),
payload=node_input,
response_schema=ApprovalDecision,
)
# ---------- 节点 3:根据审批结果分支 ----------
def approval_router(node_input):
"""根据审批决定路由。node_input 是人工审批的回复(ApprovalDecision)。"""
if node_input.approved:
return Event(route="approved")
return Event(route="rejected")
# ---------- 节点 4a:放行 ----------
def approve_and_refund(node_input) -> Event:
return Event(message=f"退款审批通过,已发起退款 ¥{node_input.amount}。退款将在 1-3 个工作日到账。")
# ---------- 节点 4b:拒绝 ----------
def reject_refund(node_input) -> Event:
return Event(message="很抱歉,您的退款申请未通过人工审批。如有疑问请联系客服。")
# ---------- 动态编排器:把流程串起来 ----------
@node(rerun_on_resume=True)
async def refund_pipeline(ctx: Context, user_request: RefundRequest) -> Event:
# 步骤 1:检查资格 + 计算金额
check = await ctx.run_node(check_and_calc, user_request)
if not check["eligible"]:
return Event(message="您的订单不符合退款条件。")
# 步骤 2:金额判断——这就是"把 if-else 还给代码"
if check["amount"] <= 1000:
# 小额:自动放行,走审批节点但直接通过
return await ctx.run_node(approve_and_refund, check)
# 大额:需要人工审批
decision = await ctx.run_node(human_approval, check)
# 决策路由(人工审批的回复就是决策)
if decision.approved:
return await ctx.run_node(approve_and_refund, check)
else:
return await ctx.run_node(reject_refund, check)
root_agent = Workflow(
name="refund_with_hitl_workflow",
edges=[("START", refund_pipeline)],
)6.5.3 流程是怎么走的
我们跟踪一次"大额退款"的完整流程:
refund_pipeline收到用户请求(比如订单 ORD-20260901-002,退款 ¥2,599)ctx.run_node(check_and_calc, ...)→ 算出金额 2599,eligibleif check["amount"] <= 1000:→ 2599 > 1000,进入大额分支(这一行 if,就是"把判断留给代码")ctx.run_node(human_approval, ...)→ 工作流暂停,发出RequestInput等人工审批- 经理在后台看到"订单 002 申请退款 ¥2599,请确认",点击通过
- 工作流恢复,
decision.approved=True ctx.run_node(approve_and_refund, ...)→ 放行退款- 流程结束
对比纯 Agent 方案:纯 Agent 可能在"要不要人工审批"这一步自己决定"金额大,我直接放行吧"——它没有权限意识。而图工作流 + HITL 把"大额必须人工"变成了不可绕过的流程约束。
6.5.4 为什么这是安全闸门
这里要强调一个贯穿全书的概念(第 11 章还会展开):HITL 不只是"体验",更是安全机制。
当"超过 1000 元必须人工审批"被固化在工作流里,它就是一个硬约束——无论模型怎么想,流程都会在审批节点停下。这是代码层面的强制,不是"提醒模型要小心"的软约束。对于金融、数据、安全相关的操作,这种硬约束是不可或缺的。
「为什么 ADK 这样设计」:把 if-else 还给代码,把判断留给 LLM
第 5 章我们第一次提出"把 if-else 还给代码,把判断留给 LLM"。这一章,这个哲学被贯彻到了极致。我们用三个例子回顾:
例一:金额判断。
if check["amount"] <= 1000:
# 自动放行这是 if-else——代码,不是 LLM。模型不需要"推理"¥500 和 ¥1200 哪个该审批,代码直接决定。确定、零成本、可测试。
例二:循环重试。
while finding:
code = await ctx.run_node(fix_code, code)
finding = await ctx.run_node(quality_check, code)这是 while 循环——代码,不是 LLM。不用让模型"决定要不要再修一次",循环条件自动控制。不会死循环,因为条件是代码定义的。
例三:人工审批。
if check["amount"] > 1000:
decision = await ctx.run_node(human_approval, check)这是流程暂停——代码强制。不依赖模型"自觉"去请示,流程结构就保证了大额必然停下等人。
这三个例子合起来,就是 ADK 2.0 动态工作流的完整哲学:凡是能确定的,都用代码确定;只有真正需要理解和生成的,才留给 LLM;而关键的、有风险的决策,让人来拍板。
这一章结束,你手里已经有三件工具:多智能体(第 4 章,处理模糊分发)、图工作流(第 5 章,确定性流程)、动态工作流 + HITL(本章,复杂控制流 + 人机协作)。到第 7 章,我们把这些能力接上整个开放生态——模型、工具、协议。
本章小结
- 静态图的局限:循环、运行时确定的分支、暂停等人,用图语法表达很笨重
- 动态工作流:用
@node+ctx.run_node()+ 普通 Python 控制流(while/if/asyncio.gather)编排,熟悉又强大 - 并行是免费的:
asyncio.gather天然表达并行,不需要图语法 - HITL 的核心 API:节点
yield RequestInput(message, payload, response_schema)→ 暂停 → 恢复 - rerunOnResume:
false(默认)恢复时节点不重跑,用户输入直接传给后继节点;true节点重跑(调用ctx.run_node的节点必须设 true) - RetryConfig:声明式重试(max_retries / retry_on / backoff),动态工作流里也可用 try/while 手写
- 主线项目:退货流程加入"超 ¥1000 人工审批",HITL 成为不可绕过的安全闸门
- 核心哲学:确定的事用代码,智能的事用 LLM,风险的事让人
练习
- 改阈值:把人工审批阈值从 ¥1000 改成 ¥5000,验证改动只有一行 if,流程自动生效。
- 加重试:给
check_and_calc节点加一个"偶尔失败"的模拟(比如 30% 概率抛异常),用RetryConfig或 try/while 让它自动重试,观察行为。 - 并行审批:假设大额退款需要"财务"和"主管"两个人同时审批,用
asyncio.gather并行执行两个human_approval节点,实现会签。 - 理解 rerun_on_resume:把
refund_pipeline的rerun_on_resume改成False,跑一次大额流程,看会发生什么,理解这个参数的意义。
下一章预告:第 7 章,我们把视角从"一个框架内部"拉回"整个生态"。模型中立(不绑 Google)、MCP、A2A、AG-UI、AGENTS.md——ADK 凭什么说自己是"开放生态"里的最佳选择。