14333 字
72 分钟

LangGraph 源码深潜:Reducer 与并行状态合并机制解剖

核心问题: 为什么 LangGraph 的 state 不是普通字典,而是一个由 State schema → Channel → Reducer 共同决定更新语义的状态模型?

源码主线: StateGraph(TravelState) → _add_schema() → _get_channels() → LastValue / BinaryOperatorAggregate / add_messages → Pregel super-step writes → channel.update(values) → merged state

前置文章: 第 6 篇《StateGraph 源码解剖》、第 7 篇《Conditional Edge 与 Router 源码解剖》

依赖基线: langgraph==1.2.7langchain-corelanggraph 依赖解析

源码基线: langchain-ai/langgraph GitHub release/tag 1.2.7;本文所有源码链接优先使用 1.2.7 tag

阅读边界: 本文只讲 State schema 中字段如何变成 channel、reducer 如何合并并行写入、为什么无 reducer 会冲突、messages 为什么需要特殊 reducer;不展开 checkpoint 持久化、interrupt 恢复和完整 Pregel 调度器实现。


0. 系列位置、前置知识与版本基线#

前两篇 LangGraph 文章已经完成了两层源码认知:

第 6 篇 StateGraph:
StateGraph 是 builder,compile 后得到 CompiledStateGraph;
node 输入完整 state,返回 partial state update;
edge 决定下一步节点。
第 7 篇 Conditional Edge:
add_conditional_edges 注册 path function;
节点执行后读取 state 并返回下一个节点;
返回 END 终止,返回多个目标可形成并行分支。

这一篇补上 LangGraph 状态模型中最容易被忽略、但最关键的机制:

多个节点同时写同一个 state key 时,LangGraph 到底如何决定最终 state?

如果把 state 当成普通 Python 字典,就很难解释下面的问题:

为什么某些字段可以被多个并行节点同时写入?
为什么某些字段一旦被多个节点同时写入就报错?
为什么 Annotated[list[str], add] 可以追加?
为什么 messages 不建议直接用 operator.add?
为什么 node 只返回 partial update,却能稳定合并回全局 state?

本篇要建立的新认知是:

StateGraph 的 state 并不是直接 dict merge。
state schema
每个字段被解析成一个 channel
每个 channel 决定该字段的更新规则
默认 LastValue:覆盖式更新,只允许每个 super-step 一个写入
Annotated reducer:BinaryOperatorAggregate,允许多个写入聚合
messages:add_messages,按消息 ID 追加或覆盖

本篇只解决:

  • State schema 如何解析 Annotated
  • Reducer 在构建期如何变成 channel;
  • node 返回的 partial update 在运行时如何被写入 channel;
  • 多节点写同一个 key 时,LangGraph 如何判断合并还是冲突;
  • 覆盖式更新与追加式更新的源码边界;
  • messages 为什么需要 add_messages 这种特殊 reducer。

本篇不展开:

  • checkpoint 如何保存 channel values;
  • interrupt 如何恢复到暂停点;
  • subgraph 如何跨图合并 state;
  • Pregel 调度器所有细节;
  • LangGraph Server 中的远程状态更新协议。

1. 本篇问题、学习目标与能力边界#

1.1 核心问题#

LangGraph 如何把 TypedDict / Annotated 定义的状态 schema,转换成有合并规则的 channel 系统,并在并行节点写入时正确合并 partial state update?

1.2 学习目标#

完成本篇后,读者必须能够:

  1. 解释为什么 StateGraph 的 state 不是普通字典,而是由 schema 与 channel 决定更新语义。
  2. 解释 Annotated[list[str], add] 如何在构建期被解析为 BinaryOperatorAggregate
  3. 解释默认字段为什么使用 LastValue,以及为什么默认字段不允许同一 super-step 多写入。
  4. 解释 reducer 在什么时候被调用:不是 node 内部,而是在 Pregel super-step 写入应用阶段。
  5. 区分覆盖式更新、追加式更新、消息合并、强制覆盖 Overwrite 的边界。
  6. 能为旅行规划助手设计哪些字段覆盖、哪些字段追加、哪些字段必须使用 add_messages

1.3 能力边界#

能力本篇是否覆盖说明
State schema 解析重点讲 _get_channels()_get_channel()Annotated 解析
reducer channel 构造重点讲 BinaryOperatorAggregate 与默认 LastValue
并行写入合并重点讲 super-step 结束后的 channel update
messages 合并add_messages 为什么比 operator.add 更适合对话历史
checkpoint 存储下一篇 durable execution 继续讲
interrupt / resumeHITL 篇继续讲
Pregel 完整调度算法部分只讲与 writes / channel update 有关的主链

2. 核心概念与最小心智模型#

2.1 一句话定义#

Reducer 是 LangGraph 中每个 state key 的“更新合并函数”,负责把一个或多个节点返回的 partial update 合并到已有 state;它不负责调度节点,也不负责判断下一条边。

更准确地说,在 LangGraph 源码里,reducer 通常不直接裸露在运行时主链,而是被封装进 channel:

没有 reducer 的字段
→ LastValue channel
→ 覆盖式更新
→ 同一 super-step 最多一个写入
Annotated[T, reducer] 字段
→ BinaryOperatorAggregate channel
→ reducer(left, right) 聚合
→ 同一 super-step 可多个写入
messages 字段
→ add_messages reducer
→ append-only + same-id overwrite + message deserialization

2.2 最小心智模型#

State schema
解析字段类型与 Annotated metadata
每个字段生成一个 channel
节点执行返回 partial update
Pregel 收集同一 super-step 的全部 writes
按 channel 分组
channel.update(values)
得到下一步 state

2.3 核心术语#

术语源码对象语义不要误解为
State schemaTypedDict / BaseModel / dataclass声明 state key 与更新语义运行时 state 本身
State keyschema 字段一个可读写的状态字段普通 dict key
ChannelBaseChannel 子类某个 key 的运行时存储与更新规则消息队列
默认 channelLastValue没有 reducer 时的覆盖更新任意次数写入都可覆盖
reducer channelBinaryOperatorAggregate用二元函数聚合多个写入只在 Python dict merge 时调用
partial updatedict[str, Any]node 返回的局部状态更新完整 state
super-stepPregel step一轮并行任务执行与写入应用边界单个 node 调用
add_messageslanggraph.graph.message.add_messagesmessage list 专用 reducer简单 list add
Overwritelanggraph.types.Overwrite强制绕过 reducer 覆盖字段默认更新方式

2.4 与相邻抽象的边界#

对象负责什么不负责什么与 Reducer 的关系
StateGraph读取 schema、注册节点边、编译图执行节点主循环构建期把字段解析成 channel
CompiledStateGraph可执行 Runnable 图定义 reducer 函数运行时通过 Pregel 使用 channel
Pregel调度 super-step、执行任务、应用 writes决定业务字段语义在写入阶段调用 channel.update
BaseChannel存储字段值与应用更新决定节点执行顺序reducer 语义的封装载体
Conditional Edge根据 state 决定下一节点合并 state update读取 reducer 合并后的 state 做路由

3. 完整执行链路#

3.1 高层链路#

观察目标:下面的代码不是为了演示 API,而是为了定位 reducer 如何从 schema 进入运行时。

from typing import Annotated
from operator import add
from typing_extensions import TypedDict
from langgraph.graph import StateGraph, START, END
class TravelState(TypedDict):
user_request: str
research_notes: Annotated[list[str], add]
def search_food(state: TravelState):
return {"research_notes": ["推荐美食:寿司、拉面、居酒屋"]}
def search_attractions(state: TravelState):
return {"research_notes": ["推荐景点:浅草寺、涩谷、上野公园"]}
def merge_plan(state: TravelState):
return {
"research_notes": state["research_notes"] + ["生成综合旅行计划"]
}
builder = StateGraph(TravelState)
builder.add_node("search_food", search_food)
builder.add_node("search_attractions", search_attractions)
builder.add_node("merge_plan", merge_plan)
builder.add_edge(START, "search_food")
builder.add_edge(START, "search_attractions")
# 这里使用等待边,表示 merge_plan 等 search_food 与 search_attractions 都完成后再运行
builder.add_edge(["search_food", "search_attractions"], "merge_plan")
builder.add_edge("merge_plan", END)
graph = builder.compile()
result = graph.invoke({
"user_request": "我想去东京玩 5 天",
"research_notes": [],
})

高层执行链路:

TravelState schema
StateGraph.__init__
_get_channels(TravelState)
research_notes → BinaryOperatorAggregate(list[str], operator.add)
compile 生成 CompiledStateGraph / Pregel
START 同时触发 search_food 与 search_attractions
两个节点都返回 {research_notes: [...]} partial update
Pregel 在 super-step 边界收集两个 writes
research_notes channel.update([food_notes, attraction_notes])
operator.add 合并 list
merge_plan 读取合并后的 research_notes
END

