第 10 期:子图与多 agent
客人有时候不是问一个订单,是一次甩过来好几个:“帮我看看这三个订单能不能改期”。 这一期要接住的其实是两件不相关的事,只是恰好在同一个场景里都用得上。
第一件事:查一个订单要做什么,不是写死的。“能不能改期“要先查订单拿出行日期, 再查政策;“退款政策是什么“可能压根不用管出行日期。这段“一句人话该翻译成 查哪几步、传什么参数“的判断,值得封装成一个独立的、可以单独测试的小 agent, 而不是在主流程里为每种问法各写一段 if/else。子图做的是这件事——跟要不要 并行没有关系,只查一个订单也用得上它。
第二件事:三个订单之间没有先后关系,没理由一个查完再查下一个。Send
做的是这件事——父图按状态里已经确定的订单号列表,一次性并行派发好几份任务,
不用模型自己决定调几次工具;它派发的对象可以是子图,也可以是任何一个普通
节点,跟“是不是子图“没有关系。
《笨办法学 Agent》上册练习 19 到 21 讲的是这两件事的合体:sub_agent 隔离
上下文(对应子图),练习 20 让模型并行调用好几次 sub_agent(对应 Send)。
这一期两个都用上,是因为这个场景两头都要——不代表子图天生是用来并行的,
也不代表并行天生要靠子图。
敲进去
第 10 期的代码在 code/ep10/,在第 9 期的基础上加四样:一个独立的
订单核对子图、一个新工具、状态里三个新字段、图里两条新边。
先是子图本身,subagent.py,一个只有两个工具(get_order、get_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.py,tools 节点之后不再是一条固定边,是一次路由判断:
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_orders)route_after_tools 照旧直接回 agent,跟
第 3-9 期一模一样。一旦 order_ids 被写进状态,它改用 Send 给
lookup_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_orders
的 order_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_reports 用
operator.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_subgraph 的 SubState
只有 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 反而多了一趟子图初始化的开销,划不来。
加分练习
- 给
order_subgraph也接一个类似cancel_order的写操作工具,然后 想清楚:如果某个并行分支里的子任务需要人工审批(interrupt), 其余还在跑的分支要不要等它?现在的实现完全没考虑这个,是故意留白 的复杂度。 - 把
check_aspect从“所有订单问同一句话“换成每个订单能问不同问题 的结构,改一下check_orders的参数形状,让order_ids和问题 一一对应而不是共用一个aspect。 - 用真机测一组对照实验:把
lookup_order改成顺序await而不是 靠Send并行,跑同一个三订单任务,对比总耗时——用真实数字量化 “并行“到底省了多少秒,不要只信直觉。 - 把
lookup_order里的ainvoke换成astream,看看父图能不能 观察到子图内部的中间步骤——想清楚为什么现在这个设计选择了看不见。