第 6 章 Google ADKDynamic WorkflowHITL

第 6 章 动态工作流与人机协作

第 6 章 动态工作流与人机协作

第 5 章的退货审批流程,是用静态图画的:一条条边、一个个分支,清清楚楚。但真实业务不是总这么"听话"——有时候需要"循环直到成功",有时候需要"根据实时情况临时决定下一步",还有时候需要"暂停下来,等一个人拍板"。这一章,我们给图工作流装上三样东西:动态能力、重试能力、和人。

6.1 静态图的局限

第 5 章的 Workflow(edges=[...]) 适合"流程结构固定"的场景。但有几个场景,静态图会力不从心:

场景一:循环。 "客服自动重试发送验证码,最多 3 次"——这是一个循环,静态图要表达循环很笨重(得画回边,而且循环次数还是写死的)。

场景二:运行时才确定的分支。 "根据用户输入的商品列表,每个商品分别走一遍质检"——商品数量用户说了算,图的结构在写代码时根本不知道有多少个分支。

场景三:暂停等人。 "退款金额超过 5000 元,必须人工审批"——工作流要暂停,等一个真人确认,然后再继续。

这些场景的共同点是:流程的控制流不是静态的,而是需要代码逻辑动态决定的。这就是动态工作流(Dynamic Workflow) 的用武之地。

6.2 动态工作流:用代码表达控制流

6.2.1 核心思想

动态工作流是"基于图的工作流的更灵活、更强大的替代方案"。它的核心洞察是:

与其用图语言(edges)笨拙地表达循环和复杂分支,不如直接在普通 Python 代码里写控制流。

ADK 的动态工作流允许你:

  • 用一个**编排器节点(orchestrator node)**承载主要逻辑
  • 在编排器里用普通的 whileif/elseforasyncio.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_checkfix_codesummarize 都是普通 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。核心机制:

  1. 节点 yield RequestInput(...) → 工作流暂停
  2. 系统收到人工输入 → 工作流恢复
  3. 恢复后,用户的输入传递给后继节点

看官方的最小示例:

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 流程是怎么走的

我们跟踪一次"大额退款"的完整流程:

  1. refund_pipeline 收到用户请求(比如订单 ORD-20260901-002,退款 ¥2,599)
  2. ctx.run_node(check_and_calc, ...) → 算出金额 2599,eligible
  3. if check["amount"] <= 1000:2599 > 1000,进入大额分支(这一行 if,就是"把判断留给代码")
  4. ctx.run_node(human_approval, ...)工作流暂停,发出 RequestInput 等人工审批
  5. 经理在后台看到"订单 002 申请退款 ¥2599,请确认",点击通过
  6. 工作流恢复,decision.approved=True
  7. ctx.run_node(approve_and_refund, ...) → 放行退款
  8. 流程结束

对比纯 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) → 暂停 → 恢复
  • rerunOnResumefalse(默认)恢复时节点不重跑,用户输入直接传给后继节点;true 节点重跑(调用 ctx.run_node 的节点必须设 true)
  • RetryConfig:声明式重试(max_retries / retry_on / backoff),动态工作流里也可用 try/while 手写
  • 主线项目:退货流程加入"超 ¥1000 人工审批",HITL 成为不可绕过的安全闸门
  • 核心哲学:确定的事用代码,智能的事用 LLM,风险的事让人

练习

  1. 改阈值:把人工审批阈值从 ¥1000 改成 ¥5000,验证改动只有一行 if,流程自动生效。
  2. 加重试:给 check_and_calc 节点加一个"偶尔失败"的模拟(比如 30% 概率抛异常),用 RetryConfig 或 try/while 让它自动重试,观察行为。
  3. 并行审批:假设大额退款需要"财务"和"主管"两个人同时审批,用 asyncio.gather 并行执行两个 human_approval 节点,实现会签。
  4. 理解 rerun_on_resume:把 refund_pipelinererun_on_resume 改成 False,跑一次大额流程,看会发生什么,理解这个参数的意义。

下一章预告:第 7 章,我们把视角从"一个框架内部"拉回"整个生态"。模型中立(不绑 Google)、MCP、A2A、AG-UI、AGENTS.md——ADK 凭什么说自己是"开放生态"里的最佳选择。