3.2 对象流转#

阶段输入类型核心函数输出类型状态变化
schema 解析type[TravelState]_get_channels()dict[str, BaseChannel]research_notes 变成 reducer channel
构建图node / edge 注册StateGraph.add_node() / add_edge()builder 内部结构记录并行起点与等待边
编译图builderStateGraph.compile()CompiledStateGraphchannel 与节点被挂入 Pregel
节点执行TravelStatesearch_food() / search_attractions()dict[str, list[str]]产生 partial update
写入应用pending writeschannel.update(values)bool根据 reducer 合并 state
下游读取merged statemerge_plan()partial update读取合并后的研究笔记

3.3 时序链路#

Caller
│ graph.invoke(initial_state)
CompiledStateGraph / Pregel
│ 初始化 channels:user_request, research_notes
START
│ 激活 search_food 与 search_attractions
Parallel nodes
│ 两个节点都写 research_notes
Pregel write application
│ research_notes channel 执行 reducer
merge_plan
│ 读取合并后的 research_notes
END

3.4 正常结束条件#

一次执行正常完成的条件是:

所有已激活节点执行完毕
所有 pending writes 被 channel 正确应用
没有新的节点被边激活
到达 END 或任务队列为空
输出 schema 对应的 state 被读取并返回

在本例中,最终 state 中:

{
"user_request": "我想去东京玩 5 天",
"research_notes": [
"推荐美食:寿司、拉面、居酒屋",
"推荐景点:浅草寺、涩谷、上野公园",
"生成综合旅行计划",
],
}

关键不是结果长什么样,而是:

两个并行节点写同一个 key 没有互相覆盖,
是因为 research_notes 字段被 reducer channel 接管了。

4. 源码地图、关键文件与阅读顺序#

4.1 核心目录#

langgraph/
├── graph/
│ ├── state.py # StateGraph、CompiledStateGraph、schema→channel 解析
│ ├── message.py # add_messages、MessagesState
│ └── graph.py # 基础 Graph 抽象
├── channels/
│ ├── base.py # BaseChannel 协议
│ ├── last_value.py # 默认覆盖式 channel
│ ├── binop.py # BinaryOperatorAggregate reducer channel
│ ├── topic.py # 多值 topic channel
│ └── ephemeral_value.py
└── pregel/
├── main.py # Pregel 可执行图主类
├── _algo.py # writes 应用、channel 更新、任务准备等算法
├── _write.py # ChannelWrite
└── _read.py # ChannelRead

4.2 关键文件#

优先级文件核心对象阅读目的
1langgraph/graph/state.pyStateGraphCompiledStateGraph看 schema 如何变成 channels
2langgraph/channels/last_value.pyLastValue看默认覆盖式更新与并发冲突
3langgraph/channels/binop.pyBinaryOperatorAggregate看 reducer 如何应用多个值
4langgraph/graph/message.pyadd_messagesMessagesState看 messages 为什么特殊
5langgraph/pregel/_algo.pywrites 应用逻辑看 reducer 在运行时什么时候被调用
6langgraph/pregel/_write.pyChannelWrite看 node 输出如何转成 channel writes

4.3 推荐阅读顺序#

1. StateGraph.__init__
2. StateGraph._add_schema
3. _get_channels / _get_channel
4. LastValue.update
5. BinaryOperatorAggregate.update
6. StateGraph.compile / attach_node
7. Pregel writes application
8. add_messages / MessagesState

4.4 不建议的阅读顺序#

不建议从 pregel/main.py 开始读。原因是 Pregel 主循环承担:

调度
checkpoint
interrupt
streaming
cache
retry
任务准备
写入应用

如果直接从 Pregel 主类读,会把 reducer、checkpoint、stream、interrupt 全部混在一起。更好的路径是:

先在 StateGraph 构建期看字段如何变成 channel,
再到 channels 看 update 行为,
最后回到 Pregel 看运行时什么时候调用 update。

5. 对象模型、继承关系与协议边界#

5.1 核心对象关系#

State schema
_get_channels(schema)
BaseChannel
├── LastValue
├── BinaryOperatorAggregate
├── Topic
└── EphemeralValue
Pregel channel update
Merged state

对于本文最重要的两个 channel:

LastValue
默认 channel
语义:一个 step 内最多一个更新;后续 step 覆盖旧值
BinaryOperatorAggregate
reducer channel
语义:使用 operator(left, right) 聚合同一 key 的多个更新

5.2 对象职责#

对象生命周期输入输出核心职责
StateGraph构建期schema / nodes / edgesbuilder解析 schema、保存节点边
CompiledStateGraph编译后graph inputoutput stateRunnable 化的可执行状态图
BaseChannel构建期创建,运行时复制updateschannel value定义某个 key 的更新规则
LastValue运行时 channelSequence[Value]last value默认覆盖更新,防并发多写
BinaryOperatorAggregate运行时 channelSequence[Value]reduced valuereducer 聚合多写入
add_messagesreducer functionold messages, new messagesmerged messages消息追加、同 ID 覆盖、反序列化
Pregel运行时tasks / writesnext state在 super-step 边界应用 channel updates

5.3 协议边界#

State schema 协议
负责:声明 state key、字段类型、reducer metadata。
不负责:执行节点、调用 reducer。
Channel 协议
负责:接收某个 key 的一批更新,并决定如何合并。
不负责:知道这些更新来自哪个业务节点。
Pregel 执行协议
负责:调度节点、收集 writes、在 step 边界调用 channel.update。
不负责:替业务字段猜测合并语义。
Reducer 函数
负责:把旧值与新值合并成一个值。
不负责:校验整个业务 state 是否合理。

5.4 稳定接口与内部实现#

类型对象文章中的使用原则
公共 APIStateGraphAnnotated[..., reducer]add_messages可用于工程代码
公共 APIOverwrite可在需要绕过 reducer 时谨慎使用
扩展接口自定义 reducer 函数必须满足二元函数契约,建议纯函数
内部实现_get_channels()_get_channel()只用于源码理解,不建议业务代码依赖
内部实现LastValueBinaryOperatorAggregate可读源码理解语义,一般不直接实例化
内部实现Pregel writes 应用函数只用于解释运行时,不作为业务扩展点

6. 源码阅读策略与证据标准#

6.1 本篇阅读策略#

先看 schema:字段怎么声明?
再看构建期:字段如何转 channel?
再看 channel:update(values) 怎么合并?
再看运行时:什么时候把多个 values 传给 update?
最后看边界:并发多写、Overwrite、messages、异常。

6.2 证据等级#

标记含义写作要求
源码事实可以由 langgraph==1.2.7 源码直接证明1.2.7 tag 源码链接
官方契约官方文档或 API Reference 明确承诺附官方文档链接
简化伪代码对真实控制流的压缩表达明确标注“伪代码”
作者推断根据调用链得出的设计理解明确使用“从调用关系可以推断”
工程建议面向项目实践的建议说明适用条件

6.3 本篇证据清单#

结论证据类型文件或文档定位
StateGraph 节点签名是 State -> Partial<State>,state key 可用 reducer 聚合官方契约Graph API / StateGraph ReferenceState / Reducers
没有 reducer 时默认覆盖更新官方契约Graph APIReducers / Default reducer
Annotated[list, add] 让 list 追加官方契约Graph APIReducers Example B
并行多节点写同一无 reducer key 会报 INVALID_CONCURRENT_GRAPH_UPDATE官方契约Error docsINVALID_CONCURRENT_GRAPH_UPDATE
LastValue 一个 step 最多接收一个值源码事实channels/last_value.pyLastValue.update()
BinaryOperatorAggregate 逐个应用 operator源码事实channels/binop.pyBinaryOperatorAggregate.update()
add_messages 通过 ID 合并消息官方契约 / 源码事实Graph API / graph/message.pyadd_messages

7. 构建期源码解剖#

本章回答:

用户在 TypedDict 里写的 Annotated[list[str], add],如何在构建期被 LangGraph 转换成运行时可调用的 reducer channel?

7.1 构建期职责#

输入归一化动作构建结果
TravelState解析类型注解schema → channels
user_request: str无 reducer,选择默认 channelLastValue(str)
research_notes: Annotated[list[str], add]识别 Annotated metadata 中的 reducerBinaryOperatorAggregate(list[str], add)
messages: Annotated[list[AnyMessage], add_messages]识别消息 reducerreducer channel / message-specific merge

7.2 构建期总链路#

StateGraph(TravelState)
__init__ 初始化 builder 容器
_add_schema(TravelState)
_get_channels(TravelState)
_get_channel(name, annotation)
_is_field_binop(annotation)
BinaryOperatorAggregate(inner_type, reducer)
self.channels[key] = channel
compile 时被挂入 Pregel channels

7.3 StateGraph.__init__() 源码解剖#

职责与所处阶段#

