1. 本章目标 #
上一章结尾留了一个悬念:同一个字段被多个节点写入时,默认是「后写覆盖先写」。 这个规则对「清洗文本」是对的,但对「对话消息」就完全错了——你要的是把新消息追加到历史后面,而不是把之前的对话全冲掉。
上一章还留了一个报错:两个并行分支同时写一个字段时,图会直接抛异常。
InvalidUpdateError: At key 'text': Can receive only one value per step.
Use an Annotated key to handle multiple values.这两个问题的答案是同一个东西:Reducer。
本章目标:
搞清楚 Reducer 如何决定「新值和旧值怎么合并」,并吃透
add_messages这个最重要的 reducer。
学完这一章,你应该能做到:
- 用
Annotated[类型, reducer函数]给状态字段指定合并规则 - 说清 reducer 什么时候被调用、参数是什么,包括一个反直觉的事实:
invoke传进去的初始值也会走 reducer - 知道字段的类型标注会决定
existing的初始值,这是TypedDict标注在 LangGraph 里唯一有运行时意义的地方 - 写出自己的 reducer:去重累加、取最大值、字典合并
- 用 reducer 让并行分支正常合流,并且清楚并行时节点读到的是什么、多个分支的值又是怎么被依次喂进 reducer 的
- 掌握
add_messages的五种行为:追加、按 id 替换、自动补 id、格式转换、用RemoveMessage删除 - 用
MessagesState快速搭出对话类的图,并按需扩展自己的槽位 - 产出:可追加 messages 的状态图,一张同时用了三种 reducer 的多轮对话图
参考文档:
2. 先复现问题 #
2.1. 默认行为:覆盖 #
给状态加一个 log 字段,让两个节点各往里写一条记录:
# 状态用 TypedDict,第 20 章的默认写法
from typing import TypedDict
# 图的三件套
from langgraph.graph import END, START, StateGraph
# 状态里只有一个列表字段,故意先不加任何 reducer
class S1(TypedDict):
# 期望它能累积多条记录
log: list[str]
# 第一个节点:往 log 里写一条
def a(state: S1) -> dict:
# 注意返回的是一个只含一个元素的列表
return {"log": ["A 来过"]}
# 第二个节点:也往 log 里写一条
def b(state: S1) -> dict:
# 同样是单元素列表
return {"log": ["B 来过"]}
# 组装成 a → b 的串行图
builder = StateGraph(S1)
# 注册 a 节点
builder.add_node("a", a)
# 注册 b 节点
builder.add_node("b", b)
# 入口边
builder.add_edge(START, "a")
# a 跑完接着 b
builder.add_edge("a", "b")
# b 跑完结束
builder.add_edge("b", END)
# 编译后的对象统一叫 graph
graph = builder.compile()
# 传一个初始值进去,看最后剩下什么
print("结果:", graph.invoke({"log": ["初始"]}))运行输出:
结果: {'log': ['B 来过']}只剩最后一条。 不光 A 来过 没了,连传进去的 初始 也被冲掉了。
这里其实发生了两次覆盖,很多人只数到一次:
['初始'] ──a 写入──► ['A 来过'] ──b 写入──► ['B 来过']
↑ ↑ ↑
初始值 第一次覆盖 第二次覆盖换句话说,你从 invoke 传进去的初始值也会被覆盖掉,它并没有什么特殊地位。这个细节在 §3.3 会展开成一条完整的机制。
这不是 bug,而是默认规则:没有指定 reducer 时,新值直接替换旧值。 对 status = "已完成" 这种字段,覆盖正是我们要的;对 log 这种需要累积的字段,就完全不对了。
而且它是悄悄发生的。第 20 章 §10 那份静默失败清单里,这一条排在最容易踩的位置:代码不报错,只是数据无声无息地少了。
2.2. 你需要的是「怎么合并」 #
问题的本质是:图不知道你想怎么合并。
旧值 ['初始'] + 新值 ['A 来过'] = ?
覆盖 → ['A 来过'] ← 默认
追加 → ['初始', 'A 来过'] ← 你想要的
去重追加、取最大、字典合并… ← 还有别的可能「怎么合并」只有你知道,所以得告诉图。告诉它的方式,就是加一个 reducer。
3. Reducer 基础 #
3.1. 一句话和语法 #
Reducer 是一个函数,它规定「这个字段的旧值和新值怎么合并成新的旧值」。
「新的旧值」听着有点绕,但说得很准:reducer 的输出会成为下一次调用的输入,像滚雪球一样一层层累积上去。如果你写过 JavaScript 的 Array.reduce 或者 Python 的 functools.reduce,那就是同一个概念——LangGraph 的 reducer,就是把图运行期间对某个字段的所有写入「归约」成一个最终值。
语法是用 Annotated 给字段挂一个函数,写出来是这个样子:
字段名: Annotated[字段类型, reducer 函数]
↑ ↑ ↑
标准库 typing 真正的类型 合并规则Annotated 来自 Python 标准库 typing,作用是「给类型标注附加额外信息」。类型检查器只看第一个参数(list[str]),LangGraph 则专门读第二个参数,把它当 reducer 用。
它只是标注,不改变运行时的值。 Annotated[list[str], operator.add] 对 Python 来说仍然是 list[str],节点里拿到的还是普通列表,不是什么包装对象。所以加不加 reducer,节点代码完全一样,下一节就能看到。
3.2. 第一个 reducer:operator.add #
对列表来说,「追加」就是列表相加,Python 标准库里现成就有 operator.add:
# operator 模块把 + - * / 这些运算符包装成了普通函数
import operator
# Annotated 用来挂 reducer
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 和 S1 的唯一区别就在下面那一行
class S2(TypedDict):
# 唯一的改动:给 log 挂上 operator.add
log: Annotated[list[str], operator.add]
# 节点代码和 §2.1 完全一致,一个字都没改
def a2(state: S2) -> dict:
# 依然只返回自己新增的那一条
return {"log": ["A 来过"]}
# 第二个节点同样没变
def b2(state: S2) -> dict:
# 依然只返回自己新增的那一条
return {"log": ["B 来过"]}
# 建图的代码也和 §2.1 一模一样
builder = StateGraph(S2)
# 注册第一个节点
builder.add_node("a", a2)
# 注册第二个节点
builder.add_node("b", b2)
# 入口边
builder.add_edge(START, "a")
# a 跑完接着 b
builder.add_edge("a", "b")
# b 跑完结束
builder.add_edge("b", END)
graph = builder.compile()
# 同样传一个初始值
print("结果:", graph.invoke({"log": ["初始"]}))运行输出:
结果: {'log': ['初始', 'A 来过', 'B 来过']}节点代码一个字没改,只是在状态定义里加了一个 Annotated,行为就从覆盖变成了追加。
这是 reducer 设计上最干净的一点:合并规则属于「状态」,不属于「节点」。 节点只管返回自己算出来的那一小段,怎么合进去是状态的事。
3.3. Reducer 到底什么时候被调用 #
这一节的实验能解开大部分疑惑,做法是把 reducer 变成一个探针。reducer 是唯一能看到「合并这个动作本身」的地方,在里面打一行 print,整个合并过程就透明了。
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 自定义 reducer:功能上等价于 operator.add,只是多打一行日志
def spy_reducer(existing, update):
"""只做追加,但把每次调用的参数打出来。"""
# !r 用 repr 打印,能看清是 [] 还是 None、是 '' 还是空
print(f" reducer 被调用: existing={existing!r} update={update!r}")
# existing or [] 是兜底写法,防止 existing 是 None 时报错
return (existing or []) + update
# 把这个探针 reducer 挂到 log 上
class S3(TypedDict):
# 类型还是 list[str],只是 reducer 换成了带日志的版本
log: Annotated[list[str], spy_reducer]
# 还是那张 a → b 的串行图
builder = StateGraph(S3)
# 用 lambda 简写节点,逻辑不重要,重点看 reducer 被调了几次
builder.add_node("a", lambda s: {"log": ["A"]})
# 第二个节点同理
builder.add_node("b", lambda s: {"log": ["B"]})
# 入口边
builder.add_edge(START, "a")
# 串行连接
builder.add_edge("a", "b")
# 出口边
builder.add_edge("b", END)
graph = builder.compile()
# 第一次实验:传初始值
print("--- 传初始值 ---")
# 期待看到 3 次 reducer 调用
print("结果:", graph.invoke({"log": ["初始"]}))
# 第二次实验:什么都不传,对比调用次数
print("--- 不传初始值 ---")
# 期待看到 2 次 reducer 调用
print("结果:", graph.invoke({}))运行输出:
--- 传初始值 ---
reducer 被调用: existing=[] update=['初始']
reducer 被调用: existing=['初始'] update=['A']
reducer 被调用: existing=['初始', 'A'] update=['B']
结果: {'log': ['初始', 'A', 'B']}
--- 不传初始值 ---
reducer 被调用: existing=[] update=['A']
reducer 被调用: existing=['A'] update=['B']
结果: {'log': ['A', 'B']}先注意调用次数:传初始值时调了 3 次(1 个初始值 + 2 个节点),不传时只调了 2 次。这个数字差,就是下面第二条结论的证据。
三个重要结论:
第一,reducer 的签名是 (existing, update)。 第一个参数是当前的累积值,第二个是节点这次返回的值,返回值则成为新的累积值。参数名叫什么无所谓(叫 old, new 也行),是位置决定含义。
第二,invoke 传进去的初始值也会走 reducer。 看第一行 existing=[] update=['初始']。图并不是「先把初始值放进状态,再开始跑」,而是把初始输入也当成一次更新送进 reducer。
这一点很反直觉,但它能解释本章后半的好几个现象,值得单独记住:
| 你以为的 | 实际的 |
|---|---|
| 初始值直接成为状态,然后节点开始改 | 初始值是「第一次更新」,和节点的更新走同一条路 |
两个直接后果,后面都会碰到:
- 你用
add_messages(§6)时,传进去的{"messages": [...]}也会被它处理,所以字符串和字典会被自动转成HumanMessage(§6.5、§7.2) - 有 reducer 的字段,你没法通过
invoke直接「设定」它的值,只能「贡献」一个值——传{"turn": 10}给一个挂了operator.add的字段,得到的是0 + 10,而如果状态里已经有 5,得到的是15而不是10(§7.3 有实测)
第三,字段没值时 existing 不是 None,而是这个类型的「零值」。 上面看到的是 [],因为字段标注成了 list[str]。换个类型就换个零值——实测六种常见类型:
| 字段类型标注 | 第一次调用时的 existing |
|---|---|
list |
[] |
dict |
{} |
set |
set() |
int |
0 |
float |
0.0 |
str |
'' |
所以用 operator.add 且完全不传初始值也不会报错:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 挂 operator.add 的列表字段
class S4(TypedDict):
# 标注成 list[str],零值会是 []
log: Annotated[list[str], operator.add]
# 只有一个节点的最小图
builder = StateGraph(S4)
# 节点往 log 里写一条
builder.add_node("a", lambda s: {"log": ["A"]})
# 入口边
builder.add_edge(START, "a")
# 出口边
builder.add_edge("a", END)
graph = builder.compile()
# 什么都不传,看会不会因为 log 没有初始值而报错
print("结果:", graph.invoke({}))运行输出:
结果: {'log': ['A']}因为 [] + ['A'] 是合法的。同理,Annotated[int, operator.add] 不传初始值也能跑,因为 0 + 1 合法:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 挂 operator.add 的整数字段,用来做计数器
class S4b(TypedDict):
# 标注成 int,零值会是 0
total: Annotated[int, operator.add]
# 同样是单节点最小图
builder = StateGraph(S4b)
# 节点每次只返回 1,靠 reducer 累加
builder.add_node("a", lambda s: {"total": 1})
# 入口边
builder.add_edge(START, "a")
# 出口边
builder.add_edge("a", END)
graph = builder.compile()
# 不传初始值,existing 会是 0
print("结果:", graph.invoke({}))运行输出:
结果: {'total': 1}这里藏着一个重要事实:在 LangGraph 里,TypedDict 的类型标注是有运行时意义的。
第 20 章 §4.3 说过「TypedDict 的标注只给 IDE 看,运行时完全不生效」——那句话在没有 reducer 的前提下是对的。但一旦挂了 reducer,标注就成了 LangGraph 推断零值的依据,写错会在运行时炸掉:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 标注写成 list,但节点返回的是数字
class SM(TypedDict):
# 标注是 list,reducer 是加法
f: Annotated[list, operator.add]
# 建一张单节点的图
builder = StateGraph(SM)
# 节点返回 5,而 existing 会被初始化成 []
builder.add_node("n", lambda s: {"f": 5})
# 入口边
builder.add_edge(START, "n")
# 出口边
builder.add_edge("n", END)
graph = builder.compile()
# [] + 5 会报错
try:
print("结果:", graph.invoke({}))
except Exception as e:
print("报错:", type(e).__name__, e)运行输出:
报错: TypeError can only concatenate list (not "int") to list报错信息里的 list 和 int 正好点明了冲突所在:existing 按标注取了 [],而 update 是 5。
写自定义 reducer 时,仍然建议加上
existing or []这层保护。 虽然实测existing不会是None,但这层判断成本很低,还能顺手兼容「手动调用 reducer 做单元测试」的场景——那时候你很可能直接传个None进去。
3.4. 节点读到的是合并后的值 #
reducer 是「写」的规则,那「读」呢?串行执行时,后面的节点读到的是已经合并好的值:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 依然是挂了 operator.add 的 log 字段
class S5(TypedDict):
# 后面 §5.2 的并行实验也会用同一个状态
log: Annotated[list[str], operator.add]
# 第一个节点:打印它读到的状态
def n1(state: S5) -> dict:
# 这里打印的是进入 n1 时的完整状态
print(" n1 读到:", state)
# 贡献一个 A
return {"log": ["A"]}
# 第二个节点:同样打印,重点看它能不能看到 n1 写的 "A"
def n2(state: S5) -> dict:
# 如果能看到 A,说明读到的是合并后的值
print(" n2 读到:", state)
# 贡献一个 B
return {"log": ["B"]}
# 建图
builder = StateGraph(S5)
# 注册第一个会打印的节点
builder.add_node("n1", n1)
# 注册第二个
builder.add_node("n2", n2)
# 入口指向 n1
builder.add_edge(START, "n1")
builder.add_edge("n1", "n2") # 串行:n1 跑完再跑 n2
# n2 跑完结束
builder.add_edge("n2", END)
graph = builder.compile()
# 传初始值,方便区分「初始」「A」「B」三段来源
print("结果:", graph.invoke({"log": ["初始"]}))运行输出:
n1 读到: {'log': ['初始']}
n2 读到: {'log': ['初始', 'A']}
结果: {'log': ['初始', 'A', 'B']}n2 读到的 log 里已经有 A 了。这和第 20 章 §2.3 那条规则一致:节点收到的永远是当前的完整状态,reducer 只是改变了「当前」这个值是怎么算出来的。
说得更精确一点:
reducer 不改变「节点读到完整状态」这条规则,它改变的是「完整状态」这个值是怎么被算出来的。 没有 reducer 时,「当前值」= 最后一次写入;有 reducer 时,「当前值」= 所有写入被归约的结果。
但这条结论只对串行成立,并行时会不一样——见 §5.3,也是本章最容易踩的坑之一。
4. 自定义 Reducer #
4.1. 规则 #
自己写一个 reducer 只需满足:
# 签名模板:第一个参数是已有的累积值,第二个是本次的更新
def my_reducer(existing, update):
# 返回合并后的新值
return ...两个参数,一个返回值。 没有别的要求:不用继承类、不用加装饰器、不用注册。和第 20 章的节点一样,它就是个普通函数,可以脱离图直接单测:
# 一个具体的 reducer:列表追加
def append_all(existing, update):
# 拷一份再拼,不改动传进来的对象
return (existing or []) + update
# 直接当普通函数调,不需要建图
assert append_all([1, 2], [3]) == [1, 2, 3]
# 空值情况也能单独验证,这就是 existing or [] 的用处
assert append_all(None, [1]) == [1]
print("两个断言都通过了")运行输出:
两个断言都通过了三条建议:
| 建议 | 原因 |
|---|---|
不要修改 existing,返回新对象 |
原地修改会让快照和历史记录出问题 |
处理 existing 为空的情况 |
用 existing or [] / or {} 兜底 |
| 保持纯函数,别在里面调 API | reducer 可能被频繁调用 |
第一条要特别小心,因为错误的写法看起来更「高效」:
# 错误示范:原地修改 existing
def bad_reducer(existing, update):
# extend 会直接改掉原列表,而这个列表可能被 checkpointer 的历史快照引用着
existing.extend(update)
# 返回的还是那个被改过的同一个对象
return existing
# 正确写法:构造新列表
def good_reducer(existing, update):
# 用 + 或 list() 拷一份,历史快照就不会被污染
return (existing or []) + update后果是历史快照会被事后篡改:你去查「第二步之后状态是什么」,看到的却是第五步的数据。这类 bug 往往要等加了 checkpointer(第 26 章)才暴露出来,排查特别费劲,所以一开始就别写原地修改。
第三条也有具体理由:reducer 的调用次数比你想的多。§3.3 已经看到初始值也会触发一次,§5.2 还会看到并行分支是逐个喂进来的。要是在 reducer 里发通知、写数据库,次数一定对不上。
还有一条隐含要求:reducer 里抛的异常会直接冒上来,不会被图包装或吞掉。
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 故意在 reducer 里抛异常,看图会怎么处理
def boom(existing, update):
raise RuntimeError("reducer 内部炸了")
class SB(TypedDict):
log: Annotated[list[str], boom]
builder = StateGraph(SB)
builder.add_node("a", lambda s: {"log": ["A"]})
builder.add_edge(START, "a")
builder.add_edge("a", END)
graph = builder.compile()
try:
graph.invoke({})
except Exception as e:
# 异常类型和消息都是 reducer 里原封不动抛出来的那个
print("报错:", type(e).__name__, e)运行输出:
报错: RuntimeError reducer 内部炸了这算好事,至少不是静默失败。但它也意味着一个写得不严谨的 reducer 就能让整张图挂掉,所以别在里面做可能失败的事,比如解析外部数据。
4.2. 例子:去重累加 #
实际需求:累积对话涉及的话题,但同一话题只记一次。
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 参数名用 old / new 也完全可以,位置决定含义
def merge_unique(old: list[str], new: list[str]) -> list[str]:
"""累加但去重,且保持首次出现的顺序。"""
# 用 list() 拷一份,避免原地修改 old(§4.1 第一条建议)
merged = list(old or [])
# 逐个检查本次的新元素
for item in new:
# 已经有了就跳过,实现去重
if item not in merged:
# 追加到末尾,因此保留的是「首次出现」的位置
merged.append(item)
# 返回全新的列表
return merged
# 挂到状态字段上
class ChatState(TypedDict):
# 类型是普通的 list[str],合并规则由 merge_unique 决定
topics: Annotated[list[str], merge_unique]
# 建一张两节点的图,第二个节点故意重复报一个已有话题
builder = StateGraph(ChatState)
builder.add_node("t1", lambda s: {"topics": ["售后", "物流"]})
builder.add_node("t2", lambda s: {"topics": ["售后", "设备"]})
builder.add_edge(START, "t1")
builder.add_edge("t1", "t2")
builder.add_edge("t2", END)
graph = builder.compile()
# 「售后」被报了两次,但结果里只该出现一次
print("结果:", graph.invoke({"topics": []}))运行输出:
结果: {'topics': ['售后', '物流', '设备']}用 set 也能去重,但会丢掉顺序。这种「业务上的细微要求」,正是自定义 reducer 的用武之地:
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class SetState(TypedDict):
# set 版本:合并规则一行就写完了,但顺序是乱的
topics: Annotated[set[str], lambda o, n: (o or set()) | set(n)]
builder = StateGraph(SetState)
builder.add_node("t1", lambda s: {"topics": {"售后", "物流"}})
builder.add_node("t2", lambda s: {"topics": {"售后", "设备"}})
builder.add_edge(START, "t1")
builder.add_edge("t1", "t2")
builder.add_edge("t2", END)
graph = builder.compile()
# 去重是做到了,但三个话题的先后顺序已经无从保证
print("结果:", graph.invoke({"topics": set()}))运行输出(集合的打印顺序每次可能不同):
结果: {'topics': {'售后', '设备', '物流'}}对「已识别话题」这种要给人看的字段,顺序一乱体验就差了(每次刷新顺序都不一样)。能用一行 lambda 解决的就用 lambda,有业务讲究的就写成命名函数。 命名函数还有个好处:merge_unique 这个名字本身就是一句文档,而 lambda o, n: ... 得读一遍才知道在干什么。
4.3. Reducer 不只用于列表 #
「reducer」这个词容易让人以为只跟列表有关,其实任何类型都能挂:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 一个状态里挂三种不同类型的 reducer
class S7(TypedDict):
# 数字累加
total: Annotated[int, operator.add]
# 取最大值(注意必须包一层,见 §4.4)
best: Annotated[int, lambda old, new: max(old, new)]
# 字典合并
conf: Annotated[dict, lambda old, new: {**(old or {}), **new}]
# 第一个节点:三个字段各写一次
def x1(state: S7) -> dict:
# total 加 10,best 竞争 3,conf 合入 a
return {"total": 10, "best": 3, "conf": {"a": 1}}
# 第二个节点:再各写一次,观察三种合并方式的差异
def x2(state: S7) -> dict:
# total 加 5,best 竞争 9(会赢),conf 合入 b
return {"total": 5, "best": 9, "conf": {"b": 2}}
# 建图
builder = StateGraph(S7)
# 注册第一个节点
builder.add_node("x1", x1)
# 注册第二个节点
builder.add_node("x2", x2)
# 入口边
builder.add_edge(START, "x1")
# 串行连接
builder.add_edge("x1", "x2")
# 出口边
builder.add_edge("x2", END)
graph = builder.compile()
# 三个字段都给初始值,方便对照
print("结果:", graph.invoke({"total": 0, "best": 0, "conf": {}}))运行输出:
结果: {'total': 15, 'best': 9, 'conf': {'a': 1, 'b': 2}}total:0 + 10 + 5 = 15best:max(max(0, 3), 9) = 9conf:两个字典合并
best 这个字段值得多看一眼。 它演示了一种很有用的模式:节点不必知道「当前最大值是多少」,只管报出自己算出的值,谁最大交给 reducer 判断。 多路打分、几个检索器比置信度的时候特别实用——各路节点互不知情,甚至可以并行跑,最后自然收敛出最大值。
常见 reducer 速查:
| 需求 | 写法 |
|---|---|
| 列表追加 | operator.add |
| 数字累加 | operator.add |
| 字符串拼接 | operator.add(零值是 '') |
| 列表去重追加 | 自定义(§4.2) |
| 取最大 / 最小 | lambda old, new: max(old, new) |
| 字典合并 | lambda old, new: {**(old or {}), **new} |
| 集合并集 | lambda old, new: (old or set()) 或 set(new) |
| 只保留最近 N 条 | lambda old, new: ((old or []) + new)[-N:] |
| 消息列表 | add_messages(§6) |
| 覆盖 | 不加 reducer,就是默认行为 |
倒数第三行那个「只保留最近 N 条」很实用,它把「裁剪」这件事直接写进了状态定义里,节点完全不用操心历史会不会太长。messages 有更专业的 add_messages + RemoveMessage 可用(§6.6);普通日志、检索结果列表,这一行就够了。
4.4. 坑:内置函数不能直接当 reducer #
看着很自然的写法:
from typing import Annotated, TypedDict
from langgraph.graph import StateGraph
# 想取最大值,直觉上直接把 max 挂上去
class Bad(TypedDict):
# 直接挂内置函数 max
best: Annotated[int, max]
# 注意:连一个节点都还没加
try:
# 报错就发生在这一行
StateGraph(Bad)
# 捕获所有异常,把类型和原文都看清楚
except Exception as e:
# 打印异常类型和信息原文
print("报错:", type(e).__name__, str(e))图还没建完就报错了——问题出在 StateGraph(Bad) 这一行:
报错: ValueError no signature found for builtin <built-in function max>原因是 LangGraph 要检查 reducer 的函数签名(确认它接受两个参数),而 max、min 这些用 C 实现的内置函数拿不到签名。
解决办法是包一层:
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class Good(TypedDict):
# 用 lambda 包一层,lambda 是 Python 函数,有签名
best: Annotated[int, lambda old, new: max(old, new)]
builder = StateGraph(Good)
# 两个节点各报一个分数,最大的那个会留下
builder.add_node("p1", lambda s: {"best": 3})
builder.add_node("p2", lambda s: {"best": 9})
builder.add_edge(START, "p1")
builder.add_edge("p1", "p2")
builder.add_edge("p2", END)
graph = builder.compile()
print("结果:", graph.invoke({"best": 0}))运行输出:
结果: {'best': 9}但这个签名检查只管「有没有签名」,不管「语义对不对」。 实测四个候选:
| 挂上去的东西 | StateGraph(...) |
实际能用吗 |
|---|---|---|
max / min |
报错 | — |
operator.add |
通过 | √ |
lambda old, new: max(old, new) |
通过 | √ |
sum |
通过 | × 运行时才炸 |
sum 这一行是个陷阱:它有签名(sum(iterable, start)),建图时一路绿灯,直到真的跑起来:
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class SumState(TypedDict):
# sum 能通过签名检查,所以这一行不会报错
total: Annotated[int, sum]
builder = StateGraph(SumState)
builder.add_node("a", lambda s: {"total": 5})
builder.add_edge(START, "a")
builder.add_edge("a", END)
# 编译也一路绿灯
graph = builder.compile()
print("建图和编译都通过了")
try:
# 真跑起来才暴露:LangGraph 是按 sum(existing, update) 调的
graph.invoke({})
except Exception as e:
print("运行时报错:", type(e).__name__, e)运行输出:
建图和编译都通过了
运行时报错: TypeError 'int' object is not iterable因为 LangGraph 是按 sum(existing, update) 调用的,相当于 sum(0, 5),而 sum 的第一个参数必须是可迭代对象。
结论:签名检查能挡住
max/min,但挡不住「签名对得上、语义不对」的函数。 判断一个函数能不能当 reducer,标准只有一条:把它当成f(旧值, 新值) -> 新的旧值读一遍,看语义通不通。operator.add通,sum不通。
4.5. 同一个状态里可以混用 #
不是所有字段都需要 reducer,混着写完全没问题:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 同一个状态里,一个字段追加、一个字段覆盖
class S8(TypedDict):
log: Annotated[list[str], operator.add] # 追加
status: str # 覆盖(默认)
# 建图
builder = StateGraph(S8)
# 第一个节点同时写两个字段
builder.add_node("p", lambda s: {"log": ["p"], "status": "处理中"})
# 第二个节点也同时写两个字段
builder.add_node("q", lambda s: {"log": ["q"], "status": "已完成"})
# 入口边
builder.add_edge(START, "p")
# 串行连接
builder.add_edge("p", "q")
# 出口边
builder.add_edge("q", END)
graph = builder.compile()
# 两个字段都给初始值,对照最终结果
print("结果:", graph.invoke({"log": [], "status": "新建"}))运行输出:
结果: {'log': ['p', 'q'], 'status': '已完成'}log 累积了两条,status 只保留最后一次写入(初始的「新建」和中间的「处理中」都被冲掉了)。这才是常态:一张图里通常只有一两个字段需要累积。
判断标准就一句话:这个字段的历史值还有用吗? 有用(消息、日志、检索结果)就加 reducer,没用(当前状态、当前步骤名)就用默认的覆盖。
照这个标准过一遍典型字段:
| 字段 | 历史值还有用吗 | 结论 |
|---|---|---|
messages 对话历史 |
有用,全都要 | add_messages |
log 执行日志 |
有用 | operator.add |
retrieved_docs 多路检索结果 |
有用,要汇总 | operator.add |
turn 轮次计数 |
要累加 | operator.add |
status 当前状态 |
没用,只关心现在 | 默认覆盖 |
current_step 当前步骤名 |
没用 | 默认覆盖 |
user_id 会话级常量 |
不变,无需保留历史 | 默认覆盖 |
error 最近一次错误 |
看需求:只看最后一个就覆盖,要全部就累加 | 视情况 |
别不管三七二十一给所有字段都加 reducer。 给 status 挂上 operator.add,它会变成一个越来越长的字符串('新建处理中已完成'),而且不会报错——又是一个静默失败。
5. 并行分支:最需要 Reducer 的地方 #
5.1. 复现上一章的报错 #
先统一一个说法。 本章讲并行时反复出现的「同一步」,就是第 20 章 §6.4 里那个 超步(superstep):同一步内能并行的节点一起跑完,统一合并结果,再进入下一步。第 22 章之后统一用「超步」这个词,两者指的是同一件事。
第 20 章 §6.4 那个报错,现在可以正面处理了。两个节点从 START 并行出发,都写 log:
from typing import TypedDict
from langgraph.graph import END, START, StateGraph
# 状态和 §2.1 一样,故意不挂 reducer
class SP(TypedDict):
log: list[str] # 故意不加 reducer
# 建图
builder = StateGraph(SP)
# 第一个节点写 A
builder.add_node("n1", lambda s: {"log": ["A"]})
# 第二个节点写 B
builder.add_node("n2", lambda s: {"log": ["B"]})
builder.add_edge(START, "n1") # 两条边都从 START 出发 = 并行
# 第二条并行边
builder.add_edge(START, "n2")
# 两个节点各自连到 END
builder.add_edge("n1", END)
# 第二个分支的出口
builder.add_edge("n2", END)
graph = builder.compile()
# 用 try 包住,因为预期会抛异常
try:
# 并行写同一个键,看会不会炸
print(graph.invoke({"log": ["初始"]}))
# 捕获异常看原文
except Exception as e:
# 打印异常类型和原文
print("报错:", type(e).__name__, str(e))运行输出:
报错: InvalidUpdateError At key 'log': Can receive only one value per step. Use an Annotated key to handle multiple values.
For troubleshooting, visit: https://docs.langchain.com/oss/python/langgraph/errors/INVALID_CONCURRENT_GRAPH_UPDATE现在这句报错完全读得懂了:同一步里收到了两个值,又没有 reducer 告诉图该怎么合并,所以它宁可报错也不猜。
顺便说一句,这个报错是件好事。对比第 20 章那四种静默失败,LangGraph 在这里选择了「大声失败」——两个并行分支的结果到底谁覆盖谁,猜错了后果很严重。
5.2. 加上 reducer 就通了 #
状态换成 §3.4 里那个带 reducer 的 S5,节点和边都不动:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# S5 的 log 挂了 operator.add
class S5(TypedDict):
log: Annotated[list[str], operator.add]
# n1 / n2 就是 §3.4 里那两个会打印的节点
def n1(state: S5) -> dict:
print(" n1 读到:", state)
return {"log": ["A"]}
def n2(state: S5) -> dict:
print(" n2 读到:", state)
return {"log": ["B"]}
builder = StateGraph(S5)
builder.add_node("n1", n1)
# 第二个节点同样复用
builder.add_node("n2", n2)
# 两条边都从 START 出发,构成并行
builder.add_edge(START, "n1")
# 第二条并行边
builder.add_edge(START, "n2")
# 两个分支各自结束
builder.add_edge("n1", END)
# 第二个分支的出口
builder.add_edge("n2", END)
graph = builder.compile()
# 这次不会报错了
print("结果:", graph.invoke({"log": ["初始"]}))运行输出:
n1 读到: {'log': ['初始']}
n2 读到: {'log': ['初始']}
结果: {'log': ['初始', 'A', 'B']}两个分支的结果都保留了下来。唯一的改动还是在状态定义里——和 §3.2 一样,节点和边一个字没动。
那么并行的两个值是怎么进 reducer 的?把 §3.3 的探针 reducer 挂上再跑一次:
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
# 还是那个会打日志的探针 reducer
def spy_reducer(existing, update):
print(f" reducer: existing={existing!r} update={update!r}")
return (existing or []) + update
class SS(TypedDict):
log: Annotated[list[str], spy_reducer]
builder = StateGraph(SS)
builder.add_node("n1", lambda s: {"log": ["A"]})
builder.add_node("n2", lambda s: {"log": ["B"]})
# 并行结构
builder.add_edge(START, "n1")
builder.add_edge(START, "n2")
builder.add_edge("n1", END)
builder.add_edge("n2", END)
graph = builder.compile()
# 数一数 reducer 被调了几次,以及每次的 existing 是什么
print("结果:", graph.invoke({"log": ["初始"]}))运行输出:
reducer: existing=[] update=['初始']
reducer: existing=['初始'] update=['A']
reducer: existing=['初始', 'A'] update=['B']
结果: {'log': ['初始', 'A', 'B']}注意 reducer 仍然只被调用了 3 次,而且是顺序累积的:['初始'] 先和 ['A'] 合并,结果再和 ['B'] 合并。这就解答了一个容易想不通的问题:
既然 reducer 只有两个参数,那三个并行分支的值怎么塞进去? 答案是不用塞——LangGraph 会一个一个地喂进来,每次都拿上一次的结果当
existing。所以 reducer 永远只需要处理两个值,不管有多少个并行分支。
这也意味着 reducer 里看不到「这一批一共有几个更新」。如果业务上需要「把同一步收到的所有值一起处理」(比如取中位数),reducer 这个形态做不到,得让各分支写各自的字段,再加一个汇总节点。
5.3. 并行时节点读到的是「分叉前的快照」 #
上面输出里有个关键细节:n1 和 n2 都读到 ['初始']。
对比串行时的情况(§3.4):
n1 读到 |
n2 读到 |
|
|---|---|---|
串行 n1 → n2 |
['初始'] |
['初始', 'A'] ← 看得到 n1 的结果 |
并行 n1 ∥ n2 |
['初始'] |
['初始'] ← 看不到 n1 的结果 |
这就是并行的本质:并行分支看到的是同一份分叉前的快照,彼此的写入互不可见。 结果要等这一步结束之后,才由 reducer 统一合并。
这条规则很容易踩坑:
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
log: Annotated[list[str], operator.add]
def n1(state: S) -> dict:
return {"log": ["A"]}
# 危险写法:并行节点里做「读-改-写」
def n2(state: S) -> dict:
# 以为能读到 n1 刚写的值,其实读不到
if "A" in state["log"]:
# 这个分支永远进不来,而且不会报错
return {"log": ["看到 A 了"]}
return {"log": ["没看到 A"]}
builder = StateGraph(S)
builder.add_node("n1", n1)
builder.add_node("n2", n2)
# 两条边都从 START 出发,n1 和 n2 并行
builder.add_edge(START, "n1")
builder.add_edge(START, "n2")
builder.add_edge("n1", END)
builder.add_edge("n2", END)
graph = builder.compile()
print("结果:", graph.invoke({"log": []}))运行输出:
结果: {'log': ['A', '没看到 A']}n1 明明写了 A,n2 却报告「没看到 A」——因为它们在同一步里,读的是同一份分叉前的快照。
这种 bug 特别难查:不报错,而且「有时候还能对」。 如果 n1 写的值恰好在上一轮就已经存在于状态里,条件判断照样能命中,测试就过了;等到某次真的依赖「本轮 n1 刚写的值」,才会悄悄走错分支。
并行节点之间不要有数据依赖。 有依赖就说明它们本来就该串行,把边改成
n1 → n2就行。
判断方法很简单:问一句「这个节点读的字段,会不会被同一步的另一个节点写?」 答案是「会」,就必须改成串行。
5.4. 别依赖并行的合并顺序 #
三个并行节点 p1 / p2 / p3 各往列表里追加一个元素,跑 10 次,结果去重后只有一种:
[('p1', 'p2', 'p3')]顺序看起来很稳定。那它由什么决定? 直觉是「按 add_node 的添加顺序」——但这是错的。做个对照实验:把节点名换成非字母顺序,再故意打乱添加顺序。
import operator
from typing import Annotated, TypedDict
from langgraph.graph import END, START, StateGraph
class S(TypedDict):
log: Annotated[list[str], operator.add]
def build(names: list[str]):
"""按给定的添加顺序建一张并行图,每个节点往 log 里写自己的名字。"""
b = StateGraph(S)
for name in names:
# 默认参数 n=name 是为了避免闭包全都捕获最后一个 name
b.add_node(name, lambda s, n=name: {"log": [n]})
# 每个节点都从 START 出发,构成并行
b.add_edge(START, name)
b.add_edge(name, END)
return b.compile()
# 添加顺序刻意写成 zebra → apple → mango
graph = build(["zebra", "apple", "mango"])
print(" 结果:", graph.invoke({"log": []})["log"])
# 跑 10 次去重,确认顺序是稳定的而不是碰巧
runs = {tuple(graph.invoke({"log": []})["log"]) for _ in range(10)}
print(" 10 次去重:", sorted(runs))运行输出:
结果: ['apple', 'mango', 'zebra']
10 次去重: [('apple', 'mango', 'zebra')]顺序完全没跟着添加顺序走,而是变成了 apple → mango → zebra——按节点名排序。再验证两组:
| 节点名 | 添加顺序 | 实际合并顺序 |
|---|---|---|
p1 / p2 / p3 |
p1, p2, p3 | p1, p2, p3 |
zebra / apple / mango |
zebra, apple, mango | apple, mango, zebra |
3rd / 1st / 2nd |
3rd, 1st, 2nd | 1st, 2nd, 3rd |
丙 / 甲 / 乙 |
丙, 甲, 乙 | 丙, 乙, 甲(按 Unicode 码位) |
结论:并行分支的合并顺序是按节点名排序的,和你 add_node 的顺序无关。(最后一行的中文顺序印证了这是纯粹的字符串排序:丙 U+4E19 < 乙 U+4E59 < 甲 U+7532,跟「甲乙丙」的常识顺序正好错开。)
这个发现让「别依赖顺序」多了一条更具体、更容易踩的理由:
改一个节点名,就可能改变并行分支的合并顺序。 把
search_web重命名成web_search这种纯整理性的改动,会悄悄把它在结果列表里的位置从前面挪到后面。代码评审的时候,没人会觉得改个名字有风险。
顺便说一下为什么要多做这个实验:第一眼看到 p1, p2, p3 就断定「按添加顺序」,是因为那个例子里两种解释给出的答案恰好相同。所以必须专门设计一个能区分两种解释的实验,否则你得到的是一个碰巧正确的预测,加上一个完全错误的因果理解。
再说一遍,这仍然是实现细节,不是官方承诺的行为。 它跟调度实现有关,换个版本、换成异步都可能变。安全的做法是让 reducer 满足「换顺序结果一样」(数学上叫可交换):
| reducer | 换顺序结果一样吗 | 并行是否可放心 |
|---|---|---|
数字累加 operator.add |
一样 | √ |
| 取最大 / 最小 | 一样 | √ |
| 列表去重累加 | 内容一样,顺序不一定 | ⚠️ |
| 字典合并 | 键不冲突时一样 | ⚠️ 键冲突时看顺序 |
列表追加 operator.add |
不一样(元素排列变了) | ⚠️ 内容全在,顺序别依赖 |
列表追加这一条要特别注意:元素不会丢,但先后顺序不该当成保证。 如果顺序有业务含义(比如日志时间线),有三个办法:
| 办法 | 做法 | 代价 |
|---|---|---|
| 改成串行 | 把 START → a、START → b 改成 START → a → b |
慢,但顺序确定 |
| 元素自带排序键 | 往列表里放 {"ts": ..., "text": ...},读的时候自己排 |
稍麻烦,但并行照旧 |
| 各写各的字段 | a 写 result_a,b 写 result_b,再加个汇总节点 |
字段变多,但完全没有顺序问题 |
第三个办法常被忽略,却往往最干净:没有共享字段,就没有合并顺序问题。
6. add_messages:最重要的 Reducer #
这是实战里用得最多的 reducer,专门用来处理消息列表。看起来它只是在「追加消息」,实际有五种行为。
6.1. 为什么消息列表需要专门的 reducer #
用 operator.add 追加消息不行吗?不行。对话场景有几个特殊需求它满足不了:
| 需求 | operator.add |
add_messages |
|---|---|---|
| 新消息追加到末尾 | √ | √ |
| 修改已有的某条消息 | × 只能追加 | √ 按 id 替换 |
| 删除某条消息(裁剪历史) | × | √ RemoveMessage |
| 消息没有 id 时自动补 | × | √ |
接受 {"role": ..., "content": ...} 字典 |
× | √ 自动转换 |
最后一条的后果比看起来严重。实测「给 messages 挂 operator.add,再传一个字符串进去」:
import operator
from typing import Annotated, TypedDict
from langchain_core.messages import AIMessage, AnyMessage
from langgraph.graph import END, START, StateGraph
class WrongState(TypedDict):
# 错误示范:消息列表用了 operator.add
messages: Annotated[list[AnyMessage], operator.add]
builder = StateGraph(WrongState)
builder.add_node("reply", lambda s: {"messages": [AIMessage("回复")]})
builder.add_edge(START, "reply")
builder.add_edge("reply", END)
graph = builder.compile()
# 传一个裸字符串,operator.add 不会做任何转换
out = graph.invoke({"messages": ["我是字符串"]})
# 打印每条的类型和内容,看看混进了什么
print("结果:", [(type(m).__name__, str(getattr(m, "content", m))) for m in out["messages"]])运行输出:
结果: [('str', '我是字符串'), ('AIMessage', '回复')]没有报错,但状态里混进了一个裸 str。 下游任何一个 m.content 都会炸成 AttributeError,而那个报错的位置,离真正的原因(状态定义里用错了 reducer)已经很远。这是本章最典型的静默失败之一。
导入方式,两条路径都行:
# 常用写法:直接从 langgraph.graph 导入
from langgraph.graph import add_messages as a1
# 等价写法:从定义它的模块导入(REMOVE_ALL_MESSAGES 只能从这里拿)
from langgraph.graph.message import add_messages as a2
# 用 is 判断是不是同一个函数对象
print("是同一个函数对象:", a1 is a2)运行输出:
是同一个函数对象: True两处是同一个函数对象,用哪个都一样。
下面逐个验证这五种行为。为了看得清楚,这里直接把 add_messages 当普通函数调用——它本来就是普通函数,两个参数一个返回值,和 §4.1 的自定义 reducer 完全一样。这个手法也值得学起来:任何 reducer 都能脱离图单独调用来摸清行为,不必每次都建一张图。
6.2. 行为一:追加 #
# 两种最常见的消息类型
from langchain_core.messages import AIMessage, HumanMessage
# 被测的 reducer
from langgraph.graph import add_messages
# 手动指定 id,方便观察哪条是哪条
existing = [HumanMessage("你好", id="h1")]
# 第二个参数是「本次更新」,注意要用列表包起来
print(add_messages(existing, [AIMessage("你好呀", id="a1")]))运行输出(略去空字段):
[HumanMessage(content='你好', id='h1'),
AIMessage(content='你好呀', id='a1', tool_calls=[], invalid_tool_calls=[])]两条都在,新的排在后面。这是最常见的情况,也是 create_agent 每一轮对话在做的事。
6.3. 行为二:同 id 会替换,不是追加 #
这是 add_messages 和 operator.add 最本质的区别。
from langchain_core.messages import AIMessage, HumanMessage
from langgraph.graph import add_messages
# 准备两条消息,各有自己的 id
existing = [HumanMessage("原始内容", id="h1"), AIMessage("回答", id="a1")]
# 更新时用同一个 id
merged = add_messages(existing, [HumanMessage("改过的内容", id="h1")])
# 遍历结果,打印类型、id 和内容
for m in merged:
# 用 type(m).__name__ 看消息类型
print(f" {type(m).__name__}(id={m.id}) {m.content}")运行输出:
HumanMessage(id=h1) 改过的内容
AIMessage(id=a1) 回答结果仍然是两条,h1 的内容被替换了,而且位置没变,还在第一条。「位置不变」这点很关键——要是替换会把消息挪到末尾,对话顺序就乱了。
这个行为是流式输出、消息修正、PII 脱敏的基础,它们都需要「改写已有消息」而不是「再追加一条」。第 10 章那个 PII middleware 能把已生成的消息脱敏,靠的就是这个机制。
记住这条规则:id 相同 = 更新,id 不同 = 追加。
还有一个容易忽略的细节:同 id 时,消息类型也会被替换掉。
from langchain_core.messages import AIMessage, HumanMessage
from langgraph.graph import add_messages
existing = [HumanMessage("原始内容", id="h1"), AIMessage("回答", id="a1")]
# 已有的 h1 是 HumanMessage,这次用同一个 id 传一条 AIMessage
r = add_messages(existing, [AIMessage("换成 AI 了", id="h1")])
for m in r:
print(f" {type(m).__name__}(id={m.id}) {m.content}")运行输出:
AIMessage(id=h1) 换成 AI 了
AIMessage(id=a1) 回答原来那条 HumanMessage 整条被换成了 AIMessage,位置还在第一个。也就是说 add_messages 按 id 做的是「整条替换」,不是「只改 content」。正常用的时候很少碰到这一点,但手写更新时如果复用了错误的 id,就可能把用户消息悄悄变成 AI 消息——历史从此就不对了,而且不会报错。
6.4. 行为三:不给 id 就自动生成 #
from langchain_core.messages import HumanMessage
from langgraph.graph import add_messages
# 第一次:从空列表开始,加一条不带 id 的消息
m1 = add_messages([], [HumanMessage("第一句")])
# 第二次:再加一条内容完全相同、同样不带 id 的消息
m2 = add_messages(m1, [HumanMessage("第一句")])
# 如果按内容去重就该是 1 条,按 id 去重就是 2 条
print("两条内容相同但 id 不同,所以都保留:", len(m2))
# 把自动生成的 id 打出来看看长什么样
for m in m2:
# 两条的 content 一样,但 id 是两个不同的 UUID
print(f" id={m.id} {m.content}")运行输出(UUID 每次跑都不一样,只看两个 id 是否不同就够了):
两条内容相同但 id 不同,所以都保留: 2
id=094d658b-49ae-46f4-9b84-cf65f0cd6e62 第一句
id=02622719-7732-499b-85b8-6e8f49cea382 第一句没有 id 的消息会被自动补一个 UUID。所以内容完全相同的两条消息也会各占一条——去重看的是 id,不是内容。
这个设定是对的:用户连着说两遍「在吗」,那就是两条消息,不该被当成重复吞掉一条。
反过来,想更新某条消息,就必须先拿到它的 id,而且只能从状态里读出来,不能自己编:
from langchain_core.messages import AIMessage, HumanMessage
from langgraph.graph import END, START, MessagesState, StateGraph
# 正确做法:从状态里取出真实 id,再用同一个 id 构造新消息
def upper(state: MessagesState) -> dict:
last = state["messages"][-1]
# id=last.id 是关键,只有 id 一致才会触发「替换」
return {"messages": [AIMessage(str(last.content).upper(), id=last.id)]}
# 对照组:忘了带 id
def upper_no_id(state: MessagesState) -> dict:
last = state["messages"][-1]
return {"messages": [AIMessage(str(last.content).upper())]}
for name, node in [("带 id", upper), ("漏了 id", upper_no_id)]:
builder = StateGraph(MessagesState)
builder.add_node("n", node)
builder.add_edge(START, "n")
builder.add_edge("n", END)
graph = builder.compile()
out = graph.invoke({"messages": [HumanMessage("hello", id="h1")]})
# 条数就是关键差别:替换是 1 条,追加是 2 条
print(f"{name}: {len(out['messages'])} 条 ->",
[str(m.content) for m in out["messages"]])运行输出:
带 id: 1 条 -> ['HELLO']
漏了 id: 2 条 -> ['hello', 'HELLO']漏了 id=last.id 时,add_messages 会自动补一个新 UUID,于是你以为在改,实际是追加了一条——历史里出现两条内容相似的消息。这就是 §9 坑表里「更新消息却变成了追加一条新的」的根因。
6.5. 行为四:自动转换格式 #
不一定要传 Message 对象,字典也行:
from langgraph.graph import add_messages
# 传一个 {"role": ..., "content": ...} 形式的字典
out = add_messages([], [{"role": "user", "content": "字典写法"}])
# 检查三件事:转成了什么类型、内容对不对、有没有补 id
print(type(out[0]).__name__, "|", out[0].content, "| 自动补了 id:", out[0].id is not None)运行输出:
HumanMessage | 字典写法 | 自动补了 id: True四种 role 的映射关系,实测如下:
role |
转成的类型 | 备注 |
|---|---|---|
"user" |
HumanMessage |
最常用 |
"assistant" |
AIMessage |
|
"system" |
SystemMessage |
|
"tool" |
ToolMessage |
还必须带 tool_call_id,否则报错 |
这解释了为什么前面章节可以写 agent.invoke({"messages": [{"role": "user", "content": "..."}]})——是 add_messages 在入口把字典转成了 HumanMessage(回想 §3.3:初始输入也会走 reducer)。
这两个机制合起来才解释得通,缺一个都说不通:
你写的: invoke({"messages": [{"role": "user", "content": "你好"}]})
│
▼ ① §3.3:初始输入也会被送进 reducer
add_messages([], [{"role": "user", ...}])
│
▼ ② §6.5:add_messages 认识字典,负责转换
[HumanMessage(content='你好', id='自动生成的 UUID')]它甚至不要求必须是列表——单独一条消息、甚至一个裸字符串都能接受:
from langchain_core.messages import HumanMessage
from langgraph.graph import add_messages
# 不包 list,直接给一条消息
print(add_messages([], HumanMessage("裸消息"))[0].content)
# 甚至直接给字符串,也会被转成 HumanMessage
r = add_messages([], "裸字符串")
print(type(r[0]).__name__, "|", r[0].content)运行输出:
裸消息
HumanMessage | 裸字符串6.6. 行为五:用 RemoveMessage 删除 #
裁剪对话历史时需要「删掉某几条」。add_messages 用一个特殊的消息类型来表达删除:
# RemoveMessage 是一个「表示删除」的特殊消息类型
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import add_messages
# 准备三条带 id 的消息
existing = [
HumanMessage("第一条", id="m1"),
AIMessage("第二条", id="m2"),
HumanMessage("第三条", id="m3"),
]
# 想删中间那条,就传一条 id 指向它的 RemoveMessage
after = add_messages(existing, [RemoveMessage(id="m2")])
# 用列表推导式只打印 id 和内容,输出干净些
print("删除 m2 后:", [(m.id, m.content) for m in after])运行输出:
删除 m2 后: [('m1', '第一条'), ('m3', '第三条')]这个设计初看有点绕:「删除」居然要靠「添加一条删除消息」来表达。但放到图里就合理了——节点只能返回「更新」,所以删除也得写成一种更新。
换个角度看会更顺:RemoveMessage 不是一条消息,而是一条指令。 它借用了消息的外壳(因为节点只能往 messages 字段里放消息),但 add_messages 看到它时执行的是「删掉这个 id」,而不是「把它追加进列表」,所以它自己不会出现在结果里。
删除和新增可以放在同一次更新里,按顺序处理:
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import add_messages
existing = [
HumanMessage("第一条", id="m1"),
AIMessage("第二条", id="m2"),
HumanMessage("第三条", id="m3"),
]
# 一次更新里同时删一条、加一条
r = add_messages(existing, [RemoveMessage(id="m1"), HumanMessage("新的", id="m4")])
print(" ", [(m.id, m.content) for m in r])运行输出:
[('m2', '第二条'), ('m3', '第三条'), ('m4', '新的')]m1 被删掉,m4 追加到末尾。这个能力就是 §8 那个 trim 节点和「摘要压缩」(§6.7)的基础。
删一个不存在的 id 会报错:
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import add_messages
existing = [
HumanMessage("第一条", id="m1"),
AIMessage("第二条", id="m2"),
HumanMessage("第三条", id="m3"),
]
# 预期会抛异常,所以用 try 包住
try:
# 这个 id 是编的,历史里没有
add_messages(existing, [RemoveMessage(id="不存在")])
# 捕获异常看原文
except Exception as e:
# 打印异常类型和原文
print("报错:", type(e).__name__, str(e))报错: ValueError Attempting to delete a message with an ID that doesn't exist ('不存在')所以要删消息,先从
state["messages"]里取真实的 id,别自己编。
这个「大声失败」是有意为之:删一条不存在的消息,通常说明 id 的来源本身就错了(比如是从上一轮的旧快照里读的),继续跑只会错得更远。对比 §6.4,「更新时漏了 id」是静默追加——同一个 add_messages,删除比更新严格得多。
6.7. 清空:REMOVE_ALL_MESSAGES #
想一次清空全部历史,不用挨个写 RemoveMessage:
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import add_messages
# 这个常量只在 langgraph.graph.message 里,不在 langgraph.graph
from langgraph.graph.message import REMOVE_ALL_MESSAGES
existing = [
HumanMessage("第一条", id="m1"),
AIMessage("第二条", id="m2"),
HumanMessage("第三条", id="m3"),
]
# 先看看它到底是什么,其实就是一个约定的魔法字符串
print("常量值:", repr(REMOVE_ALL_MESSAGES))
# 把它当 id 传给 RemoveMessage,就表示「全删」
print("清空后:", add_messages(existing, [RemoveMessage(id=REMOVE_ALL_MESSAGES)]))运行输出:
常量值: '__remove_all__'
清空后: []它只是一个字符串常量 '__remove_all__',没有别的花样。所以理论上写 RemoveMessage(id="__remove_all__") 效果一样,但别这么写:用常量才能在版本变动时不出问题,读代码的人也更容易看懂。
更实用的是「清空并重建」——同一次更新里先删光再放新的:
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import add_messages
from langgraph.graph.message import REMOVE_ALL_MESSAGES
existing = [
HumanMessage("第一条", id="m1"),
AIMessage("第二条", id="m2"),
HumanMessage("第三条", id="m3"),
]
# 一次更新里放两个元素:先全删,再放一条新的
after2 = add_messages(
existing,
# 顺序很关键:RemoveMessage 在前,新消息在后
[RemoveMessage(id=REMOVE_ALL_MESSAGES), HumanMessage("重新开始", id="new1")],
)
# 只剩新放进去的那一条
print("清空并重建:", [(m.id, m.content) for m in after2])运行输出:
清空并重建: [('new1', '重新开始')]列表里的顺序有意义,写反了结果完全不同。 实测把两个元素颠倒过来:
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import add_messages
from langgraph.graph.message import REMOVE_ALL_MESSAGES
existing = [
HumanMessage("第一条", id="m1"),
AIMessage("第二条", id="m2"),
HumanMessage("第三条", id="m3"),
]
# 错误示范:先放新消息,再清空
after3 = add_messages(
existing,
# 新消息在前,RemoveMessage 在后
[HumanMessage("重新开始", id="new1"), RemoveMessage(id=REMOVE_ALL_MESSAGES)],
)
print(" 结果:", [(m.id, m.content) for m in after3])运行输出:
结果: []全空了,刚放进去的新消息也被一起清掉了。 因为 add_messages 是按列表顺序依次处理的,REMOVE_ALL_MESSAGES 执行时会把「当前已经处理好的全部内容」都删掉,包括同一批里排在它前面的新消息。
记住顺序:先删,后加。 这个坑不报错,症状是「摘要压缩之后消息列表空了,模型开始胡说」——连摘要本身都被删掉了。
这正是摘要压缩的做法:把长历史总结成一条摘要,再一次性换掉全部原始消息。第 11 章的 SummarizationMiddleware 底层就是这么干的。
from langchain_core.messages import AIMessage, HumanMessage, RemoveMessage
from langgraph.graph import END, START, MessagesState, StateGraph
from langgraph.graph.message import REMOVE_ALL_MESSAGES
# 为了不花钱,这里用一个假模型顶替真实的 init_chat_model
class FakeModel:
def invoke(self, messages):
# 真实场景下这里是模型生成的摘要
return AIMessage(f"用户一共说了 {len(messages)} 条")
model = FakeModel()
# 摘要压缩的骨架:一次更新完成「删光 + 放摘要」
def summarize(state: MessagesState) -> dict:
# 先请模型把历史总结成一段话
summary = model.invoke([*state["messages"], HumanMessage("请总结以上对话")])
# 关键:RemoveMessage 必须在前,摘要在后
return {"messages": [
RemoveMessage(id=REMOVE_ALL_MESSAGES),
AIMessage(f"【历史摘要】{summary.content}"),
]}
builder = StateGraph(MessagesState)
builder.add_node("summarize", summarize)
builder.add_edge(START, "summarize")
builder.add_edge("summarize", END)
graph = builder.compile()
# 传三条历史进去,出来应该只剩一条摘要
out = graph.invoke({"messages": [
HumanMessage("第一句"), AIMessage("第二句"), HumanMessage("第三句"),
]})
print("压缩后:", [(type(m).__name__, m.content) for m in out["messages"]])运行输出:
压缩后: [('AIMessage', '【历史摘要】用户一共说了 4 条')]6.8. 五种行为总表 #
| 你返回什么 | add_messages 做什么 |
|---|---|
| 一条新消息(id 不存在 / 没 id) | 追加到末尾 |
| 一条消息,id 和已有的相同 | 整条替换那条,位置不变(连类型一起换) |
| 消息没有 id | 自动补 UUID |
{"role": ..., "content": ...} 字典 |
转成 对应的 Message 对象 |
| 一个裸字符串 | 转成 HumanMessage |
RemoveMessage(id="xxx") |
删除 那条(id 不存在会报错) |
RemoveMessage(id=REMOVE_ALL_MESSAGES) |
清空全部(注意它在列表里的位置) |
把「严不严格」单独拎出来,方便预判出问题时的症状:
| 操作 | 出错时的表现 |
|---|---|
| 更新时漏了 id | 静默变成追加 |
| 更新时用错了 id(指向别的消息) | 静默把那条消息覆盖掉 |
| 删除时 id 不存在 | 报错 |
REMOVE_ALL_MESSAGES 位置放错 |
静默把新消息也删掉 |
三个静默、一个报错。 所以只要写了消息更新或删除的代码,收尾一定要打印一遍 messages 的 id 和条数来核对,别只看内容。
7. MessagesState:预置的对话状态 #
7.1. 它其实只有一行 #
对话类的图几乎都要用到一个挂了 add_messages 的 messages 字段,LangGraph 干脆预置了一个。与其只信文档,不如直接把它的定义打印出来看看:
# MessagesState 是预置的状态类
from langgraph.graph import MessagesState
# __annotations__ 能看到一个 TypedDict 声明了哪些字段
print("MessagesState 注解:", MessagesState.__annotations__)运行输出:
MessagesState 注解: {'messages': ForwardRef('Annotated[list[AnyMessage], add_messages]', module='langgraph.graph.message')}就是这么一行,没有别的花样(ForwardRef 只是延迟求值的包装,内容就是那个 Annotated)。等价于你自己写:
from typing import Annotated, TypedDict
# AnyMessage 是所有消息类型的联合类型
from langchain_core.messages import AIMessage, AnyMessage, HumanMessage
# 手动导入 reducer
from langgraph.graph import END, START, MessagesState, StateGraph, add_messages
# 自己写一遍,效果和 MessagesState 完全一样
class MyState(TypedDict):
# 这一行就是 MessagesState 的全部内容
messages: Annotated[list[AnyMessage], add_messages]
# 两种状态各建一张图,行为应该分毫不差
for name, state_cls in [("MessagesState", MessagesState), ("MyState", MyState)]:
builder = StateGraph(state_cls)
builder.add_node("reply", lambda s: {"messages": [AIMessage("回复")]})
builder.add_edge(START, "reply")
builder.add_edge("reply", END)
graph = builder.compile()
out = graph.invoke({"messages": [HumanMessage("你好")]})
# 条数和消息类型都打出来对照
print(f"{name}: {len(out['messages'])} 条 ->",
[type(m).__name__ for m in out["messages"]])运行输出:
MessagesState: 2 条 -> ['HumanMessage', 'AIMessage']
MyState: 2 条 -> ['HumanMessage', 'AIMessage']两张图跑出来分毫不差。用 MessagesState 的好处,就是少写一行、少一个 import。
不过「知道它只有一行」这件事本身是有用的,因为它顺便回答了几个常见疑问:
| 疑问 | 答案 |
|---|---|
| 它会自动帮我存历史吗? | 不会。跨轮次保留状态靠 checkpointer(§8) |
| 它会自动裁剪超长历史吗? | 不会。裁剪要自己写节点(§8 的 trim) |
| 它能装非对话数据吗? | 只有 messages 一个字段,别的要自己加(§7.3) |
| 用它和自己写有性能差别吗? | 没有,同一个东西 |
也就是说,MessagesState 只帮你省了一行代码,不提供任何额外能力。 别把「用了 MessagesState 就有记忆了」当成事实。
7.2. 直接用 #
# 两种消息类型
from langchain_core.messages import AIMessage, HumanMessage
# MessagesState 可以直接从 langgraph.graph 拿
from langgraph.graph import END, START, MessagesState, StateGraph
# 一个假装在回复的节点,不调模型
def reply(state: MessagesState) -> dict:
# 打印进入节点时的消息条数,用来对照最终结果
print(" 节点读到消息数:", len(state["messages"]))
# 只返回新生成的这一条,add_messages 负责追加
return {"messages": [AIMessage("这是回复")]}
# 直接把 MessagesState 当状态类型传给 StateGraph
builder = StateGraph(MessagesState)
# 注册唯一的节点
builder.add_node("reply", reply)
# 入口边
builder.add_edge(START, "reply")
# 出口边
builder.add_edge("reply", END)
# 编译成可运行的图
graph = builder.compile()
# 传一条用户消息进去
out = graph.invoke({"messages": [HumanMessage("你好")]})
# 看最终有几条:如果是 2 就说明追加成功
print("结果消息数:", len(out["messages"]))
# 逐条打印类型和内容
for m in out["messages"]:
# 第一条是传进去的 Human,第二条是节点返回的 AI
print(f" {type(m).__name__}: {m.content}")运行输出:
节点读到消息数: 1
结果消息数: 2
HumanMessage: 你好
AIMessage: 这是回复节点只返回了一条 AIMessage,最终状态里有两条——add_messages 把它追加到了历史后面。这就是对话记忆最底层的形态。
这里的分工值得细看一下:
节点的职责: 只生成「这一轮新增的内容」 → [AIMessage("这是回复")]
状态的职责: 决定新增内容怎么并入历史 → add_messages
你的职责: 什么都不用做节点完全不用知道历史有多长,也不用自己拼列表。 对比手写对话循环里到处都是 history.append(ai_msg),差别就在于:合并规则是集中在状态定义这一处,还是散落在每一个写入点。
连字符串都能直接传(承接 §6.5 那个宽容度):
from langchain_core.messages import AIMessage
from langgraph.graph import END, START, MessagesState, StateGraph
def reply(state: MessagesState) -> dict:
print(" 节点读到消息数:", len(state["messages"]))
return {"messages": [AIMessage("这是回复")]}
builder = StateGraph(MessagesState)
builder.add_node("reply", reply)
builder.add_edge(START, "reply")
builder.add_edge("reply", END)
graph = builder.compile()
# messages 不给列表,直接给一个字符串
out = graph.invoke({"messages": "我是一个字符串"})
# 逐条打印类型和内容,确认裸字符串被转成了 HumanMessage
for m in out["messages"]:
# 第一条应是 HumanMessage,说明归一化生效了
print(f" {type(m).__name__}: {m.content}")运行输出:
节点读到消息数: 1
HumanMessage: 我是一个字符串
AIMessage: 这是回复字符串被自动转成了 HumanMessage(§6.5 的行为,作用在初始输入上)。传字典也一样:
from langchain_core.messages import AIMessage
from langgraph.graph import END, START, MessagesState, StateGraph
def reply(state: MessagesState) -> dict:
print(" 节点读到消息数:", len(state["messages"]))
return {"messages": [AIMessage("这是回复")]}
builder = StateGraph(MessagesState)
builder.add_node("reply", reply)
builder.add_edge(START, "reply")
builder.add_edge("reply", END)
graph = builder.compile()
# 字典形式,最常见的写法
out = graph.invoke({"messages": [{"role": "user", "content": "我是字典"}]})
# 逐条打印类型和内容,确认字典也被转成了 HumanMessage
for m in out["messages"]:
# role="user" 对应 HumanMessage
print(f" {type(m).__name__}: {m.content}")运行输出:
节点读到消息数: 1
HumanMessage: 我是字典
AIMessage: 这是回复三种写法(消息对象、字典、裸字符串)效果完全一样。这不是三个特性,而是一个特性的三种表现:add_messages 在处理初始输入时会做一次格式归一化。
7.3. 继承扩展自己的槽位 #
MessagesState 是个普通的 TypedDict,可以直接继承:
import operator
from typing import Annotated
from langgraph.graph import MessagesState
# §4.2 那个去重累加的 reducer
def merge_unique(old: list[str], new: list[str]) -> list[str]:
merged = list(old or [])
for item in new:
if item not in merged:
merged.append(item)
return merged
# 继承而不是重写,messages 字段连同它的 reducer 一起被继承下来
class ChatState(MessagesState):
"""继承 MessagesState 拿到 messages,再加两个自定义槽位。"""
# 轮次计数:每个节点返回 1,靠 operator.add 累加
turn: Annotated[int, operator.add]
# 已识别话题:用 §4.2 的自定义 reducer 去重累加
topics: Annotated[list[str], merge_unique]
# 三个字段都在,各自带着自己的 reducer
print("继承后的注解键:", list(ChatState.__annotations__))运行输出:
继承后的注解键: ['messages', 'turn', 'topics']这个模式很有用,因为消息历史会被裁剪、会被摘要压缩,但业务上的关键信息不能丢。第 11 章 §7.4 讲过的「给关键信息开专用槽位」,在图这一层就是这么实现的。
| 放哪 | 适合什么 | 会不会被裁剪 |
|---|---|---|
messages |
对话内容 | 会(裁剪 / 摘要) |
| 自定义槽位 | 订单号、客户名、累计金额、已识别话题 | 不会 |
但给槽位挂 reducer 会带来一个副作用,这是 §3.3 第二条结论的直接后果:你没法通过 invoke 给这种字段「设定」一个值。 实测——turn 挂了 operator.add,图里节点每次返回 1,现在在输入里传 turn=10:
import operator
from typing import Annotated
from langchain_core.messages import AIMessage, HumanMessage
from langgraph.graph import END, START, MessagesState, StateGraph
class ChatState(MessagesState):
# turn 挂了 operator.add
turn: Annotated[int, operator.add]
# 节点每轮只贡献 1
def step(state: ChatState) -> dict:
return {"turn": 1, "messages": [AIMessage("回复")]}
builder = StateGraph(ChatState)
builder.add_node("step", step)
builder.add_edge(START, "step")
builder.add_edge("step", END)
graph = builder.compile()
# 想「把轮次设为 10」,于是在输入里传 turn=10
out = graph.invoke({"messages": [HumanMessage("问题")], "turn": 10})
# 期待是 10,实际会是多少?
print(" 传 turn=10 时:", out["turn"])运行输出:
传 turn=10 时: 11得到的是 11,不是 10。 因为初始值 10 也走了 reducer(0 + 10 = 10),然后节点又贡献了 1(10 + 1 = 11)。
有 reducer 的字段只能「贡献」值,不能「设定」值。 想要「设定」语义(比如重置计数器、覆盖状态),有两个办法:
- 别给这个字段挂 reducer,用默认的覆盖行为
- 需要「既能累加又能重置」时,给 reducer 加一个哨兵值,比如
lambda old, new: 0 if new is None else (old or 0) + new——传None表示重置
想「从中间某个状态恢复」的时候最容易撞上这个坑:你从数据库读出上次的 turn=5 塞进 invoke,结果它和状态里已有的值又加了一遍。
8. 实战:可追加 messages 的状态图 #
把本章内容合成一张多轮对话图,同时用上三种 reducer。
8.1. 设计 #
状态 ChatState
├── messages : add_messages ← 对话历史,追加 + 可删除
├── turn : operator.add ← 轮次计数,累加
└── topics : merge_unique ← 已识别话题,去重累加
流程
START → detect_topic(确定性:识别话题、轮次+1)
→ respond(调模型生成回复)
→ trim(超长时用 RemoveMessage 删最早的)
→ END三个节点的分工,就是第 19 章「图套 Agent」思路的落地:
| 节点 | 确定性 | 花不花钱 | 为什么这么分 |
|---|---|---|---|
detect_topic |
完全确定 | 不花 | 关键词匹配这种规则,交给模型是浪费 |
respond |
不确定 | 花钱 | 只有「生成自然语言」这一步真的需要模型 |
trim |
完全确定 | 不花 | 裁剪是纯粹的算术 |
三个节点里只有一个调模型,另外两个都是纯 Python。这就是把确定性逻辑从 create_agent 里拆出来的直接好处。
8.2. 流程图 #
8.3. chat_state_graph.py #
"""实战:可追加 messages 的状态图"""
# 让类型标注延迟求值,可以写 list[str] 而不用 typing.List
from __future__ import annotations
# operator.add 当 reducer 用
import operator
# Annotated 用来挂 reducer
from typing import Annotated
# 从 .env 读 DEEPSEEK_API_KEY
from dotenv import load_dotenv
# 统一的模型初始化入口
from langchain.chat_models import init_chat_model
# 裁剪历史要用的「删除指令」
from langchain_core.messages import RemoveMessage
# 内存版 checkpointer,让状态能跨 invoke 保留
from langgraph.checkpoint.memory import InMemorySaver
# 图的三件套 + 预置对话状态
from langgraph.graph import END, START, MessagesState, StateGraph
# override=True 保证 .env 里的值覆盖掉系统环境变量
load_dotenv(override=True)
# 关键词 -> 话题,用最简单的规则做识别
KEYWORDS = {
# 这三个词都归到「售后」
"退货": "售后",
"退款": "售后",
"换货": "售后",
# 这两个归「物流」
"发货": "物流",
"快递": "物流",
# 人事类
"年假": "人事",
# 设备类
"滤芯": "设备",
}
# 消息超过这个数就裁剪
MAX_MESSAGES = 6
def merge_unique(old: list[str], new: list[str]) -> list[str]:
"""自定义 reducer:累加但去重,且保持首次出现的顺序。"""
# 拷一份,不原地修改 old(§4.1)
merged = list(old or [])
# 逐个检查本次新识别到的话题
for item in new:
# 已经记过的就跳过
if item not in merged:
# 没记过才追加,保留首次出现的位置
merged.append(item)
# 返回新列表
return merged
class ChatState(MessagesState):
"""继承 MessagesState 拿到 messages,再加两个自定义槽位。"""
# 轮次:每轮 detect_topic 返回 1,累加成总数
turn: Annotated[int, operator.add]
# 话题:去重累加,不会因为重复提到而膨胀
topics: Annotated[list[str], merge_unique]
def detect_topic(state: ChatState) -> dict:
"""确定性节点:从最新一条用户消息里识别话题,并把轮次 +1。"""
# 最后一条就是本轮用户刚发的消息
last = state["messages"][-1]
# content 可能是列表(多模态),统一转成字符串再匹配
text = str(last.content)
# 遍历关键词表,命中的话题都收集起来(可能一条都没有)
found = [topic for kw, topic in KEYWORDS.items() if kw in text]
# turn 返回 1,靠 operator.add 累加成总轮次
return {"turn": 1, "topics": found}
def respond(state: ChatState) -> dict:
"""调模型生成回复。系统提示里带上槽位信息。"""
# temperature=0 让回复尽量稳定
model = init_chat_model("deepseek:deepseek-v4-flash", temperature=0)
# 把话题列表拼成中文串,空列表时给一个兜底说法
topics = "、".join(state["topics"]) or "暂未识别"
# 把两个槽位的信息注入系统提示,这是槽位「被用上」的方式
system = (
# turn 来自 operator.add 累加的结果
f"你是客服助手。当前是第 {state['turn']} 轮对话,"
# topics 来自 merge_unique 去重累加的结果
f"已识别到的话题:{topics}。回答控制在两句话以内。"
)
# 系统提示放最前,后面接完整的对话历史
ai = model.invoke([{"role": "system", "content": system}, *state["messages"]])
# 只返回新生成的这一条,add_messages 负责追加
return {"messages": [ai]}
def trim(state: ChatState) -> dict:
"""超长时删掉最早的消息,用 RemoveMessage 表达『删除』。"""
# 此时 messages 里已经包含 respond 刚写入的 AI 回复
messages = state["messages"]
# 没超限就什么都不做
if len(messages) <= MAX_MESSAGES:
# 不需要裁剪时返回空字典,什么都不改
return {}
# 切出最前面多余的那几条
drop = messages[: len(messages) - MAX_MESSAGES]
# 打印一行日志,方便在输出里看到裁剪发生的时刻
print(f" [trim] 消息 {len(messages)} 条,删掉最早 {len(drop)} 条")
# 用真实的 id 构造 RemoveMessage(§6.6)
return {"messages": [RemoveMessage(id=m.id) for m in drop]}
# 用自定义状态建图
builder = StateGraph(ChatState)
# 注册话题识别节点
builder.add_node("detect_topic", detect_topic)
# 注册模型回复节点
builder.add_node("respond", respond)
# 注册裁剪节点
builder.add_node("trim", trim)
# 入口:每轮都必须先识别话题
builder.add_edge(START, "detect_topic")
# 识别完再让模型回复
builder.add_edge("detect_topic", "respond")
# 回复完检查要不要裁剪
builder.add_edge("respond", "trim")
# 裁剪完本轮结束
builder.add_edge("trim", END)
# 加 checkpointer 才能跨轮次保留状态(第 11 章)
graph = builder.compile(checkpointer=InMemorySaver())
# 直接运行本文件才执行演示
if __name__ == "__main__":
# thread_id 标识一个会话,同一个 id 才能共享状态
config = {"configurable": {"thread_id": "demo-001"}}
# 四轮对话,故意设计成能触发去重和裁剪
TURNS = [
# 命中「退货」-> 售后
"我想退货,多少天内可以?",
# 命中「换货」-> 还是售后,用来验证去重
"那换货呢?",
# 命中「发货」-> 物流,新话题
"顺便问下发货一般要多久",
# 命中「滤芯」-> 设备,这一轮会触发裁剪
"滤芯多久换一次?",
]
# 逐轮对话
for text in TURNS:
# 每轮只传新的用户消息,历史由 checkpointer + add_messages 维护
out = graph.invoke({"messages": [{"role": "user", "content": text}]}, config)
# 先回显用户这一轮说了什么
print(f"\n用户:{text}")
# 最后一条就是模型刚生成的回复
print(f"助手:{out['messages'][-1].content}")
# 把三个 reducer 的效果一起打印出来
print(
f" [状态] turn={out['turn']} topics={out['topics']} "
# 条数会体现追加和裁剪的净效果
f"messages={len(out['messages'])} 条"
)
# 分隔标题
print("\n=== 最终快照 ===")
# get_state 从 checkpointer 里读当前状态(第 11 章)
snap = graph.get_state(config)
# 累加出来的总轮次
print("turn:", snap.values["turn"])
# 去重累加出来的话题列表
print("topics:", snap.values["topics"])
# 下面逐条列出留存的消息
print("messages:")
# 逐条打印剩下的消息,验证最早的几条确实被删了
for m in snap.values["messages"]:
# 内容截断到 40 字,避免输出太长
print(f" {type(m).__name__}: {str(m.content)[:40]}")8.4. 运行结果 #
用户:我想退货,多少天内可以?
助手:您好,一般商品支持签收后7天内无理由退货,但需保持商品完好。具体请以商品页面标注的售后政策为准。
[状态] turn=1 topics=['售后'] messages=2 条
用户:那换货呢?
助手:换货一般支持签收后15天内申请,需保持商品完好。具体请以商品页面标注的售后政策为准。
[状态] turn=2 topics=['售后'] messages=4 条
用户:顺便问下发货一般要多久
助手:一般现货商品48小时内发出,预售商品以页面标注时间为准。请留意物流信息更新。
[状态] turn=3 topics=['售后', '物流'] messages=6 条
[trim] 消息 8 条,删掉最早 2 条
用户:滤芯多久换一次?
助手:一般建议每3-6个月更换一次,具体视使用频率和水质而定。请参考产品说明书中的建议周期。
[状态] turn=4 topics=['售后', '物流', '设备'] messages=6 条(回复措辞每次跑可能略有不同,但 turn / topics / messages 这三个数字是确定的——它们由 reducer 算出来,不由模型决定。)
三种 reducer 的效果很清楚:
| 字段 | reducer | 观察到的变化 |
|---|---|---|
messages |
add_messages |
2 → 4 → 6 → 6(第 4 轮触发裁剪,稳定在 6 条) |
turn |
operator.add |
1 → 2 → 3 → 4(每个节点只返回 1) |
topics |
merge_unique |
['售后'] → ['售后'] → +物流 → +设备 |
第 2 轮特别值得看:用户问「那换货呢?」,detect_topic 识别出「售后」,但这个话题已经在列表里了,merge_unique 没有重复添加。这就是自定义 reducer 的价值——要是用 operator.add,topics 会变成 ['售后', '售后', '物流', ...],注入系统提示之后就成了「已识别到的话题:售后、售后、物流」,看起来很不专业。
第 4 轮的 [trim] 那行日志也值得对一下数:
第 3 轮结束时 6 条
第 4 轮 respond 后 6 + 2 = 8 条 ← 一问一答各一条
trim 判断 8 > 6,删最早 2 条
第 4 轮结束 6 条注意 trim 是在 respond 之后跑的,所以它看到的是 8 条而不是 7 条。要是把 trim 放到 respond 之前,裁剪就会晚一轮才生效——这里节点顺序是有实质影响的。
最终快照证明裁剪真的生效了——第一轮对话已经不在了:
=== 最终快照 ===
turn: 4
topics: ['售后', '物流', '设备']
messages:
HumanMessage: 那换货呢?
AIMessage: 换货一般支持签收后15天内申请,需保持商品完好。具体请以商品页面标注的售后政策为
HumanMessage: 顺便问下发货一般要多久
AIMessage: 一般现货商品48小时内发出,预售商品以页面标注时间为准。请留意物流信息更新。
HumanMessage: 滤芯多久换一次?
AIMessage: 一般建议每3-6个月更换一次,具体视使用频率和水质而定。请参考产品说明书中的建议注意 turn 和 topics 完好无损。 消息历史被裁掉了一截,但「这是第 4 轮」「聊过售后、物流、设备」这些信息还在——这正是 §7.3 说的「给关键信息开专用槽位」的意义。
8.5. 验收清单 #
- 消息在追加:每轮
messages增长 2 条(一问一答) - 计数在累加:
turn从 1 涨到 4,而节点每次只返回1 - 去重生效:第 2 轮的「售后」没有重复出现在
topics里 - 裁剪生效:第 4 轮后
messages稳定在 6 条,不再增长 - 槽位不受裁剪影响:
turn和topics在裁剪后依然完整 - 模型能读到槽位:回复里体现了「第 N 轮」和已识别话题(通过系统提示注入)
9. 实用约定与坑 #
| 约定 | 说明 |
|---|---|
| 先问「历史值还有用吗」 | 有用加 reducer,没用用默认覆盖(§4.5) |
| 别给所有字段都挂 reducer | 给 status 挂 operator.add 会拼成一个长字符串(§4.5) |
消息字段一律用 add_messages |
别用 operator.add,会丢掉更新和删除能力,还会漏掉格式转换(§6.1) |
自定义 reducer 里加 existing or [] |
兜底成本很低,也方便单独单测(§3.3) |
自定义 reducer 不要原地修改 existing |
返回新对象,否则历史快照会被污染(§4.1) |
| 判断能不能当 reducer 用语义读一遍 | 「f(旧值, 新值) -> 新的旧值」通不通,签名检查靠不住(§4.4) |
| 删消息前先从 state 里取真实 id | 编造 id 会报错(§6.6) |
REMOVE_ALL_MESSAGES 放在列表最前面 |
放后面会把同批的新消息一起删掉(§6.7) |
| 关键业务信息放自定义槽位 | 消息会被裁剪,槽位不会(§7.3) |
| 并行节点之间不要有数据依赖 | 它们读到的是同一份快照(§5.3) |
| 给挂了 reducer 的字段传初始值要当心 | 那是「贡献一个值」,不是「设定成这个值」(§7.3) |
常见的坑按「会报错」和「不报错」分开列,不报错的那组才是真正耗时间的:
一、会报错(好排查)
| 现象 | 原因 | 处理 |
|---|---|---|
InvalidUpdateError: Can receive only one value per step |
并行分支写同一个键且无 reducer | 加 reducer(§5.1) |
ValueError: no signature found for builtin |
用了 max / min 等内置函数当 reducer |
包一层 lambda(§4.4) |
TypeError: can only concatenate list (not "int") to list |
类型标注和 reducer / 返回值不匹配 | 对齐标注和实际类型(§3.3) |
TypeError: 'int' object is not iterable |
用了 sum 当 reducer,它能过签名检查但语义不对 |
换成 operator.add(§4.4) |
Attempting to delete a message with an ID that doesn't exist |
RemoveMessage 的 id 是编的 |
用 state["messages"] 里的真实 id(§6.6) |
| reducer 里的异常直接把整张图打挂 | reducer 的异常不会被吞掉 | reducer 里别做可能失败的事(§4.1) |
二、不报错,但行为不对(静默失败)
| 现象 | 原因 | 处理 |
|---|---|---|
| 列表字段只剩最后一次写入 | 忘了加 reducer | Annotated[list[X], operator.add](§3.2) |
| 更新消息却变成了追加一条新的 | 没带 id | 从 state 取原 id(§6.4) |
| 某条历史消息内容被莫名换掉了 | 用错了 id,触发了「同 id 整条替换」 | 核对 id 来源(§6.3) |
摘要压缩后 messages 空了 |
REMOVE_ALL_MESSAGES 放在了新消息后面 |
调整顺序:先删后加(§6.7) |
messages 里混进了裸 str,下游 AttributeError |
用了 operator.add 而不是 add_messages |
换回 add_messages(§6.1) |
| 并行分支里读不到另一个分支刚写的值 | 并行看到的是分叉前快照 | 有依赖就改成串行(§5.3) |
| 改了个节点名,并行结果的顺序跟着变了 | 合并顺序按节点名排序 | 别依赖顺序,见三种解法(§5.4) |
传 {"turn": 10} 结果变成了 11 |
初始值也走 reducer,是累加不是赋值 | 去掉 reducer 或加重置哨兵(§7.3) |
status 变成了 '新建处理中已完成' |
给不该累积的字段挂了 operator.add |
去掉 reducer(§4.5) |
| 历史快照里的数据和当时不一致 | reducer 原地修改了 existing |
返回新对象(§4.1) |
三、其实是预期行为
| 现象 | 说明 |
|---|---|
初始输入的字符串 / 字典变成了 HumanMessage |
add_messages 作用在初始输入上(§3.3、§6.5) |
不传初始值也不报错,existing 是 [] / 0 / '' |
按类型标注取零值(§3.3) |
RemoveMessage 自己没出现在结果里 |
它是指令不是消息(§6.6) |
裁剪后 turn / topics 还在 |
槽位不受 messages 裁剪影响,这正是开槽位的目的(§7.3) |
口诀:
Reducer 属于状态,不属于节点。 节点只管返回自己算出来的那一小段,怎么合并是状态定义的事。
10. 练习 #
练习分两组:前五题是机制验证,跑一遍就知道有没有真看懂;后四题是能力建设,需要动手写代码。除最后一题外都不花钱。
机制验证
- 验证初始值也走 reducer:用 §3.3 的
spy_reducer跑一遍,数一数它被调用了几次,解释为什么传初始值时是 3 次、不传时是 2 次。 - 摸清零值规则:把一个字段的类型标注依次换成
list/int/str/dict,用探针 reducer 打印第一次调用时的existing,对照 §3.3 那张表。然后故意把标注和返回值写成不匹配的(标注list、返回5),看报错。 - 踩三次坑:分别试
Annotated[int, max](建图就报错)、Annotated[int, sum](建图通过、运行时报错)、Annotated[int, lambda o, n: max(o, n)](正常)。解释为什么第二个能通过检查。 - 观察并行顺序的真实规则:让三个并行节点各往同一个 list 追加一个元素,节点名故意取成
zebra/apple/mango,添加顺序也故意打乱,看合并顺序跟的是名字还是添加顺序。 - 验证
REMOVE_ALL_MESSAGES的顺序敏感:同一批更新里,分别把RemoveMessage(id=REMOVE_ALL_MESSAGES)放在新消息前面和后面,对比两种结果。
能力建设
- 改造上一章的图:给
ticket_graph.py的状态加一个log: Annotated[list[str], operator.add],让parse和assign各往里写一条记录,确认两条都在。 - 写一个 reducer:实现「只保留最近 3 条」的列表 reducer,挂到某个字段上验证。再进一步:写一个「能累加也能重置」的 reducer(传
None表示归零,§7.3 提过思路)。 - 消息更新:写一段代码,把
state["messages"]里最后一条AIMessage的内容改成大写。先故意不带 id 跑一次(会变成追加),再带上 id 跑一次(才是替换),对比两次的消息条数。 - 摘要压缩:用
REMOVE_ALL_MESSAGES+ 一条摘要消息,把 §8 的trim改成「压缩」而不是「丢弃」。注意验证turn和topics在压缩后仍然完整。
11. 本章小结 #
- 默认是覆盖,Reducer 改变合并规则。语法只有一处改动:
Annotated[类型, reducer函数],节点代码一个字不用动。这是本章反复出现的主题——合并规则属于状态,不属于节点。 - Reducer 的签名是
(existing, update),返回合并后的新值。它就是个普通函数,可以脱离图直接单测。 invoke传入的初始值也会走 reducer——这是本章最反直觉的一点。它解释了两件事:为什么字符串和字典能被自动转成消息对象,以及为什么有 reducer 的字段没法通过invoke直接「设定」值(传turn=10得到的是 11)。- 字段没值时
existing不是None,而是按类型标注取的零值:list→[]、int→0、str→''、dict→{}。这是TypedDict的类型标注在 LangGraph 里唯一有运行时意义的地方,写错会在运行时报TypeError。 max/min这类内置函数不能直接当 reducer(拿不到签名)。但签名检查只挡住「没签名」,挡不住「语义不对」——sum能通过检查却在运行时炸。判断标准是把它当f(旧值, 新值) -> 新的旧值读一遍通不通。- 并行分支是最需要 reducer 的地方:没 reducer 会
InvalidUpdateError(这是个「大声失败」,比静默猜测好),有 reducer 才能合流。多个分支的值是被依次喂进 reducer 的,所以 reducer 永远只需处理两个值。 - 并行节点读到的是分叉前的同一份快照,彼此写入互不可见——有数据依赖就该改成串行。
- 并行的合并顺序按节点名排序,和
add_node的添加顺序无关。 这意味着重命名一个节点就会改变合并顺序,而且不报错。要顺序保证就改串行、给元素带排序键、或者让各分支写各自的字段。 add_messages有五种行为:追加、按 id 整条替换、自动补 id、格式转换、RemoveMessage删除。其中「同 id 替换」是流式、消息修正、PII 脱敏的底层机制;它替换的是整条消息,连类型一起换。add_messages的四种出错方式里三种是静默的:漏 id(变追加)、用错 id(覆盖别人)、REMOVE_ALL_MESSAGES位置放错(新消息也被删)都不报错,只有「删不存在的 id」会报错。REMOVE_ALL_MESSAGES+ 一条新消息 = 摘要压缩,第 11 章SummarizationMiddleware就是这么实现的。顺序必须是先删后加,反了会连摘要一起删掉。MessagesState只是一行预置定义,不提供任何额外能力(不会自动存历史、不会自动裁剪),可以直接继承来扩展自己的槽位。- 消息会被裁剪,自定义槽位不会——关键业务信息(订单号、客户名、累计值)要放槽位。
- 本章产出:可追加 messages 的状态图——一张图里同时跑了
add_messages、operator.add和自定义的merge_unique三种 reducer,并用RemoveMessage做了历史裁剪。三个节点里只有一个调模型,另外两个是纯确定性逻辑。
到这里,状态(本章)、节点、边(上一章)三个概念就齐了。下一章把注意力转到运行层面:compile() 到底做了什么,invoke / stream / get_state 各自适合什么场景。