上一篇把 State 拆成一组通道,并用两个工具同一轮写 seen_hotel_ids 解释了 reducer 为什么不可省。但那张图仍然只有 agent 和 tools 两个节点:模型决定调什么工具,所有任务都挤在同一个循环里。
这篇把图往前推一步。用户一次要比较西安、北京和成都的酒店时,城市数量到运行时才知道;每座城市都要独立搜索、筛选,全部结束后再汇总。这里正好需要三种不同的控制手段:条件边决定这次请求能不能继续,Send 按城市动态创建并行任务,子图把「搜索 → 筛选」收成一个可复用单元。
Table of contents
Open Table of contents
一、先分清三种控制手段
条件边、Send 和子图经常一起出现,但解决的不是同一个问题。
| 机制 | 回答的问题 | 这篇里的职责 |
|---|---|---|
| 条件边 | 下一步走哪一个或哪几个已有节点 | 请求合法就继续,缺城市或日期就去补参数 |
Send | 运行时要创建多少份节点任务,每份收到什么输入 | 有几座城市,就创建几份 plan_city |
| 子图 | 一段多节点流程如何封装、复用和独立持久化 | 每座城市都跑同一套「搜索 → 筛选」 |
它们不是三种等价写法。条件边选择的是已经画好的路径;Send 创建的是同一个节点的多份运行时任务;子图则把一组节点折叠成父图里的一个节点。把三者接起来之后,父图只关心「校验、分发、汇总」,单城市内部怎么搜索和筛选留在子图里。
从运行时看,这是一段标准的 fan-out / fan-in,也就是 map-reduce:
| 阶段 | 当前 superstep 做什么 | 下一轮能看到什么 |
|---|---|---|
| 校验 | validate_request 写入错误列表 | 条件边读取完整校验结果 |
| 分发 | dispatch 根据城市数返回多份 Send | 三份 plan_city 任务被调度 |
| map | 三份子图各自搜索、筛选 | 每份任务产出一条城市计划 |
| reduce | city_plans reducer 合并三份更新 | summarize 读到完整结果集 |
| 汇总 | 统一排序并生成答案 | 进入 END |
这里的 fan-out 不是复制三份父图进程,而是在同一次 run 里生成三份节点任务;fan-in 也不是额外画出的节点,而是运行时等待本轮任务结束、按通道合并更新的同步边界。先抓住这条时间线,再看 API,就不会把「画了三条线」和「运行时创建三份任务」混为一谈。
二、条件边只负责选路
先从最窄的分支开始。validate_request 检查城市和入住日期,把错误写进 State;路由函数只读取结果并返回下一节点的名字:
from typing import Literal
def route_request(state: TripState) -> Literal["dispatch", "clarify"]:
if state.get("errors"):
return "clarify"
return "dispatch"
builder.add_conditional_edges(
"validate_request",
route_request,
{
"dispatch": "dispatch",
"clarify": "clarify",
},
)
第三个参数是路径映射:路由函数返回逻辑标签,图再把标签映射成真实节点名。标签与节点名相同时可以省略映射,但保留这一层有两个好处:路由函数不必知道图上节点以后叫什么,生成的图结构也更容易读。
返回类型里的 Literal["dispatch", "clarify"] 不只服务于类型检查,也让 LangGraph 在可视化时知道这条条件边有哪些可能出口。如果既没有 Literal,也没有显式路径映射,运行时仍可能跑通,但生成的图会缺少准确的分支信息。把候选路径声明出来,相当于同时补齐代码契约和架构图。
路由函数最好保持纯函数。不要一边判断,一边修改 state["errors"],因为条件边返回值只用于选路,状态更新应该由前面的节点显式返回。把写状态和选路径混在一起,不但很难从轨迹判断是谁改了 State,重放时也容易出现无法解释的分支差异。
还要区分「条件边」和普通边。add_edge("dispatch", "plan_city") 表示每次都去同一个节点,任务数在画图时已经固定;add_conditional_edges 可以根据当前 State 返回一个节点名、END,或者下一节要用的一组 Send。路径在运行时才确定,但候选节点仍要在编译前注册。
三、Send 不是普通分叉,而是动态扇出
如果永远只查西安和北京,可以画两条静态边;问题是用户这次传一座城市,下次可能传十座。图在 compile() 时还不知道要并行多少份,不能提前创建 plan_xian、plan_beijing 这类节点。
Send 把「去哪个节点」和「这份任务拿什么输入」装在一起。分发函数遍历运行时的 cities,每座城市返回一份任务:
from langgraph.types import Send
def fan_out_cities(state: TripState):
return [
Send(
"plan_city",
{
"city": city,
"check_in": state["check_in"],
"budget": state["budget"],
},
)
for city in state["cities"]
]
builder.add_conditional_edges(
"dispatch",
fan_out_cities,
["plan_city"],
)
这里的第二个参数不是父图完整 State,而是这一份 worker 的输入。西安任务只拿到西安,北京任务只拿到北京,彼此不需要共享候选酒店。Send 因而很适合 map-reduce:先把列表映射成若干并行任务,再把各任务的输出归并回父图。
注意到代码里的 plan_city 是下一步的节点名,而传入的字典完全匹配了下一节的 CityInput 字段。这意味着通过 Send 派发给子图的数据,不需要在子图里再做一层适配,LangGraph 会直接用这份负载去启动子图节点。这是在多节点之间隔离状态最干净的写法。
条件边直接返回 ["check_weather", "search_hotels"] 时,两个已有节点都会收到同一份父图 State,适合固定的并行分支;返回 [Send("plan_city", {...}), ...] 时,同一个节点会收到多份不同输入,任务数量也由本次数据决定。判断该用哪一种,只需问两个问题:分支数量是不是编译时已知,每个分支拿到的输入是否相同。
扇出之后最容易漏掉的仍是上一篇的 reducer。三份 plan_city 会在同一个 superstep 写 city_plans;如果这条通道没有 reducer,框架无法决定保留哪一份,仍会抛出 InvalidUpdateError:
import operator
from typing import Annotated
from typing_extensions import TypedDict
class TripState(TypedDict):
cities: list[str]
check_in: str
budget: int
errors: Annotated[list[str], operator.add]
city_plans: Annotated[list[dict], operator.add]
answer: str
operator.add 在这里可以直接使用,因为每个 worker 只返回一条新计划,不会回传父图已有的全量列表。若应用允许重复分发同一城市,或者上游会重试已经成功的任务,还要给计划加稳定 id,并改成按 id 去重的自定义 reducer;否则同一座城市可能被追加两遍。
四、子图把单城市流程收起来
plan_city 如果只是一个函数,搜索、过滤、排序都会塞进同一个节点,图上看不出内部在哪一步失败。更合适的做法是把单城市流程编译成子图:
class CityInput(TypedDict):
city: str
check_in: str
budget: int
class CityOutput(TypedDict):
city_plans: Annotated[list[dict], operator.add]
class CityState(CityInput, CityOutput):
candidates: list[dict]
city_builder = StateGraph(
CityState,
input_schema=CityInput,
output_schema=CityOutput,
)
city_builder.add_node("search_hotels", search_hotels)
city_builder.add_node("rank_hotels", rank_hotels)
city_builder.add_edge(START, "search_hotels")
city_builder.add_edge("search_hotels", "rank_hotels")
city_builder.add_edge("rank_hotels", END)
city_graph = city_builder.compile()
builder.add_node("plan_city", city_graph)
CityInput 是 Send 能传入的字段,CityOutput 只把 city_plans 交还父图;candidates 是子图私有通道,不会泄漏到父图。这和上一篇用 input_schema 阻止调用方注入 seen_hotel_ids 是同一个机制,只是这次输入、内部状态和输出三层都分开了。
子图接入父图有两种方式。父子图共享通道时,像上面一样把编译后的图直接传给 add_node,框架负责传递共享字段;两边 schema 完全不同,或者需要重命名、裁剪字段时,则在普通节点里调用子图并手动映射:
def call_city_graph(state: ParentState):
result = city_graph.invoke({
"city": state["target_city"],
"check_in": state["date"],
"budget": state["max_price"],
})
return {"city_plans": result["city_plans"]}
不要为了复用强行让父子图共用一份巨大 State。直接挂载适合字段语义一致的流程;包装调用适合边界不同的模块。字段同名却含义不同,比多写三行映射更危险。
五、并发结束后为什么只汇总一次
父图剩下的连线如下:
builder.add_node("validate_request", validate_request)
builder.add_node("dispatch", lambda state: {})
builder.add_node("plan_city", city_graph)
builder.add_node("clarify", clarify)
builder.add_node("summarize", summarize)
builder.add_edge(START, "validate_request")
builder.add_conditional_edges("validate_request", route_request)
builder.add_conditional_edges("dispatch", fan_out_cities, ["plan_city"])
builder.add_edge("plan_city", "summarize")
builder.add_edge("clarify", END)
builder.add_edge("summarize", END)
三份 plan_city 属于同一次动态扇出。它们可以并行执行,但 summarize 不会在第一份城市结果回来时就抢跑;LangGraph 会先完成这一轮所有待执行任务,把三份 city_plans 经 reducer 合并,再进入下一个 superstep。于是 summarize 看到的是完整列表,只执行一次。
这也是 map-reduce 架构下的天然风险:只要其中一份任务抛出未捕获的异常,这整个 superstep 就会直接崩溃,连带父图进程一起退出。这意味着另外两座已经查好的城市结果也会跟着报废。如果业务期望「成都查失败了,西安和北京仍能出计划」,必须在 plan_city 内部捕获异常,并返回一条特殊的包含错误原因的 city_plan 记录。
这里的「并行」描述的是图层面的可并发任务,不等于任意节点内部的阻塞调用都会自动变成高吞吐异步代码。如果 search_hotels 使用同步 HTTP 客户端,线程、连接池和 provider 限流仍会决定实际吞吐;节点内的 async、超时和取消留到后面的并发篇再拆。
六、子图落在哪一份 checkpoint 里
子图不是另起一套完全无关的运行时。父图带 checkpointer 时,默认编译的子图会继承父图的持久化能力,并为每次调用分配独立的 checkpoint namespace;因此三座城市并发执行时,各自的内部 candidates 和执行位置不会互相覆盖。
LangGraph 会自动按 thread_id:节点名:uuid 的格式拼接出这种相互隔离的内部命名空间。由于用了 Send,三份并发任务会被当作三个不同实例,底层不仅状态互不干扰,在 UI 里查看时也会分别挂在三条不同的轨迹下。
compile(checkpointer=...) 决定子图内部状态保留多久:
| 配置 | 语义 | 适合场景 |
|---|---|---|
None,默认 | 每次调用隔离,但继承父图 checkpointer,支持本次调用内中断与恢复 | 独立的单城市规划 |
True | 同一 thread 跨调用积累子图状态 | 真正需要多轮记忆的子 Agent |
False | 不保存子图 checkpoint | 无中断、失败后允许整段重跑的纯计算 |
这篇的三份城市任务应该保持默认值。checkpointer=True 的 per-thread 子图不适合同一实例被并行调用:多份任务会竞争同一 namespace。不要把「状态保留得更久」当成更安全,它改变的是状态所有权,也会改变并发边界。
七、recursion_limit 限制的是 superstep
图一旦出现环,就必须有业务退出条件。最小循环依靠「模型不再返回 tool_calls」进入 END;这篇的单城市子图没有环,但真实筛选流程可能在候选为空时放宽条件,再回到搜索节点。如果分支一直判断为「继续重试」,图就不会结束。
recursion_limit 是运行时最后一道保险:
from langgraph.errors import GraphRecursionError
try:
result = graph.invoke(
{
"cities": ["西安", "北京", "成都"],
"check_in": "2026-10-01",
"budget": 400,
"errors": [],
"city_plans": [],
},
{"recursion_limit": 40},
)
except GraphRecursionError:
# 记录 thread_id、当前节点和 state,再决定人工接管或重新规划
...
这个值放在 config 顶层,不是 configurable 里面。它限制的是一次 run 最多执行多少个 superstep,不是 Python 函数递归深度,也不是「最多调用 40 次工具」。一次 Send 扇出十份任务仍可能只占一个并行 superstep;反过来,两个节点来回跳 20 轮会快速耗尽额度。
想知道当前跑到了第几步,或者用来限制模型最后一次机会,可以在节点的参数中声明 RemainingSteps。这是一个特殊的注入通道,图越转它的值就越小。
不要把调大上限当成修复。预期十步结束却撞上限制,先查条件边是否永远返回同一路径、重试计数是否写回 State、结束分支是否真的连到 END。只有业务本来就需要长图时才提高上限,并为重试轮数、累计 token 和外部调用次数分别设业务护栏。
八、对照前面的循环
agent ↔ tools 最小循环 | 本篇父图 + 子图 |
|---|---|
模型通过 tool_calls 决定下一步 | 代码通过条件边显式决定路径 |
| 工具数量固定,调用次数由模型决定 | worker 节点固定,实例数由 Send 在运行时决定 |
所有工具结果回到 messages | 每份子图只返回 city_plans,私有状态不外泄 |
add_messages 合并轨迹 | operator.add 或自定义 reducer 合并并行结果 |
max_turns 防止模型循环 | recursion_limit 限制整张图的 superstep |
两种写法不是互斥关系。子图内部仍然可以是一个 Agent 循环,父图则用确定性分支控制业务流程。该让模型判断「用户更偏好地段还是价格」时交给模型;该判断「日期缺失不能查库存」时直接写条件边。把确定性规则也交给模型,只会让轨迹更贵、更难复现。
九、跑起来看轨迹
完整代码在 agent-in-action/part2-agent-frameworks/langgraph/06-branching-send-subgraph/。与上一层 05-state-reducer 相比,这一层把单城市搜索拆成子图,再在父图里按城市动态扇出。
cd part2-agent-frameworks/langgraph/06-branching-send-subgraph
python agent.py
输入西安、北京、成都三个城市时,轨迹的形状应该是:
validate_request cities=3, errors=[]
dispatch Send(plan_city, 西安) × 1
Send(plan_city, 北京) × 1
Send(plan_city, 成都) × 1
plan_city/西安 search_hotels → rank_hotels
plan_city/北京 search_hotels → rank_hotels
plan_city/成都 search_hotels → rank_hotels
summarize city_plans=3
日志先后次序不保证按城市输入顺序,因为三份任务是并发的;最终结果若需要稳定排序,应在 summarize 里按价格、评分或原始城市顺序显式排序,不能把 reducer 的合并顺序当成业务顺序。
十、常见问题
| 故障现象 | 排查思路 |
|---|---|
| 城市数量变了,图上仍只执行一份任务 | 用的是普通边或返回节点名,而不是按城市返回 Send 列表 |
多城市执行时报 InvalidUpdateError | 多份 worker 同一轮写 city_plans,但该通道没有 reducer |
city_plans 里同一城市重复出现 | 上游重复分发了已成功任务,reducer 又只做追加;给结果稳定 id,并按 id 去重 |
子图读不到 candidates | 把私有字段误放进了 output_schema,或搜索节点没有返回这条内部通道 |
| 父图出现子图的全部临时字段 | 父子图共用了一份过宽 State;拆分 CityInput、CityState 和 CityOutput |
summarize 结果顺序偶尔变化 | 并行完成顺序不稳定;汇总节点必须显式排序 |
设置了 recursion_limit 仍没生效 | 错放进了 configurable;它应位于 invoke config 顶层 |
提高 recursion_limit 后只是更晚报错 | 图存在无法收敛的环;先修退出条件,不要只抬上限 |
到这里,图已经从一个固定循环变成了可分流、可扇出、可嵌套的工作流,但 checkpoint 仍只被当成「停住再续」使用。下一篇沿着这份持久化状态继续往下:读取历史 checkpoint、改写某一刻的 State,再从旧分支 fork 出一条新执行路径。