StateGraph.__init__() 处于构建期,负责初始化 builder 内部容器,并立即解析 state_schemainput_schemaoutput_schema

真实源码签名#

以下签名来自 langgraph==1.2.7StateGraph.__init__(),为便于阅读只保留核心参数:

def __init__(
self,
state_schema: type[StateT],
context_schema: type[ContextT] | None = None,
*,
input_schema: type[InputT] | None = None,
output_schema: type[OutputT] | None = None,
**kwargs: Unpack[DeprecatedKwargs],
) -> None:
...

调用方与被调用方#

用户代码 StateGraph(TravelState)
StateGraph.__init__
StateGraph._add_schema(self.state_schema)
_get_channels(schema)

输入、输出与状态变化#

项目类型说明
输入type[StateT]用户定义的 TypedDict / Pydantic / dataclass schema
输出None初始化 builder 本身
状态变化self.channels保存 state key 对应 channel
状态变化self.schemas保存 schema 对应字段信息
副作用无外部副作用只是构建内存对象

细粒度伪代码#

以下为保留关键控制流的简化伪代码,不是源码逐字复制:

def __init__(self, state_schema, context_schema=None, *, input_schema=None, output_schema=None, **kwargs):
# 1. 兼容旧参数:config_schema / input / output
context_schema = resolve_context_schema(context_schema, kwargs)
input_schema = resolve_input_schema(input_schema, kwargs)
output_schema = resolve_output_schema(output_schema, kwargs)
# 2. 初始化 builder 容器
self.nodes = {}
self.edges = set()
self.branches = defaultdict(dict)
self.schemas = {}
self.channels = {}
self.managed = {}
self.waiting_edges = set()
self.compiled = False
# 3. 保存 schema 类型
self.state_schema = state_schema
self.input_schema = input_schema or state_schema
self.output_schema = output_schema or state_schema
self.context_schema = context_schema
# 4. 解析 schema
self._add_schema(self.state_schema)
self._add_schema(self.input_schema, allow_managed=False)
self._add_schema(self.output_schema, allow_managed=False)

逐段解释#

第一段处理历史兼容参数,不影响 reducer 主线,但说明 StateGraph 的构建入口承担 API 兼容职责。

第二段初始化 builder 容器,其中最重要的是:

self.channels

它最终保存每个 state key 对应的 channel。也就是说:

state key 的更新规则在构建期就已经确定,
不是运行到节点返回 update 时才临时猜测。

第三段记录 state_schemainput_schemaoutput_schema。如果不单独指定输入输出 schema,三者默认相同。

第四段调用 _add_schema(),这是 reducer 源码解剖真正的入口。

正常路径#

TravelState
StateGraph.__init__
self._add_schema(TravelState)
self.channels 得到 user_request 与 research_notes 两个 channel

关键分支与异常路径#

条件行为结果
使用旧参数 config_schemawarning 并映射到 context_schema兼容旧代码
未传 input_schema使用 state_schema输入字段与 state 相同
未传 output_schema使用 state_schema输出字段与 state 相同
schema 非合法类型warning可能无法正确解析更新

设计原因与工程影响#

构建期解析 schema 的好处是:

运行时不需要每次 node 返回 update 时再反射类型;
并行冲突可以由 channel 语义统一判断;
compile 后的 Pregel 图可以直接使用 channels。

工程影响:

修改 State schema 后必须重新 compile graph;
compile 后再改 builder 不会影响已编译图;
reducer 是图结构的一部分,不是临时执行参数。

源码证据#

7.4 _add_schema() 源码解剖#

职责与所处阶段#

_add_schema() 负责把一个 schema 类型转换为:

channels: dict[str, BaseChannel]
managed: dict[str, ManagedValueSpec]

它会把这些结果合并进 builder 的全局 self.channelsself.managed

真实源码签名#

def _add_schema(
self,
schema: type[Any],
/,
allow_managed: bool = True,
) -> None:
...

调用方与被调用方#

StateGraph.__init__
StateGraph._add_schema
_get_channels(schema)

输入、输出与状态变化#

项目类型说明
输入type[Any]state / input / output schema
输出None直接修改 builder
状态变化self.schemas[schema]记录该 schema 解析出的字段
状态变化self.channels合并普通 state channels
状态变化self.managed合并 managed values
副作用可能抛异常channel 类型冲突或 input/output 使用 managed value

细粒度伪代码#

def _add_schema(self, schema, allow_managed=True):
# 1. 同一个 schema 只解析一次
if schema in self.schemas:
return
# 2. 如果 schema 形态可疑,发出 warning
_warn_invalid_state_schema(schema)
# 3. 核心:把 schema 字段解析成 channels 与 managed values
channels, managed, type_hints = _get_channels(schema)
# 4. input_schema / output_schema 不允许 managed values
if managed and not allow_managed:
raise ValueError("Managed channels are not permitted in Input/Output schema")
# 5. 保存 schema 视图
self.schemas[schema] = {**channels, **managed}
# 6. 合并普通 channels
for key, channel in channels.items():
if key in self.channels:
if self.channels[key] != channel:
if isinstance(channel, LastValue):
# 允许默认 channel 不覆盖已有 reducer channel
pass
else:
raise ValueError("Channel already exists with a different type")
else:
self.channels[key] = channel
# 7. 合并 managed values
for key, managed_value in managed.items():
if key in self.managed and self.managed[key] != managed_value:
raise ValueError("Managed value already exists with a different type")
self.managed[key] = managed_value

逐段解释#

第一段避免重复解析同一个 schema。因为 state_schemainput_schemaoutput_schema 可能相同。

第三段是核心:

_get_channels(schema)

它把类型字段变成 channel。

第六段包含一个重要分支:如果同一个 key 已经存在 channel,而新的 channel 是默认 LastValue,不会覆盖旧 channel。这个设计避免 input/output schema 的默认字段把 state schema 中更强的 reducer 覆盖掉。

正常路径#

TravelState
_get_channels
user_request -> LastValue(str)
research_notes -> BinaryOperatorAggregate(list[str], add)
self.channels 保存两个 channel

关键分支与异常路径#

条件行为结果
schema 已解析直接返回避免重复工作
input/output schema 有 managed valueValueError防止运行时 managed 字段暴露为输入输出
同一 key channel 类型冲突ValueError防止一个 key 有两套合并语义
新 channel 是 LastValue 且旧 channel 已存在跳过保留更明确的 reducer channel

设计原因与工程影响#

_add_schema() 把 reducer 冲突提前到构建期暴露。如果同一个字段在多个 schema 中被声明为不同 reducer,框架不会等运行时才失败。

工程影响:

不要在 state_schema、input_schema、output_schema 里给同名字段写不同 reducer。
如果只想隐藏输出字段,应使用 input_schema/output_schema 控制可见字段,
而不是改变同名字段的 reducer 语义。

源码证据#

7.5 _get_channels() 源码解剖#

职责与所处阶段#

_get_channels() 是 schema 字段解析入口。它负责读取类型注解,并把每个字段分到:

普通 channel
managed value

真实源码签名#

源码具体签名可能随版本微调,心智签名如下:

def _get_channels(
schema: type[dict],
) -> tuple[
dict[str, BaseChannel],
dict[str, ManagedValueSpec],
dict[str, Any],
]:
...

调用方与被调用方#

StateGraph._add_schema
_get_channels
_get_channel(name, annotation)
_is_field_channel / _is_field_binop / _is_field_managed_value

输入、输出与状态变化#

项目类型说明
输入schema typeTypedDict / Pydantic / dataclass
输出channels每个普通 state key 的 channel
输出managed特殊 managed values
输出type_hints解析后的字段类型
状态变化纯解析函数

细粒度伪代码#

def _get_channels(schema):
# 1. 读取类型注解,保留 Annotated metadata
type_hints = get_type_hints(schema, include_extras=True)
channels = {}
managed = {}
# 2. 遍历每个字段
for name, annotation in type_hints.items():
# 3. 先判断是否是 managed value
if managed_spec := _is_field_managed_value(name, annotation):
managed[name] = managed_spec
continue
# 4. 再解析普通 channel
channel = _get_channel(name, annotation)
channels[name] = channel
return channels, managed, type_hints

逐段解释#

第一段必须保留 Annotated metadata。若 include_extras=False,下面这个字段:

research_notes: Annotated[list[str], add]

会退化成:

research_notes: list[str]

reducer 信息就丢失了。

第三段识别 managed value。managed value 不是普通 state 输入输出字段,它由运行时管理,例如剩余步数、运行时上下文等。

第四段才是普通 state key 的 channel 解析。

正常路径#

TravelState
get_type_hints(include_extras=True)
{
"user_request": str,
"research_notes": Annotated[list[str], add]
}
_get_channel("user_request", str)
_get_channel("research_notes", Annotated[list[str], add])

关键分支与异常路径#

