用 ADK 设计可执行的图工作流:并行、路由、人工审核与动态编排

2026-09-30 20 预计阅读时间: 1 分钟
来源: cloud.google.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:13 分钟

把所有工具交给一个 Agent,确实可以快速做出原型;但当任务涉及并行查询、明确的业务规则、人工审批和批量处理时,一个“大而全”的提示词很快就会变得难以测试、难以审计。

ADK Workflow 的价值在于把任务拆成节点和边:函数负责确定性规则,Agent 负责理解与生成,人负责处理例外,工作流则决定下一步运行什么。下面以退款流程为例,说明怎样选择静态图、动态节点以及两者的组合。

不要按提示词的书写顺序设计依赖

假设用户已经选择订单并点击“申请退款”。系统需要读取三类数据:

  • 订单信息;
  • 支付与拒付记录;
  • 客户过去一年的退款历史。

这三次查询都只依赖订单 ID,彼此之间没有依赖关系,因此应该同时启动。退款政策判断则需要全部查询结果,所以必须等待三条路径完成。

这就是典型的 fan-out 与 fan-in:

                       ┌─ fetch_order ────┐
order_id ── START ─────┼─ fetch_payment ──┼─ JOIN ── route_refund
                       └─ fetch_history ───┘

在 ADK 中,可以把异步函数直接作为节点,并通过嵌套元组和 JoinNode 表达这种依赖:

from google.adk import Workflow
from google.adk.workflow import JoinNode, START

async def fetch_order(node_input: str) -> dict:
    return {
        "order_id": node_input,
        "amount_usd": 42,
        "placed_days_ago": 8,
    }

async def fetch_payment(node_input: str) -> dict:
    return {"chargeback_open": False}

async def fetch_history(node_input: str) -> dict:
    return {"prior_refunds_12mo": 1}

join_case = JoinNode(name="join_case")

lookup_workflow = Workflow(
    name="refund_lookups",
    edges=[
        (START, (fetch_order, fetch_payment, fetch_history), join_case),
    ],
)

运行前需要安装与你项目兼容的 ADK 版本,并把示例查询替换成真实数据库或服务调用。不同版本的导入路径可能变化,应以当前版本的 Workflow API 为准。

JoinNode 的输出是以节点名为键的字典,因此后续节点可以读取:

order = node_input["fetch_order"]
payment = node_input["fetch_payment"]
history = node_input["fetch_history"]

一个容易忽略的细节是参数名。名为 node_input 的参数接收前一个节点的输出;其他参数名通常会绑定到 ctx.state。如果函数写成 fetch_order(order_id: str),运行时可能会尝试读取 ctx.state["order_id"],而不是自动接收上游输出。

把“理解意图”和“执行政策”交给不同角色

退款系统中至少存在两种性质完全不同的路由。

Agent 路由:理解用户想做什么

“鞋码不合适,能换一双吗?”和“请把钱退回来”表达的是不同意图。自然语言含义很难靠少量关键词稳定覆盖,可以让分类 Agent 输出结构化类别,再由一个小函数映射到图中已经声明的路径:

客户消息 → 意图分类 Agent → 路由函数
                         ├─ REFUND  → 退款流程
                         ├─ EXCHANGE → 换货流程
                         └─ CLARIFY → 追问用户

模型可以决定选择哪个类别,但不应随意发明新的流程目的地。图负责限定可到达的分支。

确定性路由:执行明确的退款规则

如果政策已经明确规定“存在未关闭拒付,或订单超过 30 天,必须拒绝”,那么不需要让模型重新解释规则。可以把政策写成纯函数:

def refund_policy(
    amount_usd: float,
    placed_days_ago: int,
    chargeback_open: bool,
    prior_refunds_12mo: int,
) -> str:
    if chargeback_open or placed_days_ago > 30:
        return "DENY"
    if amount_usd < 50 and prior_refunds_12mo < 3:
        return "AUTO_APPROVE"
    return "MANUAL_REVIEW"


def test_refund_policy() -> None:
    assert refund_policy(42, 8, False, 1) == "AUTO_APPROVE"
    assert refund_policy(120, 45, False, 0) == "DENY"
    assert refund_policy(80, 10, False, 1) == "MANUAL_REVIEW"


if __name__ == "__main__":
    test_refund_policy()
    print("policy tests passed")

保存为 policy.py 后可直接运行:

python policy.py

再用 ADK 节点把三份记录合并,并返回带路由名的事件:

from google.adk import Event


def route_refund(node_input: dict):
    case = {
        **node_input["fetch_order"],
        **node_input["fetch_payment"],
        **node_input["fetch_history"],
    }

    decision = refund_policy(
        amount_usd=case["amount_usd"],
        placed_days_ago=case["placed_days_ago"],
        chargeback_open=case["chargeback_open"],
        prior_refunds_12mo=case["prior_refunds_12mo"],
    )
    return Event(output=case, route=decision)

这种分工带来一个重要边界:代码决定是否退款,模型只负责把决定写成自然、清楚的客户通知。它既减少模型调用,也避免模型临时改变业务政策。

人工审核不是阻塞线程,而是暂停工作流

金额较高、历史情况复杂的请求可能需要人工判断。审核者可能几分钟后回复,也可能隔天处理,因此不能让服务器线程一直等待。

ADK 的 RequestInput 可以记录待处理请求并暂停当前运行。恢复后,带有 rerun_on_resume=True 的节点再次执行,并从 ctx.resume_inputs 中读取审核结果:

from pydantic import BaseModel, Field
from google.adk import Context, Event
from google.adk.events import RequestInput
from google.adk.workflow import node

