Keyboard shortcuts

Press or to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

第 10 期:子图与多 agent

客人有时候不是问一个订单,是一次甩过来好几个:“帮我看看这三个订单能不能改期”。 这一期要接住的其实是两件不相关的事,只是恰好在同一个场景里都用得上。

第一件事:查一个订单要做什么,不是写死的。“能不能改期“要先查订单拿出行日期, 再查政策;“退款政策是什么“可能压根不用管出行日期。这段“一句人话该翻译成 查哪几步、传什么参数“的判断,值得封装成一个独立的、可以单独测试的小 agent, 而不是在主流程里为每种问法各写一段 if/else。子图做的是这件事——跟要不要 并行没有关系,只查一个订单也用得上它。

第二件事:三个订单之间没有先后关系,没理由一个查完再查下一个。Send 做的是这件事——父图按状态里已经确定的订单号列表,一次性并行派发好几份任务, 不用模型自己决定调几次工具;它派发的对象可以是子图,也可以是任何一个普通 节点,跟“是不是子图“没有关系。

《笨办法学 Agent》上册练习 19 到 21 讲的是这两件事的合体:sub_agent 隔离 上下文(对应子图),练习 20 让模型并行调用好几次 sub_agent(对应 Send)。 这一期两个都用上,是因为这个场景两头都要——不代表子图天生是用来并行的, 也不代表并行天生要靠子图。

敲进去

第 10 期的代码在 code/ep10/,在第 9 期的基础上加四样:一个独立的 订单核对子图、一个新工具、状态里三个新字段、图里两条新边。

先是子图本身,subagent.py,一个只有两个工具(get_orderget_policy) 的迷你 agent,自己的状态、自己的循环,编译成一个独立对象:

class SubState(TypedDict):
    messages: Annotated[list[AnyMessage], add_messages]

def build_order_subgraph():
    llm = chat_model().bind_tools(SUB_TOOLS)

    def agent(state: SubState) -> dict:
        reply = llm.invoke([SystemMessage(SUB_PROMPT), *state["messages"]])
        return {"messages": [reply]}

    builder = StateGraph(SubState)
    builder.add_node("agent", agent)
    builder.add_node("tools", ToolNode(SUB_TOOLS))
    builder.add_edge(START, "agent")
    builder.add_conditional_edges("agent", _route, {"tools": "tools", END: END})
    builder.add_edge("tools", "agent")
    return builder.compile()

跟第 3 期的主图长得一模一样——节点、边、条件路由,一个字没多。区别只在 状态:SubState 只有 messages,没有 loaded_skills、没有跨会话记忆, 因为它不需要。这就是“子图“:一张完整的图,编译好,随时可以被另一张图 当成一个可调用的单元使用,不需要跟父图共用状态定义。不接 checkpointer—— 每次调用都是从零开始、跑完即弃的一次性任务,不需要记住上一次。

这里为什么要包一层 agent,而不是在 lookup_order 里直接写死 get_order()get_policy(topic="reschedule")?因为 topic 参数 只能是 reschedule/refund/usage 三选一,而 check_orders 收到的 aspect 是一句自然语言,“能不能改期”“退款政策是什么”“这个商品怎么 使用”——把这句话翻成正确的 topic,就是这个小 agent 存在的理由。 单独测过三个不同的 aspect,它每次都选对了:

KL-778 能不能改期        → get_order → get_policy(topic="reschedule")
KL-901 退款政策是什么     → get_order → get_policy(topic="refund")
KL-315 这个商品怎么使用   → get_order → get_policy(topic="usage")

这段“翻译“逻辑被封进子图,跟这次调用是不是并行的没有任何关系——哪怕 check_orders 一次只查一个订单、完全不走 Send,这个子图也照样有用。

然后是新工具,tools.py 里的 check_orders

@tool
def check_orders(order_ids: list[str], aspect: str, runtime: ToolRuntime) -> Command:
    """一次核对多个订单,每个订单派一个隔离的子 agent 去查,互不干扰、并行执行。
    order_ids 传订单号列表(比如 ["KL-778", "KL-901"]),至少两个才用这个工具——
    只查一个订单用 get_order/get_policy 就够了。aspect 是要核对的角度,一句话,
    比如"能不能改期"、"退款政策是什么"。"""
    ack = ToolMessage(
        f"已经并行核对 {len(order_ids)} 个订单,结果在后面那条消息里。",
        tool_call_id=runtime.tool_call_id,
        name="check_orders",
    )
    return Command(update={"messages": [ack], "order_ids": order_ids, "check_aspect": aspect})

