reducers 入门:字段的"合并规则"
阶段1(D01-06)你已经能跑通一个图、看懂 StateGraph 怎么把 State schema 拆成一个个 channel。今天进入阶段2第一课:reducer——它决定"多个节点往同一个字段写值时,这些值怎么合并"。我们从源码看清 Annotated[x, reducer] 这行注解,是怎么被 graph/state.py 一步步读出来、变成一个通道的。
Annotated[list, add] 里写的那个函数,就是在告诉 LangGraph:"这块栏用续写规矩"。今天全程记住这句就通。痛点:为什么每个字段要有自己的"合并规则"
messages 写,或者一个节点这一轮写、下一轮又写——这些值到底是后者覆盖前者,还是拼在一起?如果全 State 只有一种合并方式,就没法同时满足"当前用户要覆盖"和"聊天记录要累积"这两种完全相反的需求。(旧值, 新值) -> 合并后的值。写状态时,框架不是简单地 state[key] = 新值,而是 state[key] = reducer(旧值, 新值)。想覆盖,reducer 就写成"返回新值";想累积,就写成"旧列表 + 新列表"。"字段各有合并规则" = "字段各配一个 reducer"。官方 docstring 里就给了一个最小 reducer 的例子(graph/state.py:167):
# graph/state.py:167 —— StateGraph 类文档里的示例 reducer
def reducer(a: list, b: int | None) -> list:
if b is not None:
return a + [b]
return a
class State(TypedDict):
x: Annotated[list, reducer]
def reducer(a, b)a=旧值(累积的 list),b=这次新写入的值。返回值就是"合并后"要存回去的新状态。if b is not None: a + [b]这条规矩是"把新值追加进列表",并且顺手处理了 b is None 的情况(新值为空就原样返回,不追加)。x: Annotated[list, reducer]关键:把 reducer 挂在字段类型的注解里。list 是这个字段的值类型,reducer 是它的合并规矩。下一讲专门拆 Annotated。Annotated[x, reducer] 到底是什么
Annotated 是 Python 标准库 typing 提供的"给类型贴便利贴"的工具。Annotated[list, reducer] 的意思是:这个字段的类型是 list,另外还附带一条元数据 reducer。类型检查器只看第一个 list,而 reducer 这条元数据"搭便车"存进了一个叫 __metadata__ 的属性里,等 LangGraph 来取。
from typing import Annotatedt = Annotated[list, "我是便利贴"]t.__metadata__ → ('我是便利贴',)(一个元组,存着所有便利贴)t.__origin__ → <class 'list'>(脱掉便利贴后的真实类型)LangGraph 就是靠读
__metadata__ 这个元组,把你挂的 reducer 抠出来。
reducers = {"messages": add_messages} 一个字典。源码没这么做,而是塞进 Annotated。好处:① 字段定义和它的合并规则写在同一行,不会脱节——你删字段时不会忘删对应配置;② 类型检查器仍然把它当普通 list,IDE 补全、mypy 检查都不受影响(便利贴对类型系统"透明");③ 不需要额外的 API,用的是 Python 原生机制。代价是:初学者第一次见 Annotated[list, add_messages] 会懵——但看懂机制后就一目了然。_get_channels:把整个 schema 读成"通道表"
编译图时,LangGraph 要把你的 State 类翻译成一组 channel。入口函数是 _get_channels()(graph/state.py:1801):
# graph/state.py:1801
def _get_channels(schema):
if not hasattr(schema, "__annotations__"): # ① schema 不是"带字段的类"
return (
{"__root__": _get_channel("__root__", schema, allow_managed=False)},
{}, {},
)
type_hints = get_type_hints(schema, include_extras=True) # ② 读出所有字段的类型注解
all_keys = {
name: _get_channel(name, typ) # ③ 每个字段各算一个 channel
for name, typ in type_hints.items()
if name != "__slots__"
}
return (
{k: v for k, v in all_keys.items() if isinstance(v, BaseChannel)}, # 真正的通道
{k: v for k, v in all_keys.items() if is_managed_value(v)}, # 托管值(后面章节)
type_hints,
)
hasattr(schema,"__annotations__")判断你传的 schema 到底是不是一个"有字段的类"(TypedDict / dataclass / pydantic 都有 __annotations__)。如果不是(比如你直接传了个 int),就走 __root__ 分支——整个 State 就是一个值。get_type_hints(..., include_extras=True)★核心。include_extras=True 是关键——它让返回的类型保留 Annotated 的便利贴。如果不加这个参数,Annotated[list, add_messages] 会被剥成光秃秃的 list,reducer 就丢了!for name, typ in ...: _get_channel遍历每个字段,逐个调用 _get_channel(下一讲)把"字段类型"翻译成"通道对象"。三个返回值把结果分成三桶:真正的通道(channels)、托管值(managed,后面阶段讲)、以及原始的 type_hints(备用)。_get_channels 就是"读你的 State 类 → 吐出一张 {字段名: 通道对象} 的表"。这张表就是整个图运行时的"存储骨架"。而每个通道具体长什么样,全看下一讲的 _get_channel。_get_channel:四级分派决定用哪种通道
单个字段怎么变成通道?_get_channel()(graph/state.py:1836)是一条清晰的"四级分派"链:
# graph/state.py:1836
def _get_channel(name, annotation, *, allow_managed=True):
# 先剥掉 Required / NotRequired 外壳(TypedDict 的可选标记)
if hasattr(annotation, "__origin__") and annotation.__origin__ in (Required, NotRequired):
annotation = annotation.__args__[0]
if manager := _is_field_managed_value(name, annotation): # ① 是不是托管值?
if allow_managed:
return manager
else:
raise ValueError(f"This {annotation} not allowed in this position")
elif channel := _is_field_channel(annotation): # ② 便利贴里直接放了 channel?
channel.key = name
return channel
elif channel := _is_field_binop(annotation): # ③ 便利贴里是个 reducer 函数?
channel.key = name
return channel
fallback: LastValue = LastValue(annotation) # ④ 什么都没写 → 默认 LastValue
fallback.key = name
return fallback
剥 Required/NotRequiredTypedDict 里字段可以标 Required[...]/NotRequired[...]表示必填/选填。这跟"用哪种通道"无关,所以先脱掉这层壳,只看里面真正的类型。① _is_field_managed_value托管值是框架自己管理的特殊字段(如 RemainingSteps),本阶段先跳过,D11 会提到它"不许出现在输入/输出 schema"。② _is_field_channel如果你便利贴里直接塞了一个通道实例或通道类(如 Annotated[list, Topic(str)]),就用你指定的那个通道。③ _is_field_binop★最常见。如果便利贴最后一个是一个可调用的 reducer 函数(如 add_messages、operator.add),就包成 BinaryOperatorAggregate 通道(D09 细讲)。④ fallback = LastValue★兜底。字段什么便利贴都没贴(如 count: int),就给它 LastValue——"覆盖写"通道。这就是"不写 reducer 默认覆盖"的真相(L06 细讲)。_is_field_binop:从便利贴里认出 reducer
分支③是最常用的。它怎么判断"便利贴里那个东西是不是一个 reducer"?看 _is_field_binop()(graph/state.py:1890):
# graph/state.py:1890
def _is_field_binop(typ):
if hasattr(typ, "__metadata__"): # 有便利贴吗?
meta = typ.__metadata__
if len(meta) >= 1 and callable(meta[-1]): # 最后一张便利贴是"可调用的"吗?
sig = signature(meta[-1])
params = list(sig.parameters.values())
if (
sum(
p.kind in (p.POSITIONAL_ONLY, p.POSITIONAL_OR_KEYWORD)
for p in params
) == 2 # ★必须恰好接收 2 个位置参数
):
return BinaryOperatorAggregate(typ, meta[-1])
else:
raise ValueError(
f"Invalid reducer signature. Expected (a, b) -> c. Got {sig}"
)
return None
hasattr(typ,"__metadata__")没便利贴(普通 int)直接返回 None,交给④兜底。callable(meta[-1])取最后一张便利贴,看它是不是"可调用的"(函数/lambda)。是函数才可能是 reducer。signature(...)用 inspect.signature 读这个函数的参数签名——这是运行时"反省"函数长相的标准手段。sum(...) == 2★关键校验:数一数有几个"位置参数",必须恰好 2 个(对应 reducer 的 (旧值, 新值))。不是 2 个就当场报错。BinaryOperatorAggregate(typ, meta[-1])校验通过,把"值类型 + reducer 函数"打包成一个累加通道返回。D09 会拆开这个通道内部。lambda x: x 当 reducer,这里 sum(...) == 2 会失败,直接抛 ValueError: Invalid reducer signature. Expected (a, b) -> c。这是好事——错误在编译期就暴露,而不是等到运行时某个超步合并数据时才神秘崩溃。源码宁可"早失败、报清楚",也不让一个坏 reducer 悄悄溜进图里。顺带看分支②的判断 _is_field_channel()(graph/state.py:1862),它认的是"便利贴里直接放了通道实例或通道类":
# graph/state.py:1862(节选)
def _is_field_channel(typ):
if hasattr(typ, "__metadata__"):
for item in typ.__metadata__:
if isinstance(item, BaseChannel): # 便利贴是"通道实例"
return item
elif isclass(item) and issubclass(item, BaseChannel): # 便利贴是"通道类"
return item(typ.__origin__ if hasattr(typ, "__origin__") else typ)
return None
没写 reducer = LastValue(覆盖写)
分支④兜底给的是 LastValue。它的合并逻辑(update)解释了"为什么不写 reducer 的字段是覆盖语义、而且并行写会报错"(channels/last_value.py:56):
# channels/last_value.py:56
def update(self, values):
if len(values) == 0:
return False
if len(values) != 1: # ★同一步收到 >1 个值 → 报错
msg = create_error_message(
message=f"At key '{self.key}': Can receive only one value per step. "
"Use an Annotated key to handle multiple values.",
error_code=ErrorCode.INVALID_CONCURRENT_GRAPH_UPDATE,
)
raise InvalidUpdateError(msg)
self.value = values[-1] # 只有 1 个值 → 直接覆盖
return True
values注意参数是一个列表——框架把"本超步内所有对这个字段的写入"收集成列表,一次性交给通道 update。len(values) != 1 → 报错★LastValue 只接受"这一步恰好写了 1 次"。如果两个并行节点同一步都写了它,values 长度是 2,它不知道该听谁的,于是抛 InvalidUpdateError,并贴心提示"用 Annotated key 处理多值"。self.value = values[-1]正常情况:取那唯一的值覆盖旧值。这就是"覆盖写"的字面实现。count: int(没便利贴)→ 走 _get_channel 分支④ → LastValue(int) → 覆盖语义、且不允许同一步并行写。而 messages: Annotated[list, add_messages](有 reducer 便利贴)→ 走分支③ → 累加通道 → 追加语义、并行写能合并。D10 会专门演示"并行写同一字段",两种通道的天壤之别。👶 小白:那我什么时候该写 reducer、什么时候不写?
👨🏫 老师:判断标准很简单——这个字段会不会被"多次写入并且你想保留历史"?会(聊天记录、累计分数、日志列表)→ 写 reducer;不会(当前用户、当前步骤名、一次性结论)→ 不写,用默认 LastValue 覆盖即可。另外,只要你打算让多个并行节点同时写同一字段,就必须给它 reducer,否则 LastValue 会因为收到多值而报错。
边界:schema 传错会怎样 + 今日小结
最后看一个防御性边界。如果你传给 StateGraph 的 schema 既不是类、也不是 Annotated[...],框架不会硬崩,而是发一条警告(graph/state.py:111):
# graph/state.py:111
def _warn_invalid_state_schema(schema):
if isinstance(schema, type): # 是个类(TypedDict/dataclass/pydantic)→ OK
return
if typing.get_args(schema): # 是 Annotated[...] 之类带参数的 → 也 OK
return
warnings.warn(
f"Invalid state_schema: {schema}. Expected a type or Annotated[type, reducer]. "
"Please provide a valid schema to ensure correct updates.\n ..."
)
isinstance(schema, type)正常情况:你传的是 class State(TypedDict) 这种类,直接放行。typing.get_args(schema)兼容你直接传 Annotated[int, reducer] 当整个 State(罕见但合法),它带参数,也放行。warnings.warn边界处理姿态:传错了只警告不报错——因为后面 _get_channels 还能靠 __root__ 分支兜底跑起来。但会提醒你"更新可能不符合预期",让你自己决定要不要改。__root__ 兜底路径能勉强跑,属于"可能是你故意的"。该硬的地方硬报错、能兜的地方给条活路——这是成熟库对"错误严重程度"的分级处理。🧠 今天你应该能回答
- reducer 是什么?签名长啥样?(字段级合并函数,
(旧值, 新值) -> 合并值) Annotated[list, add_messages]拆开后,list 和 add_messages 分别落在哪?(__origin__和__metadata__)get_type_hints为什么要加include_extras=True?(不加就把便利贴 reducer 剥掉了)_get_channel的四级分派顺序?(托管值 → 通道实例 → reducer 函数 → 兜底 LastValue)- 为什么不写 reducer 的字段并行写会报错?(LastValue 收到多值抛 InvalidUpdateError)
- reducer 参数不是 2 个会怎样?(编译期直接 raise ValueError)
✋ 10 分钟动手
# 1. 读 reducer 读取的三个核心函数
sed -n '1801,1908p' libs/langgraph/langgraph/graph/state.py
# 2. 读 LastValue 的覆盖 + 并行报错逻辑
sed -n '56,67p' libs/langgraph/langgraph/channels/last_value.py
# 3. 亲手验证 Annotated 的结构
python -c "from typing import Annotated; t=Annotated[list,'x']; print(t.__origin__, t.__metadata__)"
add_messages。明天钻进 graph/message.py,逐行拆它怎么做"按 ID 去重/更新、RemoveMessage 删除、REMOVE_ALL 清空"——一个远比 operator.add 复杂的工业级 reducer。