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"],怎么统一处理? 靠模块级的小生成器 _flattenchannels/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"] → 吐出 bc
else: yield value单值原样吐出。这样无论上游投单值还是列表,出口都是一串扁平的元素
只摊平"一层"边界要点:它只拆一层。如果你投 [["x"]](列表套列表),拆一层后是 ["x"],这个内层列表会作为一个元素进 values——Topic 不做递归摊平。
数据结构:_flatten 把"混投"抹平成扁平元素流 收到的一批写入 "a"(单值) ["b","c"](列表) _flatten 摊平一层 扁平元素流a, b, c→ values.extend() UpdateType = Value | list[Value]:单值原样吐、列表拆一层 → 出口统一是扁平元素
正因写入允许"单值或列表混投",才需要 _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 决定"是否先清空" accumulate=True(跨步累积) 步1: [a] 步2: [a,b] 步3: [a,b,c] accumulate=False(每步清空) 步1: [a] 步2: [b] 步3: [c] 同样每步各投一个值: 累积模式越攒越长; 清空模式只留当步的 区别只在 update 开头那一句 "if not accumulate: 清空"
同样的投递序列,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 的场景很少,但理解它,能让你看懂"为什么并行分支的消息能自动汇聚、且处理完就清空"。

维度TopicBinaryOperatorAggregate(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 参数如何控制"一步能不能多写",以及"不持久化"到底为谁而设计。
← Day 29 累加/Delta Day 31 · 不持久化通道 →