add_messages:工业级的消息 reducer
昨天你知道 reducer 就是"字段的合并规矩"。今天拆开 LangGraph 里最重要、也最复杂的一个 reducer——add_messages。它远不止"两个列表拼一起":要按 ID 去重、按 ID 更新(改写已有消息)、按 ID 删除(RemoveMessage)、一键清空(REMOVE_ALL),还要把各种松散写法归一成标准消息对象。源码在 graph/message.py。
messages 想成一个共享的聊天记录本。add_messages 就是这本记录的"管理员"。它做四件事:来新消息就贴到最后;来的消息 ID 和本子里某条一样,就把那条原地改写(比如流式补全、纠错);来一张写着"删除 ID=x"的便签(RemoveMessage),就把 x 那条撕掉;来一张"全部清空"的便签(REMOVE_ALL),就把清空之前的都作废、只留之后的。今天就是逐行看这个管理员怎么干活。痛点:消息列表为什么不能简单 a+b
messages 如果只用 operator.add(简单拼接)会出三个问题:① 流式生成时,模型先吐半句、再吐整句,两次都写进来,你会得到重复的半句+整句;② 想修正某条历史消息(同一条重发),拼接只会又加一条,改不掉旧的;③ 想删除某条消息(撤回、裁剪上下文),拼接根本做不到。还有:节点可能返回 ("user","hi") 元组、{"role":...} 字典、或裸字符串——五花八门,得先统一。id 当成主键来管列表:ID 没见过就追加,ID 见过就更新,遇到删除标记就移除。这正是数据库"upsert(有则更新、无则插入)+ delete"的思路,搬到了内存里的消息列表上。所以它能同时解决去重、更新、删除三件事。官方 docstring 用一个例子点明了核心行为(graph/message.py:101):
# graph/message.py:101 —— docstring 里的"覆盖同 ID"示例
msgs1 = [HumanMessage(content="Hello", id="1")]
msgs2 = [HumanMessage(content="Hello again", id="1")] # 同一个 id="1"
add_messages(msgs1, msgs2)
# 结果只有一条:[HumanMessage(content='Hello again', id='1')] ← 被更新,不是追加
add_messages 都围绕这条主线在打转。柯里化包装器:为什么能写 add_messages(format=...)
你可能见过两种写法:直接 Annotated[list, add_messages],也见过 Annotated[list, add_messages(format="langchain-openai")](带括号调用)。这个"既能当函数、又能先配置再当函数"的魔法,靠一个装饰器 _add_messages_wrapper(graph/message.py:41):
# graph/message.py:41
def _add_messages_wrapper(func):
def _add_messages(left=None, right=None, **kwargs):
if left is not None and right is not None:
return func(left, right, **kwargs) # ① 正常调用:真的去合并
elif left is not None or right is not None:
msg = "Must specify non-null arguments for both 'left' and 'right'. ..."
raise ValueError(msg) # ② 只给一半 → 报错
else:
return partial(func, **kwargs) # ③ 一个都没给 → 返回"预配置好的函数"
_add_messages.__doc__ = func.__doc__
return cast(..., _add_messages)
@_add_messages_wrapper
def add_messages(left, right, *, format=None):
...
① left 和 right 都有说明是当 reducer 真的在合并两个列表,直接调真身 func。③ 都没有,只传 kwargs说明你在配置它,比如 add_messages(format="langchain-openai")。此时返回 partial(func, format=...)——一个"已经把 format 记住了"的新函数,等以后被当 reducer 调用。这就是"柯里化/预填参数"。② 只给了一个边界:只传 left 或只传 right 是没意义的(reducer 必须成对合并),直接报错,防止误用。__doc__ = func.__doc__小细节:把原函数的文档搬到包装函数上,保证 help(add_messages) 仍显示正确说明。add_messages(无参)和 make_add_messages(format=...)(工厂)。源码把它们合成一个名字,靠"参数给全就执行、参数没给全就返回预配置函数"来区分。好处是用户只需记住一个 add_messages,无论要不要配置,写在 Annotated 里都很自然。代价是包装逻辑本身有点绕——但对使用者是"零心智负担"。这是"把复杂留给库、把简单留给用户"的典型权衡。第一步:归一化——转 list、转消息对象、补 ID
进入真身后,第一件事是把 left/right 各种松散写法统一成标准消息对象(graph/message.py:187):
# graph/message.py:187
remove_all_idx = None
# ① 单条也裹成 list
if not isinstance(left, list):
left = [left]
if not isinstance(right, list):
right = [right]
# ② 各种写法统一成 BaseMessage(元组/字典/字符串 → 消息对象)
left = [message_chunk_to_message(m) for m in convert_to_messages(left)]
right = [message_chunk_to_message(m) for m in convert_to_messages(right)]
# ③ 补齐缺失的 id(用随机 uuid)
for m in left:
if m.id is None:
m.id = str(uuid.uuid4())
for idx, m in enumerate(right):
if m.id is None:
m.id = str(uuid.uuid4())
if isinstance(m, RemoveMessage) and m.id == REMOVE_ALL_MESSAGES: # 顺便记下"清空"标记的位置
remove_all_idx = idx
① 裹成 list允许你只返回单条消息(不用手动包 list),这里统一补成列表,后面逻辑就不用区分单条/多条。convert_to_messagesLangChain 提供的转换:把 ("user","hi")、{"role":"user","content":"hi"}、裸字符串 "hi" 全部转成规范的 HumanMessage/AIMessage 等对象。message_chunk_to_message流式场景里模型吐的是 Chunk(片段类型),这里把片段类型转成"完整消息类型",保证存进 State 的都是稳定对象。补 id = uuid4()★关键:ID 是这个 reducer 的"主键",没 ID 就没法去重/更新。所以缺 ID 的一律补随机 UUID。这也意味着:你自己不给 ID,就永远无法后续更新或删除这条消息(因为 ID 随机,你不知道)。记 remove_all_idx顺路扫描 right 里有没有"全部清空"标记(REMOVE_ALL_MESSAGES),记下它的位置,L05 用。return {"messages": "你好"} 或 [("assistant","hi")],进了 reducer 都会被"翻译"成标准消息。把混乱挡在门口、内部只处理规范对象——这是健壮代码的通用套路。核心:按 ID 去重 / 更新 / 删除
归一化后,进入真正的合并循环(graph/message.py:216)。这是 add_messages 的心脏:
# graph/message.py:216
merged = left.copy()
merged_by_id = {m.id: i for i, m in enumerate(merged)} # id → 在 merged 里的下标
ids_to_remove = set()
for m in right:
if (existing_idx := merged_by_id.get(m.id)) is not None: # 这个 id 已经存在
if isinstance(m, RemoveMessage):
ids_to_remove.add(m.id) # 是删除标记 → 记下待删
else:
ids_to_remove.discard(m.id) # 之前可能标了删,现在又更新 → 取消删
merged[existing_idx] = m # ★原地更新(覆盖同 ID)
else: # 这个 id 没见过
if isinstance(m, RemoveMessage):
raise ValueError(
f"Attempting to delete a message with an ID that doesn't exist ('{m.id}')"
) # 删一个不存在的 → 报错
merged_by_id[m.id] = len(merged)
merged.append(m) # ★全新消息 → 追加
merged = [m for m in merged if m.id not in ids_to_remove] # 最后统一剔除待删的
merged_by_id★核心数据结构:一个 {id: 下标} 的字典。有了它,判断"这个 ID 存不存在"是 O(1),不用每次遍历整个列表。existing_idx := ...get(m.id)海象运算符:一边取值一边判断。get 返回 None 表示 ID 没见过。已存在 + RemoveMessage记进 ids_to_remove(先不真删,最后统一删——避免边遍历边改下标出乱子)。已存在 + 普通消息★merged[existing_idx] = m 原地替换——这就是"同 ID 更新"。同时 discard 掉之前可能的删除标记(后写的更新赢)。不存在 + RemoveMessage边界:想删一个根本不存在的 ID → 直接 raise ValueError,不静默忽略(L07 详谈这个抉择)。不存在 + 普通消息追加到末尾,并在字典里登记它的下标。最后一行列表推导把标记为删除的 ID 一次性过滤掉,得到最终列表。left=[Human("Hi",id="1")];right=[AI("流式半句",id="2"), AI("完整回答",id="2"), Human("Hi改",id="1")]。过程:id=2 首次出现→追加半句;id=2 又来→原地覆盖成完整回答;id=1 已存在→覆盖成"Hi改"。
最终:
[Human("Hi改",id="1"), AI("完整回答",id="2")]——流式重复被自动合并,历史消息被成功修正。
RemoveMessage 与 REMOVE_ALL:"删除"的两种粒度
L04 已看到按 ID 删除单条(RemoveMessage(id="x"))。还有一种"一键清空",靠常量 REMOVE_ALL_MESSAGES(graph/message.py:38)实现,处理逻辑在归一化之后、合并循环之前(graph/message.py:212):
# graph/message.py:38
REMOVE_ALL_MESSAGES = "__remove_all__"
# graph/message.py:212 —— 如果 right 里出现了"清空"标记
if remove_all_idx is not None:
return right[remove_all_idx + 1 :] # ★丢掉 left 全部 + 清空标记之前的,只留其后的
RemoveMessage(id=REMOVE_ALL_MESSAGES)一条 id 特意等于 "__remove_all__" 的删除消息,代表"把历史全清了"。return right[remove_all_idx+1:]★极简且巧妙:一旦发现清空标记,整个 left(历史)作废,连 right 里清空标记之前的也作废,只返回清空标记之后的消息。相当于"从这里重开"。提前 return清空是"核选项",直接短路返回,根本不进 L04 的逐条合并循环——因为没必要再去 upsert 一堆马上要被清掉的东西。left=[m1,m2,m3],right=[RemoveMessage(id="__remove_all__"), AI("新对话开始",id="9")] → 返回 [AI("新对话开始",id="9")]。历史全清空,只留清空后的新消息。常用于"重置会话"或"上下文裁剪重开"。| 你想干什么 | 往 messages 写什么 | add_messages 的动作 |
|---|---|---|
| 加一条新消息 | 新消息(新 ID) | 追加到末尾 |
| 修正/流式补全某条 | 同 ID 的新消息 | 原地覆盖那条 |
| 删除某条 | RemoveMessage(id="x") | 移除 ID=x 那条 |
| 清空历史重开 | RemoveMessage(id=REMOVE_ALL_MESSAGES) | 丢弃全部历史 |
MessagesState 与实验性的 delta reducer
99% 的对话图不用自己写 Annotated[list, add_messages]——直接继承现成的 MessagesState(graph/message.py:372):
# graph/message.py:372
class MessagesState(TypedDict):
messages: Annotated[list[AnyMessage], add_messages]
class State(MessagesState): ... 再加自己的字段,就白得一个配好 add_messages 的 messages 字段。这是"约定优于配置"——把最常见用法固化成一行可复用的类。源码里还有个实验性的批量版 _messages_delta_reducer(graph/message.py:247),它体现了一个更深的设计约束——批处理不变性(batching-invariant):
# graph/message.py:247(docstring 节选)
"""**Experimental.** Batch reducer for use with `DeltaChannel`.
This reducer is batching-invariant, as required by `DeltaChannel`:
reducer(reducer(state, xs), ys) == reducer(state, xs + ys)
"""
# graph/message.py:292 —— 用 dict 做索引,一趟扫描完成去重/更新/删除(tombstone)
index = {m.id: i for i, m in enumerate(state_msgs) if m.id is not None}
result = list(state_msgs)
for msg in msgs:
mid = msg.id
if mid is None:
result.append(msg)
elif isinstance(msg, RemoveMessage):
if mid in index:
result[index[mid]] = None # 先置 None(墓碑),最后统一过滤
del index[mid]
elif mid in index:
result[index[mid]] = msg # 更新
else:
index[mid] = len(result); result.append(msg) # 追加
return [m for m in result if m is not None]
reducer(reducer(state, xs), ys) == reducer(state, xs + ys) 的意思是:无论新写入是"分两批来"还是"合成一批来",最终结果必须一样。为什么关键?因为在持久化/流式场景里,写入可能被拆成任意批次落盘或重放(后面阶段讲 checkpoint 会看到)。如果 reducer 对"分批 vs 合批"敏感,同样的操作在不同批次下会得出不同状态——那持久化就不可靠了。_messages_delta_reducer 是为 DeltaChannel 显式设计、用"墓碑(先置 None 最后过滤)"技巧一趟扫完,把这个性质当作硬约束写进了文档。REMOVE_ALL、删不存在 ID 报错、缺 ID 补 UUID、Chunk 转换这些它都不处理。它是为高性能批量场景做的裁剪版,日常仍用 add_messages。边界抉择 + 今日小结
RemoveMessage 指向一个不存在的 ID 时,源码选择 raise ValueError。这是一个有意的严格选择。静默跳过看似"更宽容",但会掩盖 bug——你以为删掉了某条消息,实际因为 ID 写错根本没删,后续逻辑就带着错误上下文继续跑,问题在很远的地方才爆发、极难排查。删除是"你明确认为它存在"的操作,若不存在,说明你的假设已经错了,越早报错越好。对比 L03 的"缺 ID 补 UUID"(宽容)——追加一条没 ID 的消息是常见且无害的,所以宽容;删一个不存在的东西是逻辑矛盾,所以严格。同一函数里对不同情况松紧有别,都是有理由的。👶 小白:我在节点里 return {"messages": [...]},这个列表是"我要替换全部"还是"我要追加"?
👨🏫 老师:是追加(准确说是 upsert)。你返回的列表会作为 right 传给 add_messages,跟已有的 left 按 ID 合并。想真正替换全部,得显式发 REMOVE_ALL 标记再写新的。这是新手最常踩的直觉误区——"返回列表 ≠ 覆盖列表"。
🧠 今天你应该能回答
add_messages用什么当"主键"来去重/更新/删除?(消息的id)- 为什么它比
operator.add复杂?(要处理流式重复、原地更新、删除、格式归一化) - 柯里化包装器怎么让它既能直接用又能
add_messages(format=...)?(参数给全就执行,没给全返回partial) - 同 ID 的新消息会发生什么?(原地覆盖旧的)
- 删一个不存在的 ID 会怎样?(
raise ValueError,故意严格) - REMOVE_ALL 和 RemoveMessage(id) 有何区别?(清空全部 vs 删除单条)
- "批处理不变性"是什么、为什么持久化需要它?
✋ 10 分钟动手
# 1. 读 add_messages 真身(归一化 + 合并循环)
sed -n '187,244p' libs/langgraph/langgraph/graph/message.py
# 2. 读柯里化包装器
sed -n '41,66p' libs/langgraph/langgraph/graph/message.py
# 3. 亲手体验去重/更新/删除
python -c "
from langchain_core.messages import HumanMessage, RemoveMessage
from langgraph.graph.message import add_messages
a=[HumanMessage('Hi', id='1')]
b=[HumanMessage('Hi改', id='1'), RemoveMessage(id='1')]
print(add_messages(a,b)) # 先更新再删除 → 空列表
"
add_messages 底层被包成了 BinaryOperatorAggregate 通道(还记得 D07 分支③吗)。明天钻进 channels/binop.py,看这个"累加通道"内部——空值 MISSING 哨兵、operator.add 怎么累加、以及 Overwrite 强制覆盖机制。