- 发布日期
第 04 讲 · Channel 家族详解(上):四个最常用通道
LastValue / Topic / BinaryOperatorAggregate / EphemeralValue 四大通道详解
学习目标
- 掌握
LastValue/Topic/BinaryOperatorAggregate/EphemeralValue的语义差异。 - 理解"为什么不写 reducer 时并发写会报错"。
- 学会按场景选对通道。
内置通道清单
libs/langgraph/langgraph/channels/__init__.py 导出的家族:
__all__ = (
# base
"BaseChannel",
# value types
"AnyValue",
"LastValue",
"LastValueAfterFinish",
"UntrackedValue",
"EphemeralValue",
"BinaryOperatorAggregate",
"DeltaChannel",
"NamedBarrierValue",
"NamedBarrierValueAfterFinish",
# topics
"Topic",
)
本讲讲前四个高频通道,下一讲讲屏障与特殊通道,第 6 讲单独讲 DeltaChannel。
1. LastValue:默认通道,"只能一个值"
channels/last_value.py。这是 StateGraph 里没有写 reducer 的字段默认用的通道。
def update(self, values: Sequence[Value]) -> bool:
if len(values) == 0:
return False
if len(values) != 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]
return True
关键语义:一个超步内最多接收一个值。如果同一步有两个节点都写了这个字段, update 收到长度 2 的序列 → 抛 InvalidUpdateError。
这就是新手最常踩的坑——
InvalidUpdateError: Can receive only one value per step. Use an Annotated key to handle multiple values.
原因:两个并行节点写了同一个无 reducer 的字段。解法:给该字段加 reducer(用 Annotated), 让通道知道"多个值怎么合",比如换成 BinaryOperatorAggregate 或 Topic。
2. BinaryOperatorAggregate:用二元算子折叠多值
channels/binop.py。这是 Annotated[X, operator.add] 背后的通道。
class BinaryOperatorAggregate(Generic[Value], BaseChannel[Value, Value, Value]):
"""Stores the result of applying a binary operator to the current value and each new value.
```python
import operator
total = Channels.BinaryOperatorAggregate(int, operator.add)
```
"""
- 给定二元算子
op(acc, new),把本步所有写依次折叠进当前值。 operator.add对int是累加、对list是拼接、对dict看你传什么算子。- 初始值由类型推断:
int()→0、list()→[]、dict()→{}(见 66–78 行的类型规整逻辑)。
适用:累加计数、列表追加、集合并——任何"多个并行结果要合并成一个"的字段。
add_messages(第 7 讲)本质也是一个 reducer,但它有更复杂的消息去重/更新逻辑, 用的是普通函数 reducer 而非纯二元算子。
3. Topic:发布订阅式的"消息队列"通道
channels/topic.py。和前两者最大的不同:它天生存多个值(一个列表)。
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
两个关键特性:
- 可接收单值或列表:
_flatten会把列表展平,所以你写[a, b]和分别写a、b等价。 accumulate开关:accumulate=False(默认):每步开头先清空——本步产出只给下一步消费一次,像消息队列。accumulate=True:跨步累积——像一个不断增长的日志。
适用:fan-out / fan-in 的中转、把一批结果传给下游一次性处理。 (内部 TASKS 通道就是一个 Topic,用来承载 Send,第 19 讲细讲。)
4. EphemeralValue:只活一步的"瞬时值"
channels/ephemeral_value.py。语义:存最后一个值,但不跨步保留—— 写入后只在下一步可读,再下一步若无人续写就清空。
适用:步间一次性信号,比如把"刚生成的中间产物"传给紧邻的下一个节点,但不想污染长期状态。 START 通道在函数式 API 里就用 EphemeralValue 承载输入(第 19 讲)。
四通道速查表
| 通道 | 一步多值? | 跨步保留? | 典型 reducer 写法 | 场景 |
|---|---|---|---|---|
LastValue | 否(>1 报错) | 是 | 默认(无 Annotated) | 单写者字段、当前问题/答案 |
BinaryOperatorAggregate | 是(折叠) | 是 | Annotated[T, op] | 计数、列表追加、聚合 |
Topic | 是(成列表) | 看 accumulate | 内部使用为主 | fan-out 中转、消息队列 |
EphemeralValue | 否 | 否 | 内部/高级用法 | 一次性步间信号 |
使用场景实战:并行扇出的结果聚合
假设你让 3 个分析节点并行跑,各产出一个结论,要汇总:
from typing import Annotated
from typing_extensions import TypedDict
import operator
class State(TypedDict):
question: str # LastValue:只有一个问题
findings: Annotated[list, operator.add] # BinaryOperatorAggregate:3 个结论拼起来
question用默认LastValue:单一来源,简单。findings必须加 reducer,否则 3 个并行节点同步写会触发InvalidUpdateError。
这就是"选错通道 = 运行时报错"的直接因果。
动手实验
- 定义
class S(TypedDict): x: int(无 reducer),让两个节点在同一步都return {"x": 1}, 观察InvalidUpdateError。 - 把
x改成Annotated[int, operator.add],再跑一次,观察两个值被累加成 2。 - 把字段换成
Annotated[list, operator.add],体会列表拼接。
阅读作业
- 精读
channels/binop.py的类型规整逻辑(66–78 行)与update实现。 - 阅读
channels/ephemeral_value.py全文,对比它与LastValue的update差异。
小结
LastValue单值、并发写报错;BinaryOperatorAggregate折叠多值;Topic天生多值可累积;EphemeralValue只活一步。- "并发写报错"= 用了
LastValue又有多个写者,解法是加 reducer。 - 选通道就是选"多个并行写如何合并"的策略。
下一讲:屏障通道与 AfterFinish 系列——多源同步的关键。