条件行为结果
字段是 managed value放入 managed不作为普通 channel
字段是 Annotated channel解析 metadata可能生成 reducer channel
字段是普通类型默认 LastValue覆盖式更新

设计原因与工程影响#

_get_channels() 说明 LangGraph 的 state schema 本质上是:

字段声明
+
字段更新协议声明

而不是只有字段类型。

工程影响:

字段类型决定值的形态;
Annotated metadata 决定更新语义。

源码证据#

7.6 _get_channel()Annotated 解析源码解剖#

职责与所处阶段#

_get_channel() 负责为单个字段选择 channel 类型。

对于本篇最核心的字段:

research_notes: Annotated[list[str], add]

它应该返回:

BinaryOperatorAggregate(list[str], add)

真实源码签名#

心智签名如下:

def _get_channel(
name: str,
annotation: Any,
*,
allow_managed: bool = True,
) -> BaseChannel | ManagedValueSpec:
...

调用方与被调用方#

_get_channels
_get_channel
_is_field_channel
_is_field_binop
LastValue fallback

输入、输出与状态变化#

项目类型说明
输入name: str字段名
输入annotation: Any字段类型注解,可能是 Annotated
输出BaseChannel字段对应 channel
状态变化纯解析函数

细粒度伪代码#

def _get_channel(name, annotation):
# 1. 如果 Annotated metadata 直接声明了 channel,则使用该 channel
if channel := _is_field_channel(annotation):
channel.key = name
return channel
# 2. 如果 Annotated metadata 声明了二元 reducer,则构造 BinaryOperatorAggregate
if binop_channel := _is_field_binop(annotation):
binop_channel.key = name
return binop_channel
# 3. 如果是 managed value,返回 managed spec
if managed := _is_field_managed_value(name, annotation):
return managed
# 4. 默认:LastValue
return LastValue(annotation, key=name)

逐段解释#

第一段允许高级用户直接把某个字段声明为 channel,这比 reducer 函数更底层。

第二段是最常见 reducer 路径:

Annotated[list[str], add]

add 是一个接收两个参数的 callable,因此被识别为二元 reducer。

第三段处理 managed value。

第四段是默认路径。没有 Annotated reducer 的字段都会变成:

LastValue

也就是覆盖式更新。

正常路径#

user_request: str
无 channel metadata
无 reducer metadata
LastValue(str)
research_notes: Annotated[list[str], add]
识别 add 是二元 reducer
BinaryOperatorAggregate(list[str], add)

关键分支与异常路径#

条件行为结果
metadata 是 BaseChannel直接使用 channel用户完全控制更新规则
metadata 是二元 callable构造 BinaryOperatorAggregatereducer 更新
metadata 是 managed spec返回 managed value不进入普通 state channel
没有 metadataLastValue默认覆盖
reducer 签名不符合二元函数不识别或抛错不能作为 reducer

设计原因与工程影响#

这个设计让 state schema 同时表达两类信息:

value type:字段值是什么类型
update type:字段如何被更新

工程影响:

如果某个字段可能被并行节点同时写入,必须显式声明 reducer。
否则框架不会帮你猜“追加”还是“覆盖”。

源码证据#

7.7 BinaryOperatorAggregate.__init__() 源码解剖#

职责与所处阶段#

BinaryOperatorAggregate.__init__() 负责保存 reducer 函数,并初始化 channel 当前值。

真实源码签名#

def __init__(self, typ: type[Value], operator: Callable[[Value, Value], Value]):
...

调用方与被调用方#

_get_channel
_is_field_binop
BinaryOperatorAggregate(typ, operator)

输入、输出与状态变化#

项目类型说明
输入typ字段值类型,例如 list[str]
输入operatorreducer,例如 operator.add
输出channel instancereducer channel
状态变化self.operator保存 reducer
状态变化self.value初始化为空容器或 MISSING

细粒度伪代码#

def __init__(self, typ, operator):
super().__init__(typ)
# 1. 保存 reducer 函数
self.operator = operator
# 2. 处理 typing 抽象类型,转成可实例化类型
typ = strip_required_notrequired_annotated(typ)
if typ is Sequence or MutableSequence:
typ = list
if typ is Set or MutableSet:
typ = set
if typ is Mapping or MutableMapping:
typ = dict
# 3. 尝试初始化空值
try:
self.value = typ()
except Exception:
self.value = MISSING

逐段解释#

第一段保存 reducer。运行时合并时调用的是:

self.operator(self.value, value)

第二段处理抽象类型。例如:

Sequence[str]

本身不能直接实例化,所以转换成 list

第三段尝试初始化空值。对 list 来说:

list() -> []

所以 research_notes 初始值可以是空列表。

正常路径#

Annotated[list[str], add]
BinaryOperatorAggregate(list[str], add)
self.operator = add
self.value = []

关键分支与异常路径#

条件行为结果
typ 可实例化self.value = typ()得到默认空值
typ 不可实例化self.value = MISSING等待第一条 update 设置初值
typ 是抽象集合类型转为具体容器list / set / dict

设计原因与工程影响#

初始化空值使追加型 reducer 更自然:

[] + ["food"] + ["attraction"]

但工程上仍建议给 graph 输入显式传入初始值,尤其是 list 字段:

{"research_notes": []}

这样可读性更强,也避免不同 channel 初始语义造成误解。

源码证据#


8. 运行时主链源码解剖#

本章回答:

构建完成后,一次 graph.invoke() 中,多个节点返回的 partial update 如何在 super-step 边界被收集、分组并交给 reducer 合并?

8.1 运行时入口#

调用方式公开入口核心内部入口返回类型
同步调用graph.invoke(input)Pregel sync loopoutput state
异步调用graph.ainvoke(input)Pregel async loopoutput state
流式调用graph.stream(input)Pregel stream loopstate/event chunks
批量调用graph.batch(inputs)Runnable batchlist[output state]

Reducer 的核心位置不在这些公开入口本身,而在运行时主链中:

node result
ChannelWrite
pending writes
按 channel 分组
channel.update(values)

8.2 运行时总链路#

graph.invoke(initial_state)
初始化 channel values
Pregel 根据 START 激活 search_food / search_attractions
两个节点并行执行
每个节点返回 partial state update
partial update 转成 channel writes
Pregel 收集同一个 super-step 内所有 writes
按 state key / channel 分组
LastValue.update 或 BinaryOperatorAggregate.update
下一步 state 可见
merge_plan 读取合并后的 research_notes

8.3 输入归一化源码解剖#

职责与所处阶段#

运行时输入归一化负责把用户传入的 initial state 写入对应 channels。

真实源码签名#

这部分分散在 Pregel 调用链中,心智入口为:

graph.invoke(input: InputT, config: RunnableConfig | None = None, **kwargs) -> OutputT

调用方与被调用方#

Caller
CompiledStateGraph.invoke
Pregel.invoke / stream
initialize input channels

输入、输出与状态变化#

项目类型说明
输入dictinitial state
输出channel values写入各 state key channel
状态变化channel 当前值初始化 user_requestresearch_notes
副作用可能校验失败输入包含未知 key 或类型不符合时可能失败

细粒度伪代码#

def invoke(input_state, config=None):
# 1. 准备运行配置
config = ensure_config(config)
# 2. 基于 compiled graph 的 channels 创建本次运行的 channel 副本
channels = clone_compiled_channels()
# 3. 把 initial state 写入输入 channels
for key, value in input_state.items():
if key not in channels:
handle_unknown_key(key)
continue
# 初始输入也走 channel update 语义
channels[key].update([value])
# 4. 启动 Pregel loop
return run_pregel_loop(channels, config)

逐段解释#

初始 state 不是简单赋值到字典,而是进入 channel。这样可以保持输入、节点更新、checkpoint 恢复都使用统一的 channel 值模型。

research_notes

channels["research_notes"].update([[]])

如果它是 BinaryOperatorAggregate(list, add),初始输入 [] 会成为或参与当前 channel value。

正常路径#

input_state["research_notes"] = []
research_notes channel.update([[]])
channel value = []

关键分支与异常路径#

条件行为结果
输入 key 不在 schema可能忽略或校验失败取决于输入 schema 与运行模式
输入 key 是 reducer field进入 reducer channel参与聚合
输入 key 是 LastValue field设置当前值覆盖初始值

设计原因与工程影响#

如果初始输入不走 channel,checkpoint 恢复、stream 更新、node update 会出现不同语义。统一走 channel 能减少状态不一致。

工程影响:

不要假设 state 是直接拷贝的 dict;
最终 state 是 channels 当前值的投影。

源码证据#

8.4 Config、Context 与 State 传播源码解剖#

数据边界#

数据来源生命周期下游消费者
Configgraph.invoke(..., config=...)单次运行node、Runnable、trace、checkpoint
Contextgraph.invoke(..., context=...)单次运行只读上下文node / route function runtime
Statechannels 当前值投影整个图执行所有 node 与 conditional edge
Writesnode return一个 super-step 内 pendingPregel 写入应用阶段

