Skip to content

BaseChannel 抽象:通道即状态原语

源码版本1.2.9

职责

在 LangGraph 的运行时里,「通道 (channel)」是状态 (state) 的最小单元。你写一个 StateGraph 编译出来的 Pregel 实例,内部并不直接持有用户的 dict,而是持有一组 BaseChannel 实例——每个通道一个 key,通道自己负责存值、合并写、对外暴露读、还能把当前状态序列化成检查点 (checkpoint) 形式。BaseChannel 就是所有通道类共用的抽象基类 (BaseChannel:19)。

它本身不存数据(__slots__ 里只有 keytyp,见 __slots__:22-26)。它只规定四个抽象方法子类必须实现:ValueType / UpdateType 两个类型属性,以及 update / get / from_checkpoint 三个动作 (abstract methods:60-99)。另外有 checkpoint / is_available / consume / finish / copy 五个带默认实现的方法,默认实现要么走 get() 兜底,要么返回 False 当 no-op。

换句话说,BaseChannel 把「状态怎么存、怎么变、怎么留痕」这三件事一次性钉成接口,Pregel 引擎只需要调 update / get / checkpoint 这些方法,不需要知道具体是 LastValue 还是 Topic 还是 BinaryOperatorAggregate 在背后怎么做合并。这就是为什么引擎和状态形状可以彻底解耦。

设计动机

为什么把状态做成「通道」这种对象,而不是直接拿一个 dict 在节点间传?

  • 统一写合并语义:同一个超步 (superstep) 内,多个节点可能并发往同一个 key 写。dict 没法表达「list 要追加,标量要覆盖,计数器要相加」这几套不同的合并规则;通道用一个 update(values: Sequence[Update]) -> bool 把策略收口,子类自己决定怎么 fold。
  • 不可变步内视图:节点读到的是步开始时的快照,本步别的节点写的值要等下一步才可见。get 是只读当前值,update 只在步边界由 apply_writes 统一调 (apply_writes:317-323),从机制上杜绝了竞争。
  • 可序列化:每个通道都能 checkpoint() 出一个可序列化值(checkpoint:49-58)和 from_checkpoint 重建自己 (from_checkpoint:60-65)。BaseCheckpointSaver 族就是靠这俩接口把整个线程状态写进 SQLite/Postgres 的。
  • 生命周期钩子:consumefinish 两个方法(consume / finish:101-121)给特殊通道(比如 Topic(accumulate=False)LastValueAfterFinish)一个在步边界或运行结尾改自己状态的机会,默认 no-op 表示大多数通道不需要这层语义。
  • 类型自描述:ValueType / UpdateType 两个属性(ValueType / UpdateType:28-36)让编译期和运行期都能反查「这个通道存什么、接受什么写」,StateGraph_is_field_channel:1862-1887 里就用它来识别 Annotated[..., SomeChannel] 这种声明。

关键文件

  • BaseChannel class:19-26 — 抽象基类定义,泛型参数 Value / Update / Checkpoint,__slots__ 只声明 keytyp
  • ValueType / UpdateType:28-36 — 两个抽象 property,子类用来声明存值类型和写值类型。
  • checkpoint:49-58 — 默认实现:直接返回 self.get(),空通道返回 MISSING
  • from_checkpoint:60-65 — 抽象方法:从一个 checkpoint 值构造一个等价的新通道,子类必须实现。
  • get:69-73 — 抽象读接口,空通道抛 EmptyChannelError
  • is_available:75-85 — 默认实现走 get() + 捕获 EmptyChannelError,子类一般重写成更快的判断。
  • update:89-99 — 抽象写接口,Pregel 在每个超步末尾调一次,顺序任意;返回 True 表示通道被改动。
  • consume / finish:101-121 — 生命周期钩子,默认 no-op 返回 False
  • apply_writes:315-345 — Pregel 把本步所有任务写回合并到通道的唯一入口,update / consume / finish 都在这段里被统一调度。
  • _is_field_channel:1862-1887StateGraph 编译时用 isinstance(item, BaseChannel)Annotated[..., channel] 翻译成 channel 实例。

数据流

通道的写发生在 apply_writes 里,Pregel 把本步所有任务的 writes 按通道聚合,然后调 update(apply_writes:317-323):

python
# Apply writes to channels
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():
    if chan in channels:
        if channels[chan].update(vals) and next_version is not None:
            checkpoint["channel_versions"][chan] = next_version
            # unavailable channels can't trigger tasks, so don't add them
            if channels[chan].is_available():
                updated_channels.add(chan)

# Channels that weren't updated in this step are notified of a new step
if bump_step:
    for chan in channels:
        if channels[chan].is_available() and chan not in updated_channels:
            if channels[chan].update(EMPTY_SEQ) and next_version is not None:
                checkpoint["channel_versions"][chan] = next_version
                if channels[chan].is_available():
                    updated_channels.add(chan)

注意第二段:本步没被任何任务写的通道也会收到一次 update(EMPTY_SEQ) 调用。这是 BaseChannel 接口的隐含契约——子类的 update 必须能接受空序列,通常返回 False 表示「没改动」。Topic(accumulate=False) 这种「每步清空」语义就是靠这个空调用触发的:空 update 让它把上一步的值清掉,从而在下一步不可见。

整个通道生命周期在 Pregel 一轮 tick 里的位置如下:

边界与失败

  • EmptyChannelError:get() 在通道从未被写过时抛 (get:70-73)。is_available 默认实现就靠 catch 这个异常判断可用性,子类若能更便宜地判断应重写。
  • MISSING 哨兵:checkpoint()get()EmptyChannelError 时返回 langgraph._internal._typing.MISSING 而不是 None (checkpoint fallback:55-58),这样 None 可以作为合法存值出现,不与「空」混淆。
  • update 返回 False:返回 False 表示通道没改动,Pregel 就不会更新 channel_versions,后续 prepare_next_tasks 也就不会因为这条通道触发任何节点 (update call:319);这是 Pregel 收敛判断的基础。
  • 空序列调用:apply_writes 对未变更通道也会调 update(EMPTY_SEQ)(EMPTY_SEQ update:329)。BaseChannel.update 的契约是空输入返回 False,但 Topic(accumulate=False) 会在此时清空自己并可能返回 True,所以子类不能假设「空输入 = 啥也不做」。
  • consumefinish 默认 no-op(consume / finish:101-121):返回 False 意味着大多数通道不参与生命周期管理。apply_writesbump_step 末尾会调 finish(finish call:338),只有 LastValueAfterFinish 这类特殊通道会用它做「等运行结束才可见」的语义。
  • copy 默认走 checkpoint:BaseChannel.copy 默认 self.from_checkpoint(self.checkpoint())(copy:40-47)。子类如果 checkpoint 代价高(比如要深拷贝大 list),应该重写 copy 提供更便宜的实现,Topic.copy 就是这么做的。

小结

BaseChannel 把「状态怎么存、怎么合并、怎么序列化」一次性钉成接口,Pregel 引擎只需要按 update → is_available → get → checkpoint → finish 这套固定顺序调,就能跑在任意通道实现上。要理解具体实现,看 /channels/last-value/channels/topic-binop;通道被 Pregel 读写的过程见 /pregel/algo;检查点如何把 checkpoint() 的返回值持久化见 /checkpoint/base-saver

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