这个工具自己不查任何数据。它只做一件事:把订单号列表和要核对的角度写进 状态,交给图的路由去决定接下来怎么并行处理——真正的查询发生在下面的 lookup_order 节点里。ack 这条 ToolMessage 是当场就要给的:OpenAI 协议要求每一次工具调用必须对应一条结果,不能“先欠着,扇出跑完再补“, 所以这里先垫一句“结果在后面“,把这次调用应付过去。

状态新增三个字段,state.py

class AgentState(TypedDict, total=False):
    messages: Annotated[list[AnyMessage], add_messages]
    loaded_skills: Annotated[list[str], operator.add]
    order_ids: NotRequired[list[str]]
    order_id: NotRequired[str]
    check_aspect: NotRequired[str]
    order_reports: Annotated[dict[str, str], operator.or_]

order_reports 那个 reducer 值得多看一眼:operator.or_,字典合并,不是 operator.add(列表追加)。原因是并行——lookup_order 会被 Send 同时 派出去好几份,几份并行调用会在同一个 step 里各自写一次这个字段, LangGraph 要求“同一个 key 被并发写“必须有 reducer 说明怎么合并,否则 直接报错。用字典还有一个好处:每个 lookup_order 只知道自己查的那一个 订单号,写 {订单号: 结果} 这样的单键字典,多份并行结果按键合并,互不 覆盖;旧一轮查询留下的键也不会干扰下一轮——读的时候只按当次的 order_ids 去字典里取值,取不到的键不管是不是这一轮的,自然被忽略。

最后是图,graph.pytools 节点之后不再是一条固定边,是一次路由判断:

def route_after_tools(state: AgentState) -> str | list[Send]:
    order_ids = state.get("order_ids")
    if not order_ids:
        return "agent"
    aspect = state.get("check_aspect", "")
    return [Send("lookup_order", {"order_id": oid, "check_aspect": aspect}) for oid in order_ids]


async def lookup_order(state: AgentState) -> dict:
    order_id = state["order_id"]
    aspect = state["check_aspect"]
    result = await order_subgraph.ainvoke(
        {"messages": [SystemMessage(f"只回答这一句:订单 {order_id},{aspect}")]}
    )
    return {"order_reports": {order_id: result["messages"][-1].content}}


def aggregate(state: AgentState) -> dict:
    reports = state.get("order_reports", {})
    lines = [f"{oid}:{reports.get(oid, '没查到结果')}" for oid in state.get("order_ids", [])]
    note = HumanMessage("[并行核对结果]\n" + "\n".join(lines))
    return {"messages": [note], "order_ids": []}