细粒度伪代码#

def run_node(node, channels, config, runtime):
# 1. 从 channels 读取当前 state 快照
state = read_state_snapshot(channels, node.input_schema)
# 2. 调用节点
update = node.invoke(state, config=config, runtime=runtime)
# 3. 节点只返回 partial update
if not isinstance(update, dict):
raise InvalidUpdateError("Expected dict update")
# 4. 把 update 转换成 channel writes
writes = []
for key, value in update.items():
writes.append(ChannelWrite(key, value))
return writes

逐段解释#

第一段从 channels 读 state,而不是从某个 dict 变量读。这意味着上一步 reducer 合并后的值会对当前节点可见。

第二段调用用户节点。用户节点不应该直接修改 channel,而是返回 partial update。

第三段要求 update 是可解释的状态更新。如果 node 返回非 dict,就无法知道写入哪些 state key。

第四段把 dict update 转为 writes,等待 Pregel 在 step 边界统一应用。

正常路径#

search_food(state)
{"research_notes": ["推荐美食:..."]}
ChannelWrite(channel="research_notes", value=["推荐美食:..."])

关键分支与异常路径#

条件行为结果
node 返回 dict转成 writes正常应用
node 返回非 dictInvalidUpdateError节点输出不合法
node 返回未知 key校验或忽略取决于图编译出的 output_keys
node 返回 Command同时包含 update 与 goto进入 Command 分支

设计原因与工程影响#

节点返回 partial update,而不是直接改全局 state,有两个好处:

1. 所有更新都能被 Pregel 统一排序、合并和记录。
2. 并行节点之间不会直接共享可变对象,降低竞态风险。

工程影响:

node 内部不要原地修改 state["research_notes"].append(...)
应该 return {"research_notes": [new_note]}
让 reducer 合并。

8.5 核心执行函数源码解剖:writes 应用阶段#

职责与所处阶段#

writes 应用阶段是 reducer 真正被调用的位置。

它发生在:

当前 super-step 的所有已激活节点执行完成之后,
下一轮节点被激活之前。

真实源码签名#

这部分在 Pregel 内部算法中,心智签名如下:

def apply_writes(
channels: dict[str, BaseChannel],
pending_writes: list[Write],
) -> None:
...

调用方与被调用方#

Pregel loop
execute selected tasks
collect pending writes
apply_writes
channel.update(values)

输入、输出与状态变化#

项目类型说明
输入channels当前图所有 channel
输入pending_writes当前 super-step 的节点写入
输出None直接更新 channel 内部值
状态变化channel values产生下一步 state
副作用可能抛并发写冲突LastValue 多写入时

细粒度伪代码#

def apply_writes(channels, pending_writes):
# 1. 按 channel 名分组
writes_by_channel = defaultdict(list)
for write in pending_writes:
channel_name = write.channel
value = write.value
writes_by_channel[channel_name].append(value)
# 2. 对每个被写入的 channel 应用更新
for channel_name, values in writes_by_channel.items():
channel = channels[channel_name]
# 3. 这里才真正触发 reducer 或覆盖逻辑
changed = channel.update(values)
# 4. 记录哪些 channel 发生了变化,用于触发下游节点、checkpoint、stream
if changed:
mark_channel_changed(channel_name)
# 5. 未收到更新的 channel 不变

逐段解释#

第一段按 channel 分组。这一步决定了并发写入是否会交给同一个 channel。

如果两个节点都返回:

{"research_notes": [...]}

那么会得到:

writes_by_channel["research_notes"] = [food_notes, attraction_notes]

第三段是 reducer 调用发生点:

channel.update(values)

如果 channel 是 BinaryOperatorAggregate,会逐个调用 operator。

如果 channel 是 LastValue,会检查同一 step 是否只有一个值。

正常路径#

pending_writes:
search_food -> research_notes = [food]
search_attractions -> research_notes = [attraction]
writes_by_channel:
research_notes -> [[food], [attraction]]
channel:
BinaryOperatorAggregate(list, add)
channel.update:
[] + [food] + [attraction]

关键分支与异常路径#

条件行为结果
channel 是 reducer channel调用 reducer合并多个写入
channel 是 LastValue 且一个写入覆盖正常
channel 是 LastValue 且多个写入InvalidUpdateError并发写冲突
channel 不存在抛错或忽略内部写入取决于写入类型

设计原因与工程影响#

Reducer 放在 super-step 写入阶段,而不是 node 内部,有一个关键设计原因:

只有 Pregel 运行时知道同一 super-step 里有哪些节点并行完成,
也只有它能把所有 writes 聚合到同一个 channel 后统一处理。

工程影响:

Reducer 是并行语义的一部分。
只在单链路顺序执行时看不出它的重要性;
一旦 fan-out 并行写同一个 key,reducer 就决定系统是否能运行。

源码证据#

8.6 LastValue.update() 源码解剖#

职责与所处阶段#

LastValue 是默认 channel。它的职责是保存最近一次写入值,但它不允许同一 super-step 收到多个值。

真实源码签名#

def update(self, values: Sequence[Value]) -> bool:
...

调用方与被调用方#

apply_writes
channels[key].update(values)
LastValue.update(values)

输入、输出与状态变化#

项目类型说明
输入Sequence[Value]同一 super-step 写给该 key 的所有值
输出boolchannel 是否变化
状态变化self.value更新为唯一值
异常InvalidUpdateError多个值同时写入

细粒度伪代码#

def update(self, values):
# 1. 没有更新,不改变 channel
if len(values) == 0:
return False
# 2. 默认 channel 一个 step 只能收到一个值
if len(values) != 1:
raise InvalidUpdateError(
"Can receive only one value per step. Use an Annotated key to handle multiple values."
)
# 3. 保存唯一值
self.value = values[-1]
return True

逐段解释#

第一段表示如果没有节点写这个 key,state 保持不变。

第二段是并发冲突的核心。默认字段没有 reducer,所以框架不知道多个值应该:

取第一个?
取最后一个?
拼接?
合并 dict?
取最大?
去重?

LangGraph 选择抛错,而不是猜。

第三段才是正常覆盖。

正常路径#

search_food 返回 {destination: "东京"}
destination channel values = ["东京"]
LastValue.update
state.destination = "东京"

关键分支与异常路径#

条件行为结果
values=[]返回 False不改变 state
values=[x]self.value=x覆盖旧值
values=[x, y]InvalidUpdateError并发写冲突

设计原因与工程影响#

默认覆盖语义适合:

destination
current_step
final_plan
risk_level
status

它不适合:

messages
search_results
research_notes
tool_observations
parallel_agent_outputs

工程影响:

如果字段可能被多个并行节点写入,不能使用默认 LastValue。

源码证据#

8.7 BinaryOperatorAggregate.update() 源码解剖#

职责与所处阶段#

BinaryOperatorAggregate.update()Annotated[..., reducer] 的运行时核心。它接收同一 super-step 写给同一个 key 的多个值,并使用二元 operator 逐个聚合。

真实源码签名#

def update(self, values: Sequence[Value]) -> bool:
...

调用方与被调用方#

apply_writes
channels["research_notes"].update(values)
BinaryOperatorAggregate.update(values)
self.operator(self.value, value)

输入、输出与状态变化#

项目类型说明
输入Sequence[Value]同一 key 的多个更新
输出bool是否发生变化
状态变化self.valuereducer 聚合后的值
特殊分支Overwrite可绕过 reducer 强制覆盖

细粒度伪代码#

def update(self, values):
# 1. 没有更新,不改变 channel
if not values:
return False
# 2. 如果当前没有值,用第一个 update 初始化
if self.value is MISSING:
self.value = values[0]
values = values[1:]
seen_overwrite = False
# 3. 逐个处理剩余 update
for value in values:
is_overwrite, overwrite_value = get_overwrite(value)
# 4. Overwrite 分支:绕过 reducer 直接覆盖
if is_overwrite:
if seen_overwrite:
raise InvalidUpdateError("Only one Overwrite per super-step")
self.value = overwrite_value
seen_overwrite = True
continue
# 5. 普通 reducer 分支
if not seen_overwrite:
self.value = self.operator(self.value, value)
return True

逐段解释#

第一段无更新则不变。

第二段处理初始值。如果 channel 当前还没有值,第一条 update 会成为初始值,后续值再与它聚合。

第五段是真正 reducer 逻辑:

self.value = self.operator(self.value, value)

operator.add 和 list 来说:

["food"] + ["attraction"]

得到:

["food", "attraction"]

第四段是 Overwrite 特殊分支,后面第 9 章还会讲。

正常路径#

初始 research_notes = []
values = [
["推荐美食:寿司、拉面、居酒屋"],
["推荐景点:浅草寺、涩谷、上野公园"]
]
operator.add([], [food])
[food]
operator.add([food], [attraction])
[food, attraction]

