Topic 与 BinaryOperatorAggregate:累积式通道与归约器
职责
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_tasks 的 Topic(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-809—Pregel.__init__给__pregel_tasks通道硬编码Topic(Send, accumulate=False),不允许用户覆盖。add_messages:60-66— 典型 reducer:按消息 ID 合并两个列表,新 ID 追加、旧 ID 替换。
数据流
BinaryOperatorAggregate.update 是 reducer 通道的心脏(update:123-144):
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 需要两个输入);之后每个值若不是 Overwrite 就 self.value = operator(self.value, value)。如果出现 Overwrite,之后所有非 overwrite 值都被丢弃——这保证「覆盖」语义不会被后续 reducer 又改回去。
StateGraph 怎么把 Annotated[list[AnyMessage], add_messages] 编译成 BinaryOperatorAggregate,看 _is_field_binop:1890-1908:
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):
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-143if 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-1907用inspect.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-channel。add_messages 这种 reducer 在 /graph/state-graph 里被 Annotated[..., reducer] 声明,运行时由 /pregel/algo 的 apply_writes 调进 update。Topic(accumulate=False) 作为 __pregel_tasks 通道是 Send 扇出的载体,具体见 /subgraph/send-command。
对照官方资料:LangGraph 文档 · README。