平时(没人调 check_ordersroute_after_tools 照旧直接回 agent,跟 第 3-9 期一模一样。一旦 order_ids 被写进状态,它改用 Sendlookup_order 连发 N 份任务——扇出几份,由状态里已经确定的数据决定, 不是模型当场调了几次工具。所有分支跑完,aggregate 把结果拼成一条新 消息追加进历史(不是塞进 check_orders 那条 ToolMessage 里,那条早就 用掉了),再回到 agent,模型这才第一次看到真正的结果。aggregate 顺手把 order_ids 清空,避免下一次路由判断误读上一轮的旧数据。

边接起来:

builder.add_node("lookup_order", lookup_order)
builder.add_node("aggregate", aggregate)
builder.add_conditional_edges("tools", route_after_tools, ["agent", "lookup_order"])
builder.add_edge("lookup_order", "aggregate")
builder.add_edge("aggregate", "agent")

跑起来

cd code
uv run python -m ep10.main wang t1 "帮我一起查一下 KL-778、KL-901、KL-315 这三个订单能不能改期"

你应该看到什么

DeepSeek:

[agent] 要调 check_orders({'order_ids': ['KL-778', 'KL-901', 'KL-315'], 'aspect': '能不能改期'}) (prompt_tokens=1719)
[tools] check_orders 返回:已经并行核对 3 个订单,结果在后面那条消息里。
[lookup_order] KL-778 开始(t=3458848.00)
[lookup_order] KL-901 开始(t=3458848.00)
[lookup_order] KL-315 开始(t=3458848.00)
[lookup_order] KL-315 结束(耗时 2.64s)
[lookup_order] KL-778 结束(耗时 2.80s)
[lookup_order] KL-901 结束(耗时 4.67s)
[aggregate] HumanMessage 返回:[并行核对结果]
KL-778:可以改期,该订单出行日前3天可免费改期一次。
KL-901:能,订单 KL-901 出行日为2026-09-05,出行日前3天(9月2日前)可免费改期一次。
KL-31
[agent] 回答:三个订单核对结果如下:

- **KL-778**:可以改期,出行日前3天可免费改期一次。
- **KL-901**:该单出行日是 2026-09-05,免费改期须在 9/2 前办理。今天已是 9/3,**已过免费窗口**,如需改期可能涉及费用。
- **KL-315**:可以改期,任意时间免费,只要在原出发前完成。

需要我帮您办理其中某张订单的改期吗? (prompt_tokens=1904)

三个 [lookup_order] ... 开始 打在同一个时刻,结束 分别隔了 2.64s、 2.80s、4.67s——不是排队一个个查完再报告,是真的同时发出去、各自等 各自的网络往返。而且模型没有停在“复述三条结果“:它注意到 KL-901 的 免费改期窗口是出行日前 3 天(9 月 2 日前),今天已经是 9 月 3 日, 主动指出这单已经过了免费窗口——这条推理是 agent 节点拿到并行结果 之后自己算出来的,不是任何一个 lookup_order 单独查出来的,三份 互相隔离的报告拼在一起,才够信息让它做这个判断。

本机 Ollama(qwen3:4b-instruct)跑同一个任务:

[lookup_order] KL-778 开始(t=3458877.07)
[lookup_order] KL-901 开始(t=3458877.08)
[lookup_order] KL-315 开始(t=3458877.08)
[lookup_order] KL-315 结束(耗时 7.51s)
[lookup_order] KL-778 结束(耗时 8.82s)
[lookup_order] KL-901 结束(耗时 10.52s)
[aggregate] HumanMessage 返回:[并行核对结果]
KL-778:可以免费改期一次,但仅限出行日前3天内。
KL-901:出行日前3天可免费改期一次,改期后不可再改;出行日前3天内不支持改期。
KL-315:出行日前3天可免费改期一次
[agent] 回答:KL-778:可免费改期一次,限出行日前3天内。  
KL-901 和 KL-315:出行日前3天可免费改期一次,改期后不可再改,日前3天内不支持改期。 (prompt_tokens=1786)

也是三个同时开始,但耗时依次拉长到 7.51s / 8.82s / 10.52s——三个请求 共享同一台机器上的同一个模型进程,确实在并发处理,只是本地算力撑不住 三份同时推理,越往后排队等得越久,不像调远程 API 那样几乎不受限。 最后一步也看出本机小模型的弱项:它没有做 KL-901 那条日期判断,还把 KL-901 和 KL-315 的结论混在一句话里,读起来像是同一个答案——三份 报告本身是对的,综合它们、挑出真正要紧的那一条,是本机模型这次没做好 的部分。

发生了什么

子图封装的是“这一步该怎么查“,跟并不并行是两件事。 order_subgraph 真正的价值在“敲进去“里那三行翻译表:把一句自由格式的 aspect 变成正确 的 topic 参数,这段判断不管 lookup_order 是被 Send 并行调用三次, 还是被顺序调用一次,都一样成立、一样值得封装。这一期把它们放进同一个 场景,容易让人以为“子图就是用来撑并行的“——但把 Send 那部分整个拿掉, check_orders 改成一个个顺序调用 order_subgraph,这个子图不会有 任何变化,一行代码都不用改。子图这个概念负责一段自成一体的判断逻辑, 并行是另一层完全独立的调度决定,两者能装进同一个场景,不代表它们 互相依赖。

扇出几份,是状态说了算,不是模型说了算。 上册练习 20 的并行扇出, 决定权在模型手里:它在一条消息里连续调用几次 sub_agent,调几次就扇出 几份。这一期反过来:模型只需要把“查哪几个“这件事做对——check_ordersorder_ids 参数——扇出几份是 route_after_tools 读这个列表的长度 决定的,不是模型运行时想调几次就调几次。两种设计都能做到并行,差别在 “并行的份数“这个决定权放在哪一层。

子图接进父图有两种官方写法,选哪种由 state 是不是共享决定。 一种是 builder.add_node("name", 编译好的图),直接把子图挂成父图的一个节点—— 这种写法要求父子共用同一套 state 的部分字段,子图读写的是父图状态里 同名的那些 channel,官方文档举的例子就是多个 agent 通过共享的 messages 字段互相看得见对方写了什么。另一种是在一个普通函数节点里手动 subgraph.invoke(...),父图传什么进去、子图吐什么出来,全靠这个函数 自己转换——官方文档写得很直接:“当父子的 state schema 不共享,或者 需要在两者之间转换状态时“用这种写法。这一期的 order_subgraph 用的是 第二种:它的 SubState 只有 messages,跟父图 AgentState 完全不共享 字段,而且不共享是故意的——它不该看见客人和主 agent 聊过什么,只该 看见分给它的那一句任务,这正是上册练习 19 的 sub_agent 要的隔离。如果 改用第一种写法(直接 add_node),子图会自动接上父图的 messages (跟父图共享同一路对话历史),隔离这条要求反而没法满足——选第二种 不是图偷懒,是隔离这个目标本身要求状态不共享。

Send 解决的是“扇出几份“,不是“扇出去干什么“。 真正干活的还是 lookup_order 里那次 order_subgraph.ainvoke()——Send 只负责把 {"order_id": oid, "check_aspect": aspect} 这份参数派给 lookup_order 的一次独立调用,派几份、派给谁,是路由函数的事,跟“派去之后具体干什么“ 完全解耦。这跟上册的注册表是同一个设计精神:谁负责决定,谁负责执行, 分成两层,互不掺和。

reducer 的字典合并,为的是让旧数据自动失效。 order_reportsoperator.or_ 而不是 operator.add,不只是因为字典比列表更适合按键 去重——更要紧的是,aggregate 读取时只按当次order_ids 去字典 里取值。哪怕字典里因为上一轮查询还躺着别的订单号,这一轮也不会误读到, 不需要专门写一段“清空上一轮残留“的逻辑,过滤本身就是清空。

OpenAI 协议的硬约束,决定了这条链路必须分两条消息说完。 一次工具 调用必须当场配一条 ToolMessage,晚到不行——所以 check_orders 自己 先垫一句“结果在后面“,占住这个 tool_call_id;真正的结果由 aggregate 之后追加一条独立的 HumanMessage,不去动前面那条已经用掉的工具结果。 模型看到的是两条连续的消息,不是一条分两次写的消息,但对它而言效果 一样——agent 节点只在两条消息都齐了之后才会被重新唤醒。

常见问题

lookup_order 里手动 .ainvoke() 一个图,这也算子图吗?看着不像 add_node 那种“官方“用法。 算,而且官方文档明确写了这是两种子图 接法之一。add_node("name", 编译好的图) 那种写法要求父子共用 state 的字段——子图读写的是父图状态里同名的 channel,官方举的例子是几个 agent 通过共享的 messages 字段互通有无。另一种就是这一期用的:在 普通函数节点里手动 subgraph.invoke(...),父子状态互不相通,转换 全靠这个函数自己写。官方原话是“当父子的 state schema 不共享,或者 需要在两者之间转换状态时“用这种写法——order_subgraphSubState 只有 messages,跟父图 AgentState 一个字段都不共享,而且这是故意的: 它不该看见客人跟主 agent 聊过什么。选哪种子图接法不是随手定的,是 “要不要隔离“这个需求本身决定的。

为什么不干脆在 check_orders 内部自己 asyncio.gather 三次调用, 不用 Send 可以,效果也是并行。区别在图层面看不看得见:走 Send, 每个 lookup_order 是图里一个独立的 step,理论上可以单独打断、单独接 checkpoint、单独重跑;工具函数内部自己 gather,三次调用被封在一个 不可分割的黑盒里,图只知道“这次工具调用花了多久“,看不见里面发生了 什么。这一期选 Send 是为了让“子图“和“并行“都在图这一层看得见,不是 说工具内部并行不对。

order_subgraph 为什么不接 checkpointer? 它是一次性任务——查完 就交结果,不需要记住上一次问的是什么。如果以后要给子图也接上断点续传, 需要给每次调用单独分配 thread_id,避免几个并行分支互相踩到同一个 存档,这一期没做。

check_orders 要不要像 cancel_order 那样过一道人工审批? 不用。 它只读订单和政策,不改变任何状态,跟 get_order/get_policy 是同一 个风险等级,不需要过第 5 期那道 interrupt 闸门。

只有两个订单也能用,为什么工具描述里写“至少两个才用“? 提醒模型 别在小事上绕远路:一个订单直接 get_order/get_policy 就是一次 调用,走 check_orders 反而多了一趟子图初始化的开销,划不来。

加分练习

  1. order_subgraph 也接一个类似 cancel_order 的写操作工具,然后 想清楚:如果某个并行分支里的子任务需要人工审批(interrupt), 其余还在跑的分支要不要等它?现在的实现完全没考虑这个,是故意留白 的复杂度。
  2. check_aspect 从“所有订单问同一句话“换成每个订单能问不同问题 的结构,改一下 check_orders 的参数形状,让 order_ids 和问题 一一对应而不是共用一个 aspect
  3. 用真机测一组对照实验:把 lookup_order 改成顺序 await 而不是 靠 Send 并行,跑同一个三订单任务,对比总耗时——用真实数字量化 “并行“到底省了多少秒,不要只信直觉。
  4. lookup_order 里的 ainvoke 换成 astream,看看父图能不能 观察到子图内部的中间步骤——想清楚为什么现在这个设计选择了看不见。