Skip to content

Topic 与 BinaryOperatorAggregate:累积式通道与归约器

源码版本1.2.9

职责

LastValue 只会覆盖,但实际 agent 状态里很多字段是要累积的:消息列表要追加、计数器要相加、待办要合并。LangGraph 给这类需求提供两条路:Topic(Topic class:23)是一个 pub/sub 风格的列表通道,BinaryOperatorAggregate(BinaryOperatorAggregate class:65)是一个接受任意二元运算符 (reducer) 的通用通道。两者都直接或间接继承自 BaseChannel(BaseChannel:19)。

BinaryOperatorAggregate 是用户写的 Annotated[list, add_messages] 这种声明背后的实现——StateGraph 编译时调用 _is_field_binop(_is_field_binop:1890-1908),把 Annotated 元数据里最后那个双参可调用对象当成 reducer,实例化成 BinaryOperatorAggregate(typ, reducer)。每次 update 都把新值用 operator 合并进当前值(update:123-144)。

Topic 更多用在内部通道:每个 Pregel 实例都自带一个名为 __pregel_tasksTopic(Send, accumulate=False)(TASKS channel:809)用来装 Send 扇出。它的特点是 accumulate=False 时每步清空——上一步的 Send 列表被 prepare_next_tasks 消费完就丢掉,下一步重新收集。

两者共同点是「update 接受多个值并按某种规则 fold」,这是它们和 LastValue 的根本区别——后者只接受一个值。

设计动机

  • reducer 是状态合并的通用接口:不同字段的合并语义千差万别——messages 要按 ID 去重追加、counter 要相加、tags 要并集。与其给每种都做一个 channel 类,不如提供一个通用通道接二元运算符 operator: (Value, Value) -> Value,让用户自己定义合并规则。add_messages(add_messages:60-66)就是这种 reducer:按消息 ID 合并,新 ID 追加、旧 ID 替换。
  • 类型擦除后还原:_strip_extras(_strip_extras:22-28)把 Annotated / Required / NotRequired 这些 typing 包装剥掉,再用 typ() 实例化默认值。collections.abc.Sequence 这种抽象类型不能直接实例化,所以在 typ fallback:83-88 里映射回 list / set / dict(typ instantiation:82-92)。
  • Topic.accumulate 双形态:accumulate=True 时是一个跨步累积的列表——update 把新值 extend 进去(update:77-85);accumulate=False 时每步先把旧值清空再 extend,变成「只看本步产生的值」。后者正是 TASKS 通道需要的:Send 一步内有效,被消费完就作废。
  • _flatten 支持混合写:Topic.update 接受 Value | list[Value] 混合序列(_flatten:15-20),所以节点可以一次写一条 Send 也可以一次写 [Send, Send],通道统一拍平。
  • Overwrite 旁路:reducer 通道还支持一种「直接覆盖」语义:Overwrite 数据类(Overwrite:980)或字典 {"__overwrite__": value}(OVERWRITE key:95)被 _get_overwrite 识别(_get_overwrite:31-51),在 update 里绕过 reducer 直接赋值。每步只允许出现一次 Overwrite,第二次会抛 InvalidUpdateError(overwrite dedup:132-138)。

关键文件

  • Topic class:23-41 — 类定义,泛型参数 Sequence[Value] / Value | list[Value] / list[Value],构造函数接 accumulate 开关。
  • _flatten:15-20 — 工具函数:把混合 Value | list[Value] 序列拍平成 Iterator[Value]
  • Topic.copy:56-61 — 重写 copy 直接复用 values 列表的副本,不走 checkpoint 序列化。
  • Topic.checkpoint / from_checkpoint:63-75 — checkpoint 直接返回 self.values 列表;from_checkpoint 还兼容旧版 tuple 形式。
  • Topic.update:77-85 — 核心逻辑:accumulate=False 先清空再 extend,accumulate=True 直接 extend;空 flat 序列返回 False
  • Topic.get / is_available:87-94 — 空列表抛 EmptyChannelError;is_available 重写成 bool(self.values)
  • Binop class:65-92 — 类定义 + __init__:接 operator 二元可调用对象,用 typ() 初始化存值。
  • _get_overwrite:31-51 — 识别三种 Overwrite 表示形式(数据类、sentinel-keyed dict、JSON 还原后的 dict)。
  • _operators_equal:54-62 — lambda 比较特殊处理:同名 lambda 视为相等,避免每次重编译图都把通道判为变更。
  • Binop.update:123-144 — reducer 的真正入口:首个值用来初始化,后续值逐个调 operator(self.value, value),Overwrite 旁路例外。
  • _is_field_binop:1890-1908 — 编译期识别 Annotated[..., reducer] 元数据里的双参可调用对象,实例化成 BinaryOperatorAggregate
  • TASKS channel:805-809Pregel.__init____pregel_tasks 通道硬编码 Topic(Send, accumulate=False),不允许用户覆盖。
  • add_messages:60-66 — 典型 reducer:按消息 ID 合并两个列表,新 ID 追加、旧 ID 替换。

数据流

BinaryOperatorAggregate.update 是 reducer 通道的心脏(update:123-144):

