Day 30 / 共 60 天 · 阶段 5 Channels 通道
Topic:发布订阅式的列表累积通道
昨天的 BinaryOperatorAggregate 靠一个算子把新值"折"进旧值。今天的 Topic 换了个思路——它是一个发布订阅(PubSub)主题:多个生产者往主题里"投"消息,通道只管把它们摊平、按顺序排进一个列表。更关键的是它有一个 accumulate 开关:开着就跨步累积、关着就每步清空——这让它既能当"消息队列"又能当"一次性收件箱"。全程只有 95 行,是最短的通道之一,却藏着 _flatten、"每步清空"这两个精巧设计。
📍 你在 60 天里的位置(阶段 5:Channels 通道 D27-32)
D27 通道抽象→
D28 LastValue→
D29 累加/Delta→
D30 Topic→
D31 不持久化→
D32 同步屏障
💡 类比先兜住(延续"公司信箱"世界观)
如果 LastValue 是"只贴最新一张便利贴"、binop 是"往留言簿后面接着写",那 Topic 就是一块 公告板 + 一个收件箱:谁都能往上贴纸条(发布),贴上去就按到达顺序排成一列。
accumulate=True 时公告板永不擦除、越贴越多;accumulate=False 时它是"当天收件箱"——每过一天(一个超步)自动清空,只留今天新收到的。它和 binop 最大的不同:binop 需要你给个算子来决定"怎么合并",Topic 不需要,它天生就是"排队进列表"。L01
为什么已有 binop 还要 Topic?
🤔 痛点:
Annotated[list, operator.add] 不也能攒列表吗?为什么还专门造个 Topic?
因为两者的"心智模型"不同。binop 是"reducer 累加",你得懂算子;Topic 是"发布订阅",语义更贴近"多个节点各自往一个话题投消息"。而且 Topic 多了一个 binop 没有的能力——每步自动清空。先看类头(channels/topic.py:23-32):# topic.py:23
class Topic(
Generic[Value],
BaseChannel[Sequence[Value], Value | list[Value], list[Value]],
):
"""A configurable PubSub Topic.
Args:
typ: The type of the value stored in the channel.
accumulate: Whether to accumulate values across steps.
If `False`, the channel will be emptied after each step.
"""
__slots__ = ("values", "accumulate")
三泛型不再全相同Day 27 讲过通道有三个类型参数 [Value 存/读, Update 写]。这里读出来是 Sequence[Value](一个列表)、但写进来的是 Value | list[Value]——可以投单个值、也可以投一个列表。这就是 _flatten 存在的理由(L03)。__slots__ = ("values","accumulate")两个字段:values 是那个不断排队的列表,accumulate 是"跨步累积/每步清空"的开关。"A configurable PubSub Topic"docstring 直接点名它是发布订阅主题——多生产者投递、按到达顺序进列表,不去重、不合并。💡 本质:Topic = "多写入者共享的一个有序列表"binop 的核心是"算子",Topic 的核心是"顺序追加 + 可清空"。当你想表达"这一步产生的一堆待办/事件,交给下游一次性处理,处理完就丢"时,
Topic(accumulate=False) 比 binop 更贴切——binop 没法"处理完自动清空"。L02
状态只有一个列表 values
构造函数简单到极点(channels/topic.py:36-41):
# topic.py:36
def __init__(self, typ: type[Value], accumulate: bool = False) -> None:
super().__init__(typ)
# attrs
self.accumulate = accumulate
# state
self.values = list[Value]() # 起点永远是空列表 []
accumulate 默认 False关键默认值:默认每步清空。这符合"事件/待办"的直觉——不特意开累积,就当"一次性收件箱"用。想跨步攒着,才显式传 accumulate=True。self.values = list()初值直接是 [],不像 binop 还要用 typ() 去猜零值——因为 Topic 的容器永远是 list,typ 只描述"列表里装的元素类型",不影响容器本身。📝 真实值例子
Topic(str, accumulate=True) → values=[],装的是字符串。投 "a" 后 values=["a"],再投 ["b","c"] → values=["a","b","c"]。注意第二次投的是列表,也被摊平进去了(下一讲揭晓)。🅰 设计取舍①:为什么 Topic 的容器写死是 list,而 binop 要动态推断类型?
因为两者的定位不同。binop 要当"通用 reducer 通道",得支持 int 求和、set 求并、dict 合并等各种类型,所以必须用
typ() 动态求零值(Day 29)。而 Topic 的语义本身就是"排队进一个有序列表"——顺序是它的灵魂,set/dict 这些无序或去重的容器根本不符合"发布订阅按到达顺序"的语义。把容器写死成 list 反而让代码更简单、语义更聚焦。这是"专用工具做减法"的取舍:放弃通用性,换来零配置和语义清晰。L03
_flatten:把"单值和列表混投"抹平
🤔 一个节点投单个值
"a"、另一个节点投一批 ["b","c"],怎么统一处理?
靠模块级的小生成器 _flatten(channels/topic.py:15-20):# topic.py:15
def _flatten(values: Sequence[Value | list[Value]]) -> Iterator[Value]:
for value in values:
if isinstance(value, list):
yield from value # 是列表 → 逐个吐出(摊平一层)
else:
yield value # 是单值 → 原样吐出
for value in values引擎把这一步收到的所有写入攒成一个序列传进来。序列里每一项可能是单值、也可能是列表(因为 UpdateType 是 Value | list[Value])。isinstance(value, list): yield from value遇到列表就拆开一层,把里面的元素一个个吐出去。["b","c"] → 吐出 b、c。else: yield value单值原样吐出。这样无论上游投单值还是列表,出口都是一串扁平的元素。只摊平"一层"边界要点:它只拆一层。如果你投 [["x"]](列表套列表),拆一层后是 ["x"],这个内层列表会作为一个元素进 values——Topic 不做递归摊平。正因写入允许"单值或列表混投",才需要 _flatten 在进 values 前抹平差异
💡 为什么用生成器(yield)而不是先拼一个大列表?生成器惰性产出、不额外占内存:
tuple(_flatten(values))(见 L05)一次遍历就地物化,中途不需要构造多个临时列表。对"一步可能收到成千上万条消息"的场景,这点内存友好很实在。这也是 Python 里处理"把嵌套结构流式摊平"的地道写法。L04
accumulate 开关:跨步累积还是每步清空
Topic 最有特色的能力藏在 update 的头几行(channels/topic.py:77-85):
# topic.py:77
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
if not self.accumulate不累积时,每次 update 一开始就把 values 重置为空列表——上一步的内容全丢。相当于"收件箱每天早上清空"。updated = bool(self.values)细节:清空前先看"本来有没有东西"。有 → 这次算"发生了变化"(哪怕这步没投新值,"从有到无"也是变化,下游要感知)。accumulate=True 时跳过清空累积模式下这个 if 整段不执行,values 保留上一步内容,新值直接 extend 追加到末尾——公告板越贴越长。同样的投递序列,accumulate 决定 values 是累积还是每步归零
🅰 设计取舍②:为什么"每步清空"要做成通道的内建能力,而不是让节点自己清?
如果让节点手动清空,就得在"下游读完之后、下一步写入之前"找一个准确的时机去清——这在并行、多节点的图里几乎不可能协调对。把"每步清空"下沉到通道的 update 生命周期里(每个超步 update 被调一次,开头就清),时机天然正确、无需协调。这正是 Day 27"通道负责状态语义、节点只管读写"分工的体现:清空时机是状态语义的一部分,理应由通道掌管。代价是通道类型变多了(要区分累积/不累积),但换来节点代码零负担。
L05
update 全流程 + get 的空判定
把 update 和 get 连起来看完整取用链路(channels/topic.py:77-91):
# topic.py:77
def update(self, values):
updated = False
if not self.accumulate:
updated = bool(self.values)
self.values = list[Value]()
if flat_values := tuple(_flatten(values)): # 摊平后物化成 tuple
updated = True
self.values.extend(flat_values)
return updated # 告诉引擎"这步这个通道变没变"
# topic.py:87
def get(self) -> Sequence[Value]:
if self.values:
return list(self.values) # 返回一份拷贝,防外部改内部
else:
raise EmptyChannelError # 空列表 = 通道"没值"
flat_values := tuple(...)海象运算符:摊平并物化成 tuple,同时判断是否非空。非空才 extend、才标记 updated=True。空投递(没实际值)就不动列表。return updatedupdate 的返回值是给引擎看的"这一步这个通道到底变没变"——决定下游订阅它的节点要不要被触发(Day 27 契约)。累积模式下没新值就返回 False;清空模式下只要之前有东西就返回 True(因为清空本身是变化)。get: if self.values读取时:列表非空就返回一份 list(...) 拷贝;空列表则抛 EmptyChannelError——即"空列表"被当作"通道无值",读它的节点会被跳过。返回拷贝而非本体防御性设计:返回 list(self.values) 而不是 self.values,避免下游节点不小心 .append() 直接改到通道内部状态。⚠ 边界:空列表 == 无值,别指望读到
[]
注意 get 里 if self.values 对空列表为 False → 抛 EmptyChannelError。也就是说你永远读不到一个"空的 Topic 列表 []"——空就等于没有。对 accumulate=False 的 Topic 尤其要留心:某一步没有任何生产者投递时,它是"空 → 无值",订阅它的下游节点这一步不会被触发,而不是收到一个 []。想区分"空批"和"没批"?Topic 帮不了你,语义上它俩就是一回事。L06
checkpoint:整张表存、带向后兼容
Topic 是要持久化的(对比明天的 EphemeralValue 就不存)。它的存/取很直白(channels/topic.py:63-75):
# topic.py:63
def checkpoint(self) -> list[Value]:
return self.values # 直接把整个列表交出去存档
# topic.py:66
def from_checkpoint(self, checkpoint: list[Value]) -> Self:
empty = self.__class__(self.typ, self.accumulate)
empty.key = self.key
if checkpoint is not MISSING:
if isinstance(checkpoint, tuple):
# backwards compatibility
empty.values = checkpoint[1] # ← 旧版本存的是 (guard, values) 二元组
else:
empty.values = checkpoint # 新版本直接就是 list
return empty
checkpoint 返回整个 values和 binop 不同(Day 29 的 DeltaChannel 只存哨兵),Topic 老老实实把整张列表存进 checkpoint。简单直接,代价是长列表存档会大。from_checkpoint 新建再灌值恢复时造一个新通道实例(带上原来的 typ、accumulate、key),再把存档的列表灌回 values。这是所有通道 from_checkpoint 的统一套路。isinstance(checkpoint, tuple)向后兼容分支:早期版本的 Topic 存的是二元组(历史上带过别的字段),所以读到 tuple 就取 [1] 当列表。新版本直接存 list。这段代码保证老存档能被新代码读起来。💡 为什么要留 tuple 兼容分支?因为 checkpoint 是持久化到数据库的——用户升级了 LangGraph 版本,但数据库里还躺着旧格式的存档。如果不兼容,升级即"丢历史对话"。这几行
isinstance(checkpoint, tuple) 就是"格式演进时不抛弃老数据"的承诺。Day 39 讲序列化时你会看到,这种"多格式并存的读兼容"在整个持久化层反复出现。L07
Topic vs BinaryOperatorAggregate + 小结
👶 小白:Topic 和昨天的 Annotated[list, operator.add] 到底该用哪个?
👨🏫 老师:日常 State 字段里你几乎总是用 Annotated[list, add](binop),因为它是 reducer 的标准写法、还能配 add_messages。Topic 更多是 LangGraph 内部机制在用——比如 Send 动态扇出(Day 16)时用一个 Topic 通道收集要投递的任务。你自己直接实例化 Topic 的场景很少,但理解它,能让你看懂"为什么并行分支的消息能自动汇聚、且处理完就清空"。
| 维度 | Topic | BinaryOperatorAggregate(binop) |
|---|---|---|
| 合并方式 | 固定:摊平后顺序追加进 list | 可配:用你给的算子 operator.add 等折叠 |
| 容器类型 | 永远 list | 任意(int/set/dict/list…由 typ 推断) |
| 每步清空 | ✅ accumulate=False 支持 | ❌ 只能靠 Overwrite 手动重置 |
| 投递形态 | 单值或列表都行(_flatten 抹平) | 按算子语义,通常单值 |
| 存档 | 整张列表存进 checkpoint | 同左(DeltaChannel 例外,只存哨兵) |
| 心智模型 | 发布订阅主题 / 收件箱 | reducer 累加器 |
一句话记法:需要"每步自动清空"或"发布订阅"语义 → Topic;需要任意类型的 reducer 累加 → binop。
🧠 今日小结自测
- Topic 和 binop 的核心区别是什么?(Topic 固定"排队进 list + 可每步清空",binop 靠算子折叠任意类型)
- _flatten 解决什么问题?(把"单值和列表混投"摊平成扁平元素流,只摊一层,用生成器省内存)
- accumulate=False 时 update 做了什么特殊动作?(开头把 values 重置为空列表,实现每步清空)
- 为什么"每步清空"由通道管而不是节点管?(时机是状态语义的一部分,通道生命周期天然对时机,节点无从协调)
- Topic 读到空列表会怎样?(抛 EmptyChannelError,空==无值,下游不被触发)
- from_checkpoint 里的 tuple 兼容分支为何存在?(读起旧格式存档,格式演进不丢历史数据)
✋ 10 分钟动手
# 1. 通读 Topic 全文(95 行)
sed -n '1,95p' libs/langgraph/langgraph/channels/topic.py
# 2. 重点对照 _flatten 与 update 的清空分支
sed -n '15,20p;77,85p' libs/langgraph/langgraph/channels/topic.py
# 3. 亲手验证两种模式的差别
python - <<'PY'
from langgraph.channels.topic import Topic
acc = Topic(str, accumulate=True)
acc.update(["a"]); acc.update([["b","c"]]) # 单值/列表混投
print("累积:", acc.get()) # ['a','b','c']
tmp = Topic(str, accumulate=False)
tmp.update(["x"]); print("步1:", tmp.get()) # ['x']
changed = tmp.update([]) # 空投递
print("步2 changed?", changed) # True(因为清空了旧的)
PY
🔮 明日预告 · Day 31 EphemeralValue & UntrackedValueTopic 会老实存档。明天看两个"故意不完整持久化"的通道:
EphemeralValue(只活一步、下一步自动清)和 UntrackedValue(内存里有值,但 checkpoint() 永远返回 MISSING、根本不进存档)。它们把 Day 27 埋的"内存态 ≠ 存档态"这个伏笔玩到极致——你会看到 guard 参数如何控制"一步能不能多写",以及"不持久化"到底为谁而设计。