从最小循环到 create_agent,前面四篇一直在用 State——messages、seen_hotel_ids、orders 三个字段,加一个 add_messages。但 State 到底是什么,add_messages 做了什么,从来没有展开过。架构概述里那句「state 怎么长,后面各有一篇」,这篇来还。先回头讲透,再往前推一步:当同一轮里两个工具同时写一个字段时,reducer 是唯一能让图不崩的机制。
需要先说明一点:代码层面这篇建在 04-create-agent 之上,而不是退回 01-min-loop。原因是到第四篇时 seen_hotel_ids 和 orders 都已经在 State 里了,正好是这篇要解剖的现成素材——三篇下来一直在用,却一直没讲。顺着系列读的人会从代码目录 01 跳到 05,中间三层的 diff 叙事断开,但概念上更顺:这篇回填的是前四篇共同的地基。
Table of contents
Open Table of contents
一、State 不是一份字典,是一组通道
最容易形成的误解,是把 State 当成一个被节点整体替换的字典。实际上,当你写下 StateGraph(State) 时,State 的每一个字段都被框架拆成一条独立的通道。节点返回的不是「新的 State」,而是一份针对若干通道的写请求;没被提及的通道原样保留。
这解释了 Command(update=...) 为什么只写两个字段,其余字段却不受影响。在 SqliteSaver 那篇里,search_hotels 返回的是:
return Command(update={
"seen_hotel_ids": seen, # 只写白名单这条通道
"messages": [_tool_message(...)], # 只写消息这条通道
})
orders 通道没有被提及,因此保持原值。通道之间各算各的账,互不覆盖。
正因如此,SqliteSaver 那篇第 74 行那句一带而过的话才有了着落:「没有指定 reducer,每次 update 都会直接替换上新的完整数据」。没有 reducer 的通道,语义是「最后写入者胜出」——谁在这一步写了它,它就变成谁给的值。但更关键的一点当时没讲:在没人写过之前,这条通道根本不存在。这就是为什么第一次 invoke 必须带上空列表和空字典:
graph.invoke({
"messages": [...],
"seen_hotel_ids": [], # 不写这行,后续 state.get("seen_hotel_ids") 拿到的是 None
"orders": {},
}, config)
不给初值,通道从未被创建,state.get("seen_hotel_ids") 返回的不是空列表,而是 None。book_hotel 拿着 None 去查白名单,自然查无所获。
把一个 superstep 拆开看,State 的更新分三段完成。节点开始执行时读到的是同一份稳定快照;执行期间各节点只提交局部写请求,不能直接看到同轮其他节点刚写的值;等这一批节点都结束后,运行时再按通道收集更新、调用 reducer,生成下一轮可见的新快照。这也是为什么并行节点不能靠「谁先执行完」互相传值——有先后依赖就要拆到两个 superstep,用边连接。
| 本轮写入情况 | 无 reducer 的通道 | 有 reducer 的通道 |
|---|---|---|
| 没有节点写 | 保留旧值 | 保留旧值 |
| 只有一份更新 | 用新值替换旧值 | 执行 reducer(existing, update) |
| 同时收到多份更新 | 抛出 InvalidUpdateError | 逐份归并成下一版值 |
Annotated[list[str], dedup_ids] 也可以按这张表来读:list[str] 是通道里保存的数据类型,dedup_ids 是这条通道收到更新时使用的合并函数。TypedDict 负责描述 State 的形状,Annotated 附带运行时元数据;它不会改变 Python 值本身的类型。
二、add_messages 到底做了什么
最小循环的对照表里写过一行光秃秃的 messages.append → add_messages,但没展开。add_messages 做的不只是追加。
先说为什么 messages 非要有 reducer不可。在 agent ↔ tools 的循环里,agent 节点返回 {"messages": [ai_message]},tools 节点返回 {"messages": [tool_message]}。这两个写发生在不同的 superstep 里。如果 messages 没有 reducer,语义就是「最后写入者胜出」——tools 一写,agent 那条消息就被整段抹掉,历史随之断层。add_messages 的第一层作用,是把每一步的写请求追加到既有列表末尾,而不是替换。跨步累积,历史才留得住。
第二层作用同样重要:它按消息 id 去重。新消息的 id 与既有消息重合时,add_messages 不会追加,而是覆盖那条旧消息。这正是后续在中间件里改写既有消息的基础——不改 id,框架就认得它,更新而非新增。每个节点只知道自己产出了什么,不知道全量历史;reducer 替它把新旧拼成一份完整的 messages。
三、同一轮两个工具写同一个字段
前面四篇的实跑轨迹里,并发一直都在。最小循环的轨迹里,模型一轮同时调了 search_location 和 search_hotels,ToolNode 在同一个 superstep 里执行了两个 tool_call。之所以相安无事,是因为这两个工具写的通道不冲突——天气查询只回灌 messages,酒店搜索才写 seen_hotel_ids。
但只要同一轮里有两个工具写同一条没有 reducer 的通道,图就会崩。把场景往前推一步:用户问两个城市的酒店,模型同一轮调了两次 search_hotels,两次都返回 Command(update={"seen_hotel_ids": ...})。跑起来,LangGraph 抛出:
InvalidUpdateError: At key 'seen_hotel_ids': Can receive only one value per step.
Use an Annotated key to handle multiple values.
报错的含义很直白:一个没有 reducer 的通道,在一个 superstep 里只能收到一份写请求。框架不知道两份值该取谁、该怎么合并,于是拒绝猜测,直接报错。这不是 bug,是一道安全闸——宁可崩,也不静默丢数据。
最直接的修法是给 seen_hotel_ids 挂上 operator.add:
import operator
class State(TypedDict):
messages: Annotated[list, add_messages]
seen_hotel_ids: Annotated[list[str], operator.add] # 原来没有 reducer
orders: dict[str, dict]
operator.add 对列表就是拼接:existing + update。两个 search_hotels 各自返回本次新搜到的 id,reducer 把两份拼成一条。但这里有个前提要改——每个工具不能再返回全量列表,只能返回本次新增的部分,否则越拼越长。
四、自定义 reducer:去重与合并
operator.add 能让图不崩,但只会机械拼接。两次搜索本来就可能命中同一家酒店,上游重试或重复事件也可能再次提交同一批 id;如果 reducer 只做 existing + update,白名单会不断积累重复值。因此 seen_hotel_ids 需要一个去重的 reducer:
def dedup_ids(existing: list[str], update: list[str]) -> list[str]:
if existing is None:
existing = []
seen = set(existing)
return existing + [x for x in update if x not in seen]
orders 的情况更直接:它是 dict,而 operator.add 调用的是 +,两个 dict 相加会直接抛 TypeError。这一步没有捷径,必须自己写合并函数:
def merge_orders(existing: dict, update: dict) -> dict:
if existing is None:
existing = {}
return {**existing, **update} # 新订单覆盖同键,其余保留
写 reducer 时要检查四件事。第一,签名是 (existing, update),返回合并后的新值,不要原地修改 existing。第二,函数必须纯净,不能在里面写数据库、发请求或修改模块变量;checkpoint 恢复、状态回放和调试都可能再次计算它。第三,并行更新不应依赖完成顺序,合并函数至少要满足结合律,若业务不关心顺序,最好也满足交换律。第四,是否需要幂等要由数据语义决定:消息列表允许重复内容,但酒店白名单和订单通常要按稳定 id 去重。
当前的列表、字典通道通常会从对应类型的空值开始,但 reducer 若声明了可空类型、接收历史数据,或者会被独立单测,仍应明确处理 None。这里保留防御分支,是为了让函数契约完整,而不是假设每次首次调用都一定传 None。
需要区分两件容易混淆的事:interrupt() 发生时,本 superstep 尚未提交的 State 更新会随 checkpoint 回滚,恢复后重新计算,并不等于 reducer 必然把同一份已提交值再追加一次;真正危险的是节点里已经发生的外部副作用,以及应用层确实重复送达的更新。前者靠幂等键,后者才靠去重 reducer。这条边界和 中断续跑里「写操作必须放在 interrupt() 之后」说的是同一类重放风险,但保护层不同。
改完之后,State 长成这样:
class State(TypedDict):
messages: Annotated[list, add_messages]
seen_hotel_ids: Annotated[list[str], dedup_ids] # 去重,重复更新安全
orders: Annotated[dict, merge_orders] # dict 合并,+ 不可用
五、对外只收 messages:input schema 分离
到目前为止,调用方每次 invoke 都得传齐三个字段,包括 seen_hotel_ids 和 orders。但这两个是内部状态——白名单由 search_hotels 写入,订单由 book_hotel 写入。如果调用方从入口就能塞进 seen_hotel_ids,等于绕过了 最小循环里那道「hotel_id 必须来自搜索结果」的校验。
LangGraph 允许把输入 schema 和内部 schema分开。定义一个只含 messages 的 InputState,传给 StateGraph:
class InputState(TypedDict):
messages: Annotated[list, add_messages]
class State(InputState):
seen_hotel_ids: Annotated[list[str], dedup_ids]
orders: Annotated[dict, merge_orders]
# 对外只收 messages;seen_hotel_ids / orders 是内部通道,调用方碰不到
builder = StateGraph(State, input_schema=InputState)
input_schema 控制的是 invoke 能收什么;内部通道仍然在 State 里供节点读写,但调用方无从注入。output_schema 同理,控制 invoke 返回什么——不写就默认和 State 一致。节点也可以有自己的 input_schema,只读自己需要的字段,这是后面子图那篇会用到的话题。
六、对照手搓看差异
手搓 01 / 02 | 本篇 |
|---|---|
模块级 seen_hotel_ids / ORDERS | State 通道,随图流转、随 Checkpointer 落盘 |
messages.append 手动追加 | add_messages reducer:追加 + 按 id 去重覆盖 |
| 同一轮并发写同一变量,自己加锁 | reducer 替你合并,无锁 |
| 续跑时重放,外部写靠幂等键 | reducer 保持纯净,重复更新按业务键去重 |
| 调用方理论上能改任何全局变量 | input_schema 把内部通道挡在入口之外 |
循环形状没变。变的是状态的归属:手搓里状态散落在模块变量和函数局部变量里,框架把它们收编成通道,每条通道有自己的合并规则。reducer 不是可选的装饰,是通道在并发和重放下保持一致性的唯一手段。
七、跑起来看轨迹
完整代码在 agent-in-action/part2-agent-frameworks/langgraph/05-state-reducer/。与上一层 04-create-agent 做 Compare Files,改动集中在三处:State 三个字段都挂上 reducer、search_hotels 改为只返回新增 id、StateGraph 加上 input_schema。
cd part2-agent-frameworks/langgraph/05-state-reducer
python agent.py
建议分两次跑。第一次先把 seen_hotel_ids 的 reducer 去掉,复现 InvalidUpdateError;第二次挂回 dedup_ids,看它正常合并。用 graph.get_state(config).values 打印各通道的值,能看到 seen_hotel_ids 去重后的完整列表和 orders 合并后的订单字典——这就是 reducer 跑完之后各通道的最终状态。
—— 同一轮两次 search_hotels ——
AIMessage 要调 ['search_hotels', 'search_hotels']
ToolMessage {"city": "西安", "hotels": [{"id": "HT-002", ...}]}
ToolMessage {"city": "北京", "hotels": [{"id": "HT-004", ...}]}
—— get_state ——
seen_hotel_ids: ['HT-002', 'HT-004'] # 去重后拼接,无重复
orders: {} # 尚未下单
两次搜索各写一份 id,dedup_ids 合并后白名单里只剩两条,没有重复。如果这一 superstep 在提交前触发 interrupt(),本轮通道更新不会先落进 checkpoint;续跑虽然会重新执行节点,成功后仍只提交一轮结果。若应用层把已经成功的搜索更新再次送达,dedup_ids 才会发挥幂等合并作用,把相同 id 挡掉。
八、常见问题
| 故障现象 | 排查思路 |
|---|---|
InvalidUpdateError: Can receive only one value per step | 同一 superstep 里多个节点或多个 tool_call 写了同一条无 reducer 的通道;给该字段挂 reducer |
挂了 operator.add 之后列表越来越长 | 工具返回的是全量列表而非新增部分;改成只返回本次新值 |
| 重试或重复事件后白名单出现重复 id | reducer 用的是 operator.add 而非按业务 id 去重;换成 dedup_ids 这类纯函数 reducer |
orders 字段报 TypeError: unsupported operand type(s) for + | operator.add 对 dict 不可用;dict 字段必须写自定义合并 reducer |
state.get("seen_hotel_ids") 返回 None | 第一次 invoke 没给该通道初值;无 reducer 的通道在被写之前不存在 |
调用方从入口塞进了 seen_hotel_ids | 没设 input_schema,内部通道暴露给了外部;加 input_schema=InputState 收窄入口 |
通道模型讲清之后,下一步就是把图从「agent ↔ tools 来回」推到非线性流转:按 state 走不同的边、运行时扇出、子图嵌套。下一篇讲条件边、Send 和子图——这也是 架构概述里点名 Send 和子图之后,一直欠着的那篇。