Day 08 / 共 60 天 · 阶段2 状态与数据流

add_messages:工业级的消息 reducer

昨天你知道 reducer 就是"字段的合并规矩"。今天拆开 LangGraph 里最重要、也最复杂的一个 reducer——add_messages。它远不止"两个列表拼一起":要按 ID 去重、按 ID 更新(改写已有消息)、按 ID 删除(RemoveMessage)、一键清空(REMOVE_ALL),还要把各种松散写法归一成标准消息对象。源码在 graph/message.py

📍 你在 60 天里的位置(阶段2:状态与数据流 · 共 6 天)
阶段1 入门 D07 reducers D08 add_messages D09 累加通道 D10 并行写 D11 输入输出 D12 Pydantic 阶段3 控制流
💡 先用一个类比兜住今天messages 想成一个共享的聊天记录本add_messages 就是这本记录的"管理员"。它做四件事:来新消息就贴到最后;来的消息 ID 和本子里某条一样,就把那条原地改写(比如流式补全、纠错);来一张写着"删除 ID=x"的便签(RemoveMessage),就把 x 那条撕掉;来一张"全部清空"的便签(REMOVE_ALL),就把清空之前的都作废、只留之后的。今天就是逐行看这个管理员怎么干活。
L01

痛点:消息列表为什么不能简单 a+b

🤔 痛点聊天 Agent 的 messages 如果只用 operator.add(简单拼接)会出三个问题:① 流式生成时,模型先吐半句、再吐整句,两次都写进来,你会得到重复的半句+整句;② 想修正某条历史消息(同一条重发),拼接只会又加一条,改不掉旧的;③ 想删除某条消息(撤回、裁剪上下文),拼接根本做不到。还有:节点可能返回 ("user","hi") 元组、{"role":...} 字典、或裸字符串——五花八门,得先统一。
💡 本质:add_messages = 带"主键(ID)"的合并它把每条消息的 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')]  ← 被更新,不是追加
大白话ID 相同 = "还是那条消息的新版本",所以覆盖;ID 不同 = "一条全新消息",所以追加。整个 add_messages 都围绕这条主线在打转。
L02

柯里化包装器:为什么能写 add_messages(format=...)

你可能见过两种写法:直接 Annotated[list, add_messages],也见过 Annotated[list, add_messages(format="langchain-openai")](带括号调用)。这个"既能当函数、又能先配置再当函数"的魔法,靠一个装饰器 _add_messages_wrappergraph/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) 仍显示正确说明。
💡 设计取舍①:为什么要搞这个"两用"包装,不直接分成两个函数? 朴素做法是给用户两个 API:add_messages(无参)和 make_add_messages(format=...)(工厂)。源码把它们合成一个名字,靠"参数给全就执行、参数没给全就返回预配置函数"来区分。好处是用户只需记住一个 add_messages,无论要不要配置,写在 Annotated 里都很自然。代价是包装逻辑本身有点绕——但对使用者是"零心智负担"。这是"把复杂留给库、把简单留给用户"的典型权衡。
L03

第一步:归一化——转 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 都会被"翻译"成标准消息。把混乱挡在门口、内部只处理规范对象——这是健壮代码的通用套路。
L04

核心:按 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 一次性过滤掉,得到最终列表。
控制流:right 里每条消息走哪个分支 取 right 里一条 m m.id 在 merged_by_id 里吗? 存在 RemoveMessage → 记入待删集合 普通消息 → 原地覆盖更新 不存在 RemoveMessage → raise ValueError 普通消息 → 追加到末尾 循环结束后:把"待删集合"里的 ID 从 merged 里过滤掉
图注:四个分支 = 存在/不存在 × 删除/普通。这就是 upsert+delete 的完整真值表。
📝 真实值走一遍 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")]——流式重复被自动合并,历史消息被成功修正。
L05

RemoveMessage 与 REMOVE_ALL:"删除"的两种粒度

L04 已看到按 ID 删除单条RemoveMessage(id="x"))。还有一种"一键清空",靠常量 REMOVE_ALL_MESSAGESgraph/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)丢弃全部历史
L06

MessagesState 与实验性的 delta reducer

99% 的对话图不用自己写 Annotated[list, add_messages]——直接继承现成的 MessagesStategraph/message.py:372):

# graph/message.py:372
class MessagesState(TypedDict):
    messages: Annotated[list[AnyMessage], add_messages]
你只要 class State(MessagesState): ... 再加自己的字段,就白得一个配好 add_messagesmessages 字段。这是"约定优于配置"——把最常见用法固化成一行可复用的类。

源码里还有个实验性的批量版 _messages_delta_reducergraph/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 最后过滤)"技巧一趟扫完,把这个性质当作硬约束写进了文档。
注意 docstring 明说它不是 add_messages 的完全等价REMOVE_ALL、删不存在 ID 报错、缺 ID 补 UUID、Chunk 转换这些它都不处理。它是为高性能批量场景做的裁剪版,日常仍用 add_messages
L07

边界抉择 + 今日小结

⚠️ 边界:为什么"删除不存在的 ID"要报错,而不是静默跳过? 回看 L04: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))   # 先更新再删除 → 空列表
"
明日预告 · Day 09add_messages 底层被包成了 BinaryOperatorAggregate 通道(还记得 D07 分支③吗)。明天钻进 channels/binop.py,看这个"累加通道"内部——空值 MISSING 哨兵、operator.add 怎么累加、以及 Overwrite 强制覆盖机制。
← Day 07 reducers 入门 Day 09 · BinaryOperatorAggregate →