REVIEW_ID = "refund:review"


class ReviewDecision(BaseModel):
    approve: bool = Field(description="是否批准退款")
    note: str = Field(default="", description="审核理由")


@node(rerun_on_resume=True)
async def escalate_to_human(ctx: Context, node_input: dict):
    answer = ctx.resume_inputs.get(REVIEW_ID)

    if answer is None:
        yield RequestInput(
            interrupt_id=REVIEW_ID,
            message=(
                f"是否批准订单 {node_input['order_id']} 的 "
                f"${node_input['amount_usd']} 退款?"
            ),
            payload=node_input,
            response_schema=ReviewDecision,
        )
        return

    reviewed_case = {
        **node_input,
        "reviewer_note": answer.get("note", ""),
    }
    yield Event(
        output=reviewed_case,
        route="AUTO_APPROVE" if answer["approve"] else "DENY",
    )

恢复机制要求工作流运行记录能够持久化,应用也需要提供一个安全的审核界面或回调入口。生产环境还应补上:

  • 审核身份与权限检查;
  • 决策时间、审核人和理由的审计日志;
  • 重复提交与超时处理;
  • 敏感字段脱敏;
  • 审核规则发生变化时的版本记录。

批量任务与动态调查是两种不同问题

如果输入是一批退款案例,而且每一项都执行同一段逻辑,可以使用 parallel_worker=True。它会为列表中的每个元素运行一次节点,并按原始顺序收集结果:

from google.adk import Workflow
from google.adk.workflow import START, node


@node(parallel_worker=True)
def review_case(node_input: dict) -> dict:
    decision = refund_policy(
        amount_usd=node_input["amount_usd"],
        placed_days_ago=node_input["placed_days_ago"],
        chargeback_open=node_input["chargeback_open"],
        prior_refunds_12mo=node_input["prior_refunds_12mo"],
    )
    return {
        "order_id": node_input["order_id"],
        "decision": decision,
    }


def collect_decisions(node_input: list[dict]) -> dict:
    return {"decisions": node_input}


batch_review = Workflow(
    name="batch_review",
    edges=[(START, review_case, collect_decisions)],
)

它和 fan-out 的区别值得明确:

模式 分发方式 汇总结果
Fan-out + JoinNode 不同节点执行不同工作 以节点名为键的字典
Parallel worker 同一节点处理列表中的每一项 保持输入顺序的列表

而“调查结果又产生新的调查任务”属于另一类问题。此时,下一步工作在输入到达前无法完整画出,更适合在动态节点中使用 Python 调度:

import asyncio
from google.adk.workflow import node

HANDLERS = {
    "AUTO_APPROVE": approve_notice,
    "MANUAL_REVIEW": escalate_to_human,
    "DENY": denial_notice,
}


@node(rerun_on_resume=True)
async def refund_flow(ctx, node_input):
    order, payment, history = await asyncio.gather(
        ctx.run_node(fetch_order, node_input, use_sub_branch=True),
        ctx.run_node(fetch_payment, node_input, use_sub_branch=True),
        ctx.run_node(fetch_history, node_input, use_sub_branch=True),
    )

    case = order | payment | history
    decision = refund_policy(
        amount_usd=case["amount_usd"],
        placed_days_ago=case["placed_days_ago"],
        chargeback_open=case["chargeback_open"],
        prior_refunds_12mo=case["prior_refunds_12mo"],
    )

    await ctx.run_node(
        HANDLERS[decision],
        case,
        use_as_output=True,
    )

use_sub_branch=True 可以把并发子节点的事件放在不同分支中,便于追踪;use_as_output=True 则把处理器结果作为父节点输出,避免重复发出结果。

动态节点恢复时可能从函数顶部重新执行,但已经完成的 ctx.run_node 调用可以从会话历史中取得记录结果。要利用这一点,应把支付、写库、发邮件等副作用放进独立子节点,而不是直接夹在父函数里,否则恢复执行可能造成重复扣款或重复通知。

静态图还是动态 Python:用可预见性来判断

选择方式时,可以先问一句:在请求到达之前,能否画出所有可能连接?不需要知道实际会走哪条路径,只需要知道有哪些合法分支、循环和阶段。

适合静态图的情况包括:

  • 并行查询结束后统一汇总;
  • 从固定规则中选择批准、拒绝或人工审核;
  • 固定的“生成—检查—修改”循环;
  • 需要让运营、审计或开发人员直观看到全部路径。

适合动态节点的情况包括:

  • 一个结果会产生数量不定的后续调查;
  • 需要递归探索关联交易或配送问题;
  • Python 循环、队列或并发控制比边列表更自然;
  • 模型可以建议调查方向,但代码必须限制深度、数量和预算。

实践中不必二选一。常见做法是用静态图描述主流程,只在“调查”这一阶段嵌入动态节点。这样既保留整体可视性,也能处理输入驱动的开放任务。

上线前的责任划分清单

一个稳健的退款工作流可以采用下面的边界:

  • 图:声明合法路径和阶段依赖;
  • 普通函数:执行金额、时间和次数等明确政策;
  • Agent:识别自然语言意图,撰写客户可读通知;
  • 人工审核者:处理政策无法自动覆盖的例外;
  • 动态节点:限制并调度运行时才出现的后续工作。

如果任务很小且开放,一个带工具的 Agent 可能已经足够。只要流程开始涉及并发、审计、审批、批处理或严格规则,就应该把控制权从提示词中移出来,放进可测试的函数和可检查的图中。图工程真正解决的不是“怎样多调用几个 Agent”,而是明确谁有权决定下一步,以及这个决定能否被测试、恢复和追踪。


相关推荐