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

reducers 入门:字段的"合并规则"

阶段1(D01-06)你已经能跑通一个图、看懂 StateGraph 怎么把 State schema 拆成一个个 channel。今天进入阶段2第一课:reducer——它决定"多个节点往同一个字段写值时,这些值怎么合并"。我们从源码看清 Annotated[x, reducer] 这行注解,是怎么被 graph/state.py 一步步读出来、变成一个通道的。

📍 你在 60 天里的位置(阶段2:状态与数据流 · 共 6 天)
阶段1 入门 D07 reducers D08 add_messages D09 累加通道 D10 并行写 D11 输入输出 D12 Pydantic 阶段3 控制流
💡 先用一个类比兜住今天 把 State 的每个字段想成公司墙上的一块公告栏。节点从不当面说话,都是"往公告栏贴、从公告栏读"。而reducer 就是每块栏的"贴纸规矩":有的栏是"擦掉重写"(新值覆盖旧值,比如"当前用户");有的栏是"往下续写"(新值追加进列表,比如"聊天记录")。你在 Annotated[list, add] 里写的那个函数,就是在告诉 LangGraph:"这块栏用续写规矩"。今天全程记住这句就通。
L01

痛点:为什么每个字段要有自己的"合并规则"

🤔 痛点一个图里往往有多个节点,都会往 State 写东西。问题来了:如果两个节点在同一个超步里都往 messages 写,或者一个节点这一轮写、下一轮又写——这些值到底是后者覆盖前者,还是拼在一起?如果全 State 只有一种合并方式,就没法同时满足"当前用户要覆盖"和"聊天记录要累积"这两种完全相反的需求。
💡 本质:reducer = 字段级的合并函数LangGraph 的答案是:让每个字段自带一个合并函数(reducer)。签名固定是 (旧值, 新值) -> 合并后的值。写状态时,框架不是简单地 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
大白话没有 reducer 的世界里,"写"就是"盖章覆盖"——后来的把先来的抹掉。reducer 让你可以自定义"写"到底是覆盖、追加、求和、取最大值……随便你。它是 LangGraph 数据流最核心的一块拼图。
L02

Annotated[x, reducer] 到底是什么

Annotated 是 Python 标准库 typing 提供的"给类型贴便利贴"的工具。Annotated[list, reducer] 的意思是:这个字段的类型是 list,另外还附带一条元数据 reducer。类型检查器只看第一个 list,而 reducer 这条元数据"搭便车"存进了一个叫 __metadata__ 的属性里,等 LangGraph 来取。

📝 在 Python 里亲手验证一下 from typing import Annotated
t = Annotated[list, "我是便利贴"]
t.__metadata__('我是便利贴',)(一个元组,存着所有便利贴)
t.__origin__<class 'list'>(脱掉便利贴后的真实类型)
LangGraph 就是靠读 __metadata__ 这个元组,把你挂的 reducer 抠出来。
数据结构:Annotated[list, add_messages] 拆开长什么样 Annotated[ list , add_messages ] __origin__ = list 字段的"值类型" __metadata__ = (add_messages,) 元组,存所有"便利贴" LangGraph 读 __metadata__ 的最后一个元素 → 认出它是 reducer
图注:Annotated 把"值类型"和"便利贴(reducer)"打包在一起,二者分别落在 __origin__ 和 __metadata__。
💡 设计取舍①:为什么 reducer 挂在类型注解里,而不是单开一个配置字典? 朴素做法可能是让你另写 reducers = {"messages": add_messages} 一个字典。源码没这么做,而是塞进 Annotated。好处:① 字段定义和它的合并规则写在同一行,不会脱节——你删字段时不会忘删对应配置;② 类型检查器仍然把它当普通 list,IDE 补全、mypy 检查都不受影响(便利贴对类型系统"透明");③ 不需要额外的 API,用的是 Python 原生机制。代价是:初学者第一次见 Annotated[list, add_messages] 会懵——但看懂机制后就一目了然。
L03

_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
L04

_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_messagesoperator.add),就包成 BinaryOperatorAggregate 通道(D09 细讲)。
④ fallback = LastValue★兜底。字段什么便利贴都没贴(如 count: int),就给它 LastValue——"覆盖写"通道。这就是"不写 reducer 默认覆盖"的真相(L06 细讲)。
控制流:一个字段的类型注解 → 走哪条分支 → 得到哪种通道 字段类型注解 annotation 剥掉 Required/NotRequired ① 托管值?ManagedValue ② 是通道?你指定的通道 ③ 是 reducer?BinaryOperator… ④ 都不是LastValue(覆盖) 从左到右依次尝试,第一个命中的就用它;全落空才走 ④ 兜底
图注:四级分派是"短路"的——从①试到④,命中即返回。绝大多数字段落在③(写了 reducer)或④(没写=覆盖)。
L05

_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 会拆开这个通道内部。
⚠️ 边界/坑:reducer 参数个数写错,编译期就炸 如果你手滑写了个只接一个参数的函数 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
对比:②认的是"你已经给了现成通道",③认的是"你只给了个函数,帮你包成通道"。顺序上先试②再试③——因为你显式指定通道的意图,优先级高于"给个函数让框架猜"。
L06

没写 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 会因为收到多值而报错。

L07

边界: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__ 分支兜底跑起来。但会提醒你"更新可能不符合预期",让你自己决定要不要改。
💡 设计取舍②:为什么这里是"警告"而不是"报错"? 前面 L05 的 reducer 签名错是直接报错(raise),这里 schema 类型不对却只是警告(warn)。区别在于:reducer 签名错必然导致合并逻辑跑不通,是"死错",越早炸越好;而 schema 传成一个奇怪对象,框架还有 __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__)"
明日预告 · Day 08:今天的 reducer 里,最著名的就是 add_messages。明天钻进 graph/message.py,逐行拆它怎么做"按 ID 去重/更新、RemoveMessage 删除、REMOVE_ALL 清空"——一个远比 operator.add 复杂的工业级 reducer。
← Day 06 对话图心智 Day 08 · add_messages 深入 →