python
def update(self, values: Sequence[Value]) -> bool:
    if not values:
        return False
    if self.value is MISSING:
        self.value = values[0]
        values = values[1:]
    seen_overwrite: bool = False
    for value in values:
        is_overwrite, overwrite_value = _get_overwrite(value)
        if is_overwrite:
            if seen_overwrite:
                msg = create_error_message(
                    message="Can receive only one Overwrite value per super-step.",
                    error_code=ErrorCode.INVALID_CONCURRENT_GRAPH_UPDATE,
                )
                raise InvalidUpdateError(msg)
            self.value = overwrite_value
            seen_overwrite = True
            continue
        if not seen_overwrite:
            self.value = self.operator(self.value, value)
    return True

读这段要注意两个分支:第一个值若通道是空的,直接当初始值(跳过一次 reducer 调用,因为 reducer 需要两个输入);之后每个值若不是 Overwriteself.value = operator(self.value, value)。如果出现 Overwrite,之后所有非 overwrite 值都被丢弃——这保证「覆盖」语义不会被后续 reducer 又改回去。

StateGraph 怎么把 Annotated[list[AnyMessage], add_messages] 编译成 BinaryOperatorAggregate,看 _is_field_binop:1890-1908:

python
def _is_field_binop(typ: type[Any]) -> BinaryOperatorAggregate | None:
    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
            ):
                return BinaryOperatorAggregate(typ, meta[-1])
            else:
                raise ValueError(
                    f"Invalid reducer signature. Expected (a, b) -> c. Got {sig}"
                )
    return None

它用 inspect.signature 校验最后一个 meta 必须是双参可调用对象,签名不匹配直接抛错。这就是为什么 add_messages(left, right) 这种两参函数能当 reducer,而单参函数会被拒。

Topic.update 走的是另一条路径,逻辑更短(update:77-85):

python
def update(self, values: Sequence[Value | list[Value]]) -> bool:
    updated = False
    if not self.accumulate:
        updated = bool(self.values)
        self.values = list[Value]()
    if flat_values := tuple(_flatten(values)):
        updated = True
        self.values.extend(flat_values)
    return updated

注意 accumulate=False 时的返回值:如果旧值非空,即使新值序列空也返回 True(因为清空本身就是一种改动),Pregel 会更新 channel_versions。这和 LastValue 空序列返回 False 的行为相反,反映了不同通道对「空写」的不同语义。

reducer 通道在一次 Pregel 超步里的位置:

边界与失败

  • Overwrite 三种识别形式并存:overwrite forms:44-50 同时接受 Overwrite 数据类、{"__overwrite__": value} 字典、以及 JSON 还原后的 {"type": "__overwrite__", "value": ...} 字典。第三个分支是为了让经过 LangGraph API server 的 orjson 序列化往返后语义不丢。
  • 每步只能有一个 Overwrite:overwrite dedup:132-138 第二次出现直接抛 InvalidUpdateError,因为「覆盖」语义在一步内多次出现是矛盾的。
  • Overwrite 之后的非 overwrite 值被静默丢弃:discard after overwrite:142-143 if not seen_overwrite 决定是否调 reducer,一旦 seen_overwrite=True,后续值既不入 reducer 也不报错。这是有意的——「覆盖」意味着不想保留之前的合并结果,也不需要后续值。
  • lambda reducer 等价性比较:_operators_equal:54-62 把所有 lambda 视为相等,因为 Python 的 lambda 每次 <lambda> 都是新对象,按身份比会导致 __eq__ 永远不相等,影响图缓存。
  • Topic(accumulate=False) 空写返回 True:clear returns True:79-80——只要旧值非空,清空就视为改动。这会让 apply_writes 更新 channel_versions 从而触发依赖该通道的节点。如果用户自己用 Topic,要意识到这一点:它不像 LastValue 那样静默处理空写。
  • Topic 旧 checkpoint 兼容:tuple compatibility:70-74 检查 isinstance(checkpoint, tuple),旧版本的 checkpoint 存的是 (typ, values) 元组,新版改成直接存 list。这段是为了不破坏旧存档。
  • reducer 签名强制双参:reducer signature check:1894-1907inspect.signature 数 positional 参数,不是两个就报错。带 *args 这种可变参数的函数不会被接受——这是把「reducer 必须严格 (a, b) -> c」钉死。
  • add_messages 是个工厂:_add_messages_wrapper:41-57@_add_messages_wrapper 装饰,实际函数可以在 add_messages(left, right) 直接调用,也可以 add_messages(format="langchain-openai") 返回 partial。这种用法是 reducer 模式的常见扩展点——给 reducer 加运行时参数。

小结

BinaryOperatorAggregate 是「带 reducer 的字段」的统一实现,Topic 是「列表累积」的两种形态:跨步累积或每步清空。它们和 /channels/last-value 构成了 LangGraph 通道系统的三件套,接口都来自 /channels/base-channeladd_messages 这种 reducer 在 /graph/state-graph 里被 Annotated[..., reducer] 声明,运行时由 /pregel/algoapply_writes 调进 updateTopic(accumulate=False) 作为 __pregel_tasks 通道是 Send 扇出的载体,具体见 /subgraph/send-command

对照官方资料:LangGraph 文档 · README