关键分支与异常路径#

条件行为结果
无更新返回 Falsestate 不变
初始值 MISSING第一条 update 作为初值后续再 reducer
普通值调用 operator聚合
一个 Overwrite直接覆盖绕过 reducer
多个 OverwriteInvalidUpdateError防止覆盖冲突

设计原因与工程影响#

BinaryOperatorAggregate 的设计让 LangGraph 不需要知道业务字段语义,只要用户提供二元函数即可。

工程影响:

Reducer 必须是纯函数。
Reducer 最好满足结合律。
并行场景下最好不要依赖严格顺序。
如果输出顺序重要,建议在 value 中带 source/order 字段,然后在 merge node 排序。

源码证据#

8.8 子步骤调度与等待边源码解剖#

职责与所处阶段#

本例中:

builder.add_edge(["search_food", "search_attractions"], "merge_plan")

表示 merge_plan 需要等待两个上游节点都完成后再运行。

真实源码签名#

def add_edge(self, start_key: str | list[str], end_key: str) -> Self:
...

调用方与被调用方#

用户代码 add_edge([A, B], C)
StateGraph.add_edge
self.waiting_edges.add(((A, B), C))
compile 时转成 barrier / waiting edge 结构

输入、输出与状态变化#

项目类型说明
输入list[str]多个起点节点
输入str目标节点
输出Selfbuilder 链式调用
状态变化self.waiting_edges记录等待边

细粒度伪代码#

def add_edge(self, start_key, end_key):
if isinstance(start_key, str):
validate_single_edge(start_key, end_key)
self.edges.add((start_key, end_key))
return self
# 多起点边:等待所有 start 节点完成
for start in start_key:
if start == END:
raise ValueError("END cannot be a start node")
if start not in self.nodes:
raise ValueError("Need to add_node first")
if end_key == START:
raise ValueError("START cannot be an end node")
if end_key != END and end_key not in self.nodes:
raise ValueError("Need to add_node end first")
self.waiting_edges.add((tuple(start_key), end_key))
return self

逐段解释#

单起点边表示:

A 完成后,触发 B。

多起点边表示:

A 和 B 都完成后,触发 C。

这对 reducer 很重要。否则如果分别写:

builder.add_edge("search_food", "merge_plan")
builder.add_edge("search_attractions", "merge_plan")

可能表达的是两个独立触发路径,而不是“等待两个都完成再合并”。

正常路径#

START
├─ search_food
└─ search_attractions
等待两者都完成
merge_plan

关键分支与异常路径#

条件行为结果
start_key 是字符串普通边单节点完成即触发
start_key 是列表waiting edge所有起点完成才触发
start 是 END抛错END 不能作为起点
end 是 START抛错START 不能作为终点

设计原因与工程影响#

Reducer 解决的是:

多个写入如何合并?

等待边解决的是:

什么时候进入下游合并节点?

两者不能混淆。

工程影响:

并行检索 + 汇总节点场景中,通常需要:
1. reducer 聚合并行结果;
2. waiting edge 保证汇总节点等待所有分支完成。

8.9 运行时主链总结#

Initial State
Channels initialized
Parallel nodes return partial updates
Pending writes grouped by channel
LastValue / BinaryOperatorAggregate / add_messages.update
Merged state visible to downstream nodes

9. 关键分支、异常与边界#

本章回答:

当多个节点写同一个 key、需要覆盖 reducer 字段、使用 messages 字段或 reducer 设计不当时,LangGraph 如何分流、失败或提供边界能力?

9.1 分支矩阵#

分支类型触发条件核心函数结果
默认覆盖单节点写无 reducer keyLastValue.update()覆盖旧值
并发冲突多节点同 step 写无 reducer keyLastValue.update()InvalidUpdateError
reducer 聚合多节点写 reducer keyBinaryOperatorAggregate.update()调用 operator 合并
强制覆盖update 值是 OverwriteBinaryOperatorAggregate.update()绕过 reducer 覆盖
消息合并写入 messagesadd_messages()按 ID 追加或覆盖
等待合并多上游到一个节点waiting edge等上游都完成再执行

9.2 同步与异步分支#

维度同步路径异步路径
入口invoke() / stream()ainvoke() / astream()
节点执行同步 callable 或线程执行async callable / coroutine
reducer 调用step 写入应用阶段async step 写入应用阶段
reducer 本身普通同步函数仍然是同步二元函数
降级行为sync node 正常运行sync node 可能通过 Runnable 适配

Reducer 本身不应该写成异步函数。原因是 channel update 是状态合并阶段,不是业务 I/O 阶段。

9.3 Batch、Stream、Parallel 或路由分支#

Reducer 与这几种执行模式的关系:

Batch:
每个 graph input 独立运行一份 graph state,reducer 只在单次运行内部合并。
Stream:
可以看到每步 updates 或 values,但 reducer 仍在 step 写入应用阶段生效。
Parallel:
reducer 最关键的使用场景;多个节点同 step 写同一 key。
Router:
条件边读取的是 reducer 合并后的 state;如果 state 合并失败,路由不会执行。

9.4 异常分类#

异常类别抛出位置是否可恢复处理策略是否反馈上层
并发写冲突LastValue.update()给字段加 reducer 或拆分 key
reducer 签名错误schema 解析期修正为二元函数
reducer 内部异常BinaryOperatorAggregate.update()视情况修 reducer 或输入数据
多个 OverwriteBinaryOperatorAggregate.update()保证同 step 只有一个覆盖
node 返回非 dictnode output 解析修正 node 返回 partial update
state 类型不一致reducer / channel update视情况加 Pydantic 校验或修 reducer

9.5 并发写冲突源码解剖#

职责与所处阶段#

并发写冲突是 LangGraph 防止不确定状态合并的保护机制。

真实源码签名#

冲突通常在:

LastValue.update(values: Sequence[Value]) -> bool

中抛出。

调用方与被调用方#

apply_writes
LastValue.update([value_from_node_a, value_from_node_b])
InvalidUpdateError

输入、输出与状态变化#

项目类型说明
输入多个 values同一 super-step 写同一 key
输出抛异常
状态变化无或回滚不产生不确定 state
副作用执行失败上层收到错误

细粒度伪代码#

def update(self, values):
if len(values) == 0:
return False
if len(values) > 1:
raise InvalidUpdateError(
"At key X: Can receive only one value per step. "
"Use an Annotated key to handle multiple values."
)
self.value = values[0]
return True

逐段解释#

问题不是“Python 不会合并字典”,而是“业务语义不明确”。

比如两个并行节点同时写:

{"destination": "东京"}
{"destination": "大阪"}

框架无法知道应该:

取东京?
取大阪?
报冲突?
让用户选择?
拆成候选目的地列表?

所以默认选择失败。

正常路径#

如果只有一个写入:

values = ["东京"]
state.destination = "东京"

关键分支与异常路径#

条件行为结果
一个值覆盖正常
多个值抛错防止不确定合并
需要多个值添加 reducer显式声明合并规则

设计原因与工程影响#

这体现了 LangGraph 的工程哲学:

框架不猜业务合并语义。

工程影响:

并行 fan-out 前必须检查多个分支是否会写同一个 key。

源码证据#

9.6 Overwrite 分支源码解剖#

职责与所处阶段#

Overwrite 用于绕过 reducer,对 reducer 字段进行强制覆盖。

例如:

research_notes: Annotated[list[str], add]

默认会追加。但某些节点可能希望重置:

return {"research_notes": Overwrite([])}

真实源码签名#

Overwrite 分支在 BinaryOperatorAggregate.update() 内部通过 _get_overwrite() 检测。

调用方与被调用方#

BinaryOperatorAggregate.update(values)
_get_overwrite(value)
if is_overwrite: self.value = overwrite_value

输入、输出与状态变化#

项目类型说明
输入Overwrite(value)强制覆盖值
输出boolchannel changed
状态变化self.value被覆盖为 overwrite value
异常多个 overwrite抛并发冲突

细粒度伪代码#

def update(self, values):
seen_overwrite = False
for value in values:
is_overwrite, overwrite_value = get_overwrite(value)
if is_overwrite:
if seen_overwrite:
raise InvalidUpdateError("Can receive only one Overwrite value per super-step")
self.value = overwrite_value
seen_overwrite = True
continue
if not seen_overwrite:
self.value = self.operator(self.value, value)
return True

逐段解释#

Overwrite 是一个逃生口:当字段默认是追加,但某个业务节点需要重置时,可以绕过 reducer。

但同一个 super-step 里多个 overwrite 仍然冲突,因为框架仍然不知道哪个覆盖值优先。

正常路径#

state.research_notes = [old]
return {research_notes: Overwrite([])}
state.research_notes = []

关键分支与异常路径#

条件行为结果
一个 Overwrite直接覆盖正常
多个 Overwrite抛错防止覆盖冲突
Overwrite 后还有普通值取决于 channel 实现应避免同 step 混用

