发布日期

第 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_SENDprepare_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.pyChannelRead(25–91 行)封装"主动读":调 config[CONFIG_KEY_READ](select, fresh)。 条件边的路由函数读 state,就是经这条路。

local_read 与 fresh=True:步内可见性的唯一特例

_algo.pylocal_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,实现"节点完成后额外写某通道"。

动手实验

  1. 写一个节点 return {"x": 1},在它的条件边路由里 print(state["x"]),确认能读到 1(fresh=True 生效)。
  2. _assemble_writes 打断点,让节点返回一个 Send,观察它被组装成 (TASKS, Send)
  3. local_read 打断点,对比 fresh=True/False 两次读到的值差异。

阅读作业

  • 精读 _read.pyPregelNode(97–298 行)与 ChannelRead(25–91 行)。
  • 精读 _write.pyChannelWrite_assemble_writes(172–192 行)。
  • 精读 _algo.pylocal_read(188–224 行)与 _proc_input(1348–1392 行)。

小结

  • 节点不直接碰通道:读经 CONFIG_KEY_READ/ChannelRead,写经 CONFIG_KEY_SEND/ChannelWritetask.writes 缓冲。
  • Send 写入 TASKS 通道,驱动下一步 fan-out。
  • local_read(fresh=True) 是步内可见性的唯一特例:让节点读到自己本任务刚写的值(路由必需)。

下一讲:重试、超时与错误路由——给节点加健壮性。