- 发布日期
第 13 讲 · 节点读写机制:PregelNode / ChannelRead / ChannelWrite
学习目标
- 理解节点是怎么"读输入"和"写输出"的(通过注入的回调,而非直接碰通道)。
- 看懂
task.writes缓冲的全链路。 - 理解
local_read(fresh=True)这个"步内可见性"的唯一特例。
PregelNode:Actor 的容器
libs/langgraph/langgraph/pregel/_read.py,类 PregelNode(97 行起)。 第 8 讲编译时已为每个节点生成了它。它打包了:
triggers:订阅的触发通道(如branch:to:X)。channels:读哪些状态通道。bound:你写的函数本体。writers:一组ChannelWrite,负责把输出写回通道。
执行时,PregelNode.node 把它们拼成一个可运行序列:
RunnableSeq(bound, *writers)
即:先跑你的函数,再把结果交给 writers 写通道。
写:ChannelWrite → CONFIG_KEY_SEND → task.writes
pregel/_write.py,类 ChannelWrite(46 行起)。它不直接改通道,而是调用 config 里注入的 CONFIG_KEY_SEND 回调:
flowchart TD
A[节点 bound 返回 dict] --> B[ChannelWrite._write]
B --> C["_assemble_writes: dict → (channel, value) 列表"]
C --> D["config CONFIG_KEY_SEND(writes)"]
D --> E["task.writes.extend(...) 缓冲!"]
E --> F[runner.commit → loop.put_writes]
F --> G[checkpoint_pending_writes 持久化]
G --> H[after_tick → apply_writes 真正落通道]
_assemble_writes(172–192 行)有个关键分支:如果值是 Send,它会被组装成 (TASKS, Send)——写进内部 TASKS 通道,下一步 prepare_next_tasks 消费它生成 PUSH 任务 (第 19 讲 fan-out 的实现)。
注意 CONFIG_KEY_SEND 是 prepare_single_task 注入的(第 10 讲)—— 它指向 task.writes.extend。所以节点的所有写先进 task.writes 缓冲, 这正是第 11 讲"步内不可变"的实现根源。
读:ChannelRead 与三层读机制
节点读状态有三层,理解它们能解释很多"读到的值不对"的疑惑:
1. 超步级读(组装节点输入)
_algo.py 的 _proc_input(1348–1392 行)在 Plan 阶段,从当前通道快照读出节点订阅的通道值, 组装成节点的 input。这是节点函数 state 参数的来源——它是本超步开始时的不可变快照。
2. 任务级读(注入 CONFIG_KEY_READ)
prepare_single_task 给任务注入 CONFIG_KEY_READ = partial(local_read, ...)。 节点内部若要主动读状态(而非靠参数),走的就是它。
3. 节点内读(ChannelRead.do_read)
pregel/_read.py 的 ChannelRead(25–91 行)封装"主动读":调 config[CONFIG_KEY_READ](select, fresh)。 条件边的路由函数读 state,就是经这条路。
local_read 与 fresh=True:步内可见性的唯一特例
_algo.py 的 local_read(188–224 行)。一般情况下,节点读的是"本超步开始时的快照" (步内不可变)。但有一个特例——fresh=True:
flowchart LR
Q[节点 X 本步刚写了 a=1<br/>但在 task.writes 缓冲里] --> R{读 a?}
R -->|fresh=False 普通读| OLD[读到旧快照值]
R -->|fresh=True| NEW[模拟应用本任务的 writes<br/>读到 a=1]
local_read(fresh=True) 会把本任务尚未 commit 的 writes 临时合并到快照上再返回—— 于是节点能在同一超步内读到自己刚写的值。
最典型的用途:条件边路由函数需要看到本节点刚产生的输出来决定下一跳。 比如 agent 节点写了新消息,紧接着的路由要根据这条新消息判断走 tools 还是 END—— 靠 fresh=True 才读得到。
这是对第 1 讲"步内不可变"黄金法则的唯一受控破例:仅限"读自己本任务的写", 不会让你读到别的并行节点的未提交写。
写入路径全景(串起第 11/12/13 讲)
sequenceDiagram
participant Fn as 节点函数
participant CW as ChannelWrite
participant TW as task.writes(缓冲)
participant R as PregelRunner
participant L as PregelLoop
participant Ch as 通道本体
Fn->>CW: 返回 {"x": 1}
CW->>TW: CONFIG_KEY_SEND → extend
Note over TW: 步内只在缓冲, 通道不变
R->>L: commit → put_writes(缓冲)
L->>Ch: after_tick → apply_writes 合并
Note over Ch: 此刻才 bump 版本, 下一步可见
使用场景
- 条件边读不到刚写的值:确认路由是否在写它的节点之后(同一节点的 writer 链里), 框架用
fresh=True保证可见;若你在另一个并行节点里想读,是读不到的(符合 BSP 语义)。 - 理解 Send 的输入:PUSH 任务的 input 是
Send.arg,不是全图 state—— 因为它走_assemble_writes的 Send 分支,第 19 讲细讲。 - 自定义 writer:高级场景可往
PregelNode.writers追加ChannelWrite,实现"节点完成后额外写某通道"。
动手实验
- 写一个节点
return {"x": 1},在它的条件边路由里print(state["x"]),确认能读到 1(fresh=True 生效)。 - 在
_assemble_writes打断点,让节点返回一个Send,观察它被组装成(TASKS, Send)。 - 在
local_read打断点,对比fresh=True/False两次读到的值差异。
阅读作业
- 精读
_read.py的PregelNode(97–298 行)与ChannelRead(25–91 行)。 - 精读
_write.py的ChannelWrite与_assemble_writes(172–192 行)。 - 精读
_algo.py的local_read(188–224 行)与_proc_input(1348–1392 行)。
小结
- 节点不直接碰通道:读经
CONFIG_KEY_READ/ChannelRead,写经CONFIG_KEY_SEND/ChannelWrite→task.writes缓冲。 Send写入TASKS通道,驱动下一步 fan-out。local_read(fresh=True)是步内可见性的唯一特例:让节点读到自己本任务刚写的值(路由必需)。
下一讲:重试、超时与错误路由——给节点加健壮性。