设计原因与工程影响#

Reducer 字段不是永远只能 reducer。Overwrite 提供了 reset 能力。

工程建议:

Overwrite 应只用于明确的重置节点,
不要在普通并行节点中混用 Overwrite 与 reducer update。

9.7 messages 为什么通常需要特殊 reducer#

职责与所处阶段#

messages 是 Agent 系统中最常见的状态字段。它看起来像 list,但不应该简单使用:

Annotated[list[AnyMessage], operator.add]

更推荐:

from langgraph.graph.message import add_messages
class State(TypedDict):
messages: Annotated[list[AnyMessage], add_messages]

真实源码签名#

def add_messages(
left: Messages,
right: Messages,
*,
format: Literal["langchain-openai"] | None = None,
) -> Messages:
...

调用方与被调用方#

State schema Annotated[list[AnyMessage], add_messages]
BinaryOperatorAggregate(..., add_messages)
channel.update(values)
add_messages(left, right)

输入、输出与状态变化#

项目类型说明
输入old messages当前消息列表
输入new messages节点新增或更新消息
输出merged messages合并后的消息列表
状态变化messages追加新消息,同 ID 覆盖旧消息

细粒度伪代码#

def add_messages(left, right, *, format=None):
# 1. 归一化输入为 list
left_list = normalize_to_message_list(left)
right_list = normalize_to_message_list(right)
# 2. 把 dict / tuple 等反序列化为 LangChain Message 对象
left_msgs = convert_to_messages(left_list)
right_msgs = convert_to_messages(right_list)
# 3. 确保消息有 ID;没有 ID 时生成 ID
for msg in left_msgs + right_msgs:
if msg.id is None:
msg.id = generate_uuid()
# 4. 按 ID 建索引
merged = list(left_msgs)
index_by_id = {msg.id: i for i, msg in enumerate(merged)}
# 5. 新消息:没有 ID 冲突就 append,有 ID 冲突就覆盖
for msg in right_msgs:
if msg.id in index_by_id:
merged[index_by_id[msg.id]] = msg
else:
index_by_id[msg.id] = len(merged)
merged.append(msg)
# 6. 可选格式转换
if format == "langchain-openai":
merged = convert_to_openai_message_format(merged)
return merged

逐段解释#

第一、二段说明 add_messages 不只是 list 拼接。它还支持把下面这种输入转成 message 对象:

{"type": "human", "content": "hello"}

第四、五段是关键:同 ID 消息会覆盖,而不是追加。

这对 human-in-the-loop 很重要。比如人工修改某条 AIMessage,如果用 operator.add,会新增一条修正版;如果用 add_messages,可以更新原消息。

正常路径#

old messages:
[HumanMessage(id=1), AIMessage(id=2)]
new messages:
[ToolMessage(id=3)]
merged:
[HumanMessage(id=1), AIMessage(id=2), ToolMessage(id=3)]

同 ID 覆盖:

old:
[AIMessage(id=2, content="bad answer")]
new:
[AIMessage(id=2, content="fixed answer")]
merged:
[AIMessage(id=2, content="fixed answer")]

关键分支与异常路径#

条件行为结果
new message ID 不存在append追加
new message ID 已存在replace覆盖旧消息
输入是 dict反序列化Message 对象
输入格式非法抛错消息无法合并

设计原因与工程影响#

messages 需要特殊 reducer,因为它同时承担:

对话历史追加
工具结果追加
人工修正覆盖
消息反序列化
供应商格式适配

工程影响:

Agent State 里的 messages 字段默认应使用 add_messages,
不要简单用 operator.add。

源码证据#

9.8 Retry、Fallback 与恢复边界#

机制适用条件不适用条件幂等要求
Retryreducer 内偶发处理失败,且 update 幂等reducer 有副作用reducer 必须纯函数
Fallbackreducer 设计错误导致线上失败状态已经部分提交且无法回滚需有 checkpoint 边界
Repairstate 字段内容不合规但可修复channel 语义错误repair 节点返回合法 update

Reducer 不应该访问网络、数据库或全局可变状态。否则 retry / replay / checkpoint 恢复时会变得不可预测。

9.9 停止条件与保护上限#

Reducer 本身不决定图何时停止,但它影响图能否进入下一步。

正常结束:所有 writes 被成功应用,并最终到达 END。
提前结束:条件边返回 END。
人工中断:interrupt 在节点内暂停,需要 checkpointer。
框架保护:recursion_limit 防止无限循环。
异常失败:channel.update 抛出 InvalidUpdateError 或 reducer 异常。

9.10 能力边界#

容易误判的能力实际提供者本篇对象的真实职责
并行节点调度Pregelreducer 只合并结果,不调度任务
下游何时执行edge / waiting edgereducer 只决定 key 如何合并
状态持久化checkpointerreducer 只生成当前 channel value
消息语义修正add_messages普通 list reducer 不懂消息 ID
冲突业务决策业务节点 / verifierreducer 不应替业务做冲突判断

10. 扩展机制与框架协作#

本章回答:

框架允许在哪里插入自定义状态合并行为?Reducer 如何与 StateGraph、Pregel、MessagesState、ToolNode、Checkpoint 协作?

10.1 扩展点总览#

扩展点扩展方式执行时机可修改内容约束
自定义 reducerAnnotated[T, reducer_fn]channel.update单个 key 合并语义必须是二元函数,建议纯函数
直接 channelAnnotated[T, BaseChannel]channel lifecycle更底层更新规则高级用法,业务少用
add_messagesAnnotated[list[AnyMessage], add_messages]messages update消息追加/覆盖适合对话历史
Overwrite返回 Overwrite(value)reducer update强制覆盖 reducer 字段同 step 只能一个 overwrite
waiting edgeadd_edge([A, B], C)调度阶段等待多个上游不负责合并值
checkpointcompile(checkpointer=...)super-step 边界保存 channel values不改变 reducer 语义

10.2 自定义 reducer 源码解剖#

职责与所处阶段#

自定义 reducer 允许业务定义字段合并语义。

例如去重追加:

def dedupe_add(left: list[str], right: list[str]) -> list[str]:
result = list(left)
for item in right:
if item not in result:
result.append(item)
return result
class TravelState(TypedDict):
research_notes: Annotated[list[str], dedupe_add]

真实源码签名#

自定义 reducer 必须符合:

def reducer(left: Value, right: Value) -> Value:
...

调用方与被调用方#

StateGraph schema parser
BinaryOperatorAggregate(typ, dedupe_add)
channel.update(values)
dedupe_add(left, right)

输入、输出与状态变化#

项目类型说明
输入left当前 channel value
输入right新 update value
输出Value合并后 value
状态变化channel value替换为 reducer 返回值

细粒度伪代码#

def dedupe_add(left, right):
# 1. 不修改原对象,先复制
merged = list(left)
# 2. 逐个处理新值
for item in right:
if item not in merged:
merged.append(item)
# 3. 返回新对象
return merged

逐段解释#

第一段避免原地修改。Reducer 作为框架运行时的一部分,最好保持纯函数语义。

第二段体现业务合并规则:追加但去重。

第三段返回合并结果,由 BinaryOperatorAggregate 保存到 self.value

正常路径#

left = ["寿司", "拉面"]
right = ["拉面", "居酒屋"]
["寿司", "拉面", "居酒屋"]

关键分支与异常路径#

条件行为结果
right 为空返回 left copy不新增
有重复跳过去重
reducer 抛异常graph run 失败上层处理

设计原因与工程影响#

自定义 reducer 是 StateGraph 的主要扩展点之一,但它只适合处理“字段合并”,不适合处理复杂业务流程。

工程影响:

如果合并逻辑需要访问数据库、调用模型、读取其他 state key,
不要写进 reducer,应该写成 merge node。

10.3 扩展调用链#

Node returns update
ChannelWrite
Pregel groups writes by channel
BinaryOperatorAggregate.update
custom reducer(left, right)
channel.value = merged
Downstream node reads merged state

10.4 与相邻框架组件的协作#

相邻组件输入协议输出协议协作边界
ToolNodemessages / tool_callsToolMessage updates通常依赖 messages: add_messages 合并
Conditional Edgemerged statenext node name读取 reducer 后的 state 做路由
Checkpointchannel valuesstate snapshot保存合并后的 channel 值
Interruptpaused stateresumed update恢复后 update 仍按 reducer 合并
Subgraphparent/child statemapped updates共享 key 时 reducer 语义要一致

10.5 公共扩展接口与内部实现#

业务代码可以依赖:
Annotated[T, reducer_fn]
add_messages
MessagesState
Overwrite
StateGraph.add_edge([A, B], C)
业务代码避免依赖:
_get_channels
_get_channel
BinaryOperatorAggregate.value
Pregel 内部 pending_writes 结构
channel.update 调用顺序细节

10.6 自定义扩展示例:多工具检索结果聚合#

示例只展示扩展契约,不重新实现框架:

from typing import Annotated
from typing_extensions import TypedDict
class SearchResult(TypedDict):
source: str
title: str
score: float
def merge_search_results(
left: list[SearchResult],
right: list[SearchResult],
) -> list[SearchResult]:
by_key = {(r["source"], r["title"]): r for r in left}
for item in right:
key = (item["source"], item["title"])
if key not in by_key or item["score"] > by_key[key]["score"]:
by_key[key] = item
return sorted(
by_key.values(),
key=lambda x: x["score"],
reverse=True,
)
class TravelState(TypedDict):
search_results: Annotated[list[SearchResult], merge_search_results]

说明:

  1. 扩展点接收当前值与新值。
  2. 扩展点只处理一个字段,不读取整个 state。
  3. 扩展点必须返回同类型值。
  4. 异常会导致图执行失败。
  5. 若需要外部资源,应移到 node 中处理。

10.7 选择扩展还是重写流程#

条件选择 reducer选择 merge node / 更底层流程
只是追加列表
只是去重聚合
需要按分数排序可以视复杂度
需要访问多个 state key
需要调用 LLM 判断
需要工具查询数据库
需要人工确认

11. 工程决策与适用场景#

11.1 适用场景#

场景是否推荐 reducer原因
并行检索结果聚合多个检索节点写同一个 results key
多 Agent 研究笔记聚合多个专家产生 notes
对话 messages需要追加与同 ID 覆盖
单一目的地 destination应保持单值覆盖,冲突应显式处理
最终计划 final_plan通常只保留最终版本
错误列表 errors多节点都可能报告错误
高风险动作结果视情况多个动作结果应结构化聚合,不要盲目 add

11.2 工程决策表#

决策点推荐选择前提风险
多节点写 listAnnotated[list[T], add] 或自定义 reducer只需追加顺序可能不应被依赖
messagesadd_messages对话历史 / tool results不要用普通 add 替代
单值字段默认 LastValue同 step 只有一个写入并行写入会报错
需要重置 reducer 字段Overwrite(value)明确 reset 节点多 overwrite 会冲突
检索结果去重自定义 reducer合并逻辑纯函数化复杂度上升
需要跨字段判断merge node读取完整 state不要塞进 reducer

11.3 性能、可靠性与安全边界#

性能:
reducer 会在每个相关 super-step 执行,复杂 reducer 会增加状态合并开销。
可靠性:
reducer 必须稳定、纯函数、可重复执行,避免随机性和外部副作用。
安全:
reducer 不应做权限判断;高风险决策应放在显式 verifier / policy node。
可观测性:
必须记录哪些节点写了哪个 key、合并前 values、合并后 state,便于排查并发冲突。

12. 常见误区与源码纠正#

12.1 误区:LangGraph 的 state 就是普通 dict#

错误原因:

用户在节点里看到的 state 确实像 dict,因此容易以为节点返回也是普通 dict merge。

源码事实:

state schema 会被解析成 channels;
node 返回的 dict 会转成 channel writes;
每个 key 的 channel 决定如何应用 update。

工程影响:

把 state 当普通 dict,会在并行分支时遇到不可解释的冲突或覆盖。

12.2 误区:没有 reducer 时多个写入会自动取最后一个#

错误原因:

Python dict update 通常是后者覆盖前者,很多人把这个经验迁移到 LangGraph。

源码事实:

默认 LastValue 一个 super-step 只能接收一个值;
多个并行写入会抛 INVALID_CONCURRENT_GRAPH_UPDATE。

工程影响:

如果希望多个分支都写同一个 key,必须显式声明 reducer。

12.3 误区:operator.add 适合所有 list 字段#

错误原因:

operator.add 简单直接,能解决 append 问题。

源码事实:

operator.add 只是 list 拼接,不去重、不按 ID 覆盖、不做类型转换。

工程影响:

对 messages 应用 operator.add 会导致人工修改消息时无法覆盖旧消息。

12.4 误区:Reducer 可以写复杂业务逻辑#

错误原因:

Reducer 能拿到 old/new 值,看起来可以做很多事。

源码事实:

Reducer 是 channel update 的一部分,应是二元纯函数。

工程影响:

把 LLM 调用、数据库访问、权限判断写入 reducer,会破坏 replay、retry 和可观测性。

12.5 误区:Reducer 决定下一个节点#

错误原因:

Reducer 发生在 state 更新阶段,而路由也读取 state,容易混淆。

源码事实:

Reducer 只合并 state;
Conditional Edge 才决定下一个节点。

工程影响:

需要路由时写 route function,不要让 reducer 返回控制流。

12.6 误区:并行合并后顺序一定稳定#

错误原因:

本地测试中 list 顺序可能看起来稳定。

源码事实:

并行任务、调度、写入收集和未来版本都可能影响顺序假设。

工程影响:

如果顺序重要,应在 value 中携带 sourceranktimestamp,然后在 merge node 显式排序。


13. 最终心智模型与掌握检查#

13.1 构建期心智模型#

TypedDict State
get_type_hints(include_extras=True)
普通字段 → LastValue
Annotated[T, reducer] → BinaryOperatorAggregate
messages → add_messages reducer
channels 保存到 StateGraph
compile 挂入 Pregel

13.2 运行时心智模型#

Node reads current state
Node returns partial update
Update becomes channel write
Pregel groups writes by channel per super-step
channel.update(values)
Reducer or LastValue produces next state
Downstream node reads merged state

13.3 分支与异常心智模型#

正常路径:
单节点写 LastValue,或多节点写 reducer key。
分支路径:
多节点并行写同一 reducer key,channel 聚合。
可恢复异常:
INVALID_CONCURRENT_GRAPH_UPDATE,给 key 增加 reducer 或拆分 key。
不可恢复异常:
reducer 内部错误、node 返回非法 update。
保护上限:
reducer 不控制 recursion_limit;循环保护由 Pregel runtime 处理。

13.4 一句话总结#

LangGraph 通过 State schema 在构建期把每个 state key 转换成 channel,运行时把节点 partial update 收集为 channel writes,并在 super-step 边界调用 LastValueBinaryOperatorAggregateadd_messages 完成状态合并;因此 state 不是普通字典,而是带更新协议的状态模型。

13.5 掌握检查#

  • 能解释 Annotated[list[str], add] 如何被解析成 reducer channel。
  • 能解释默认字段为什么是覆盖式更新。
  • 能解释多节点写同一个默认字段为什么会报错。
  • 能说明 reducer 在 super-step 写入应用阶段被调用。
  • 能写出 BinaryOperatorAggregate.update() 的核心伪代码。
  • 能写出 LastValue.update() 的核心伪代码。
  • 能解释 add_messages 为什么不是普通 list add。
  • 能区分 reducer 与 waiting edge 的职责。
  • 能判断旅行规划助手中哪些字段该覆盖、哪些字段该追加。
  • 能说明 reducer 不适合承载复杂业务逻辑。

14. 参考资料与下一篇衔接#

14.1 官方概念文档#

  1. LangGraph Graph API:State、Reducers、Messages、Nodes、Edges
    https://docs.langchain.com/oss/python/langgraph/graph-api

  2. LangGraph INVALID_CONCURRENT_GRAPH_UPDATE 错误文档
    https://docs.langchain.com/oss/python/langgraph/errors/INVALID_CONCURRENT_GRAPH_UPDATE

  3. LangGraph PyPI 项目说明
    https://pypi.org/project/langgraph/

14.2 官方 API Reference#

  1. StateGraph
    https://reference.langchain.com/python/langgraph/graph/state/StateGraph

  2. add_messages
    https://reference.langchain.com/python/langgraph/graph/message/add_messages

  3. CompiledStateGraph
    https://reference.langchain.com/python/langgraph/graphs

14.3 官方源码#

  1. langgraph/graph/state.py
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py

  2. langgraph/channels/last_value.py
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/channels/last_value.py

  3. langgraph/channels/binop.py
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/channels/binop.py

  4. langgraph/graph/message.py
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/message.py

  5. langgraph/pregel/_algo.py
    https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/_algo.py

14.4 下一篇衔接#

下一篇进入:

第 9 篇:Checkpoint / Thread / Durable Execution 源码解剖

需要继续回答:

checkpoint 在哪个 super-step 边界保存?
thread_id 如何区分不同会话执行线?
StateSnapshot 保存哪些 channel values?
get_state / get_state_history 如何读取历史状态?
checkpoint 如何支撑 interrupt、恢复、time travel 和线上故障排查?
LangGraph 源码深潜:Reducer 与并行状态合并机制解剖
https://jupiter-ws.cn/posts/agent-frameworks/langgraph-reducer-parallel-state-deep-dive/
作者
Jupiter
发布于
2026-03-12
许可协议
CC BY-NC-SA 4.0