Day 23 / 共 60 天 · 阶段 4 Pregel 执行引擎(核心深水区)

PregelNode:一个节点的"读—算—写"三段式

Day 22 说执行器跑的不是"你的函数",而是包了一层的 PregelNode。这一层负责:跑之前从 channel 好输入,跑完把输出 回 channel。今天拆 _read.pyPregelNode/ChannelRead_write.pyChannelWrite——看清你写的那个普通函数,是怎么被"读器 + 你 + 写器"串成一根 Runnable 链条的,以及 Send 到底怎么变成通道里的一条写入。

📍 你在 60 天里的位置
①入门 D1-6· ②状态 D7-12· ③控制流 D13-18· ④Pregel D19-26· ⑤通道 D27-32· ⑥持久化 D33-40· ⑦中断 D41-46· ⑧函数式 D47-52· ⑨预制件 D53-58· ⑩收官 D59-60
D19 BSP模型 D20 主循环 D21 任务准备 D22 执行器 D23 读写 D24 IO映射 D25 重试超时 D26 调试画图
L01

你的函数被包了一层:为什么需要 PregelNode

🤔 痛点:我 add_node 传的就是个普通函数 def foo(state): return {...},它怎么知道去哪读 state、又怎么把返回值写进 channel? 你写的函数只认识"一个 state 字典进、一个更新字典出"。可引擎的世界里没有"state 字典"这个现成东西——状态是拆散在一个个 channel 里的。所以必须有人在你函数之前把相关 channel 读出来拼成 state 喂给你,在你函数之后把你的返回值拆开写回对应 channel。这个"前后夹层"就是 PregelNode 干的活。

一个 PregelNode 本质是三段 Runnable 串起来的链条:

类比:翻译官 + 你 + 记录员 读器 ChannelRead(翻译官):把 channel 里的原始数据翻译成你函数能懂的 state 字典,递给你;你的函数 bound(你):按业务逻辑算出一个更新;写器 ChannelWrite(记录员):把你的更新拆开、按字段写回对应 channel。你只需要"看着 state 说话",读进来、写出去这些脏活都被两头的夹层接管了。
💡 本质:PregelNode 把"通道世界"和"函数世界"隔开 你的业务函数活在"state 字典"的世界里,引擎活在"channel + 版本号"的世界里。PregelNode 是两个世界之间的适配器:它让你写普通函数,同时让引擎拿到它需要的 channel 读写。注释原文说得直白(_read.py:98-100)——它"不会被图当 runnable 直接调用,而是作为造出 PregelExecutableTask 所需组件的容器"。
L02

PregelNode 的字段:一个节点由什么构成

先认识它的核心字段(_read.py:97-126):

# _read.py:97
class PregelNode:
    channels: str | list[str]      # 102 读哪些通道当输入(str=单值;list=拼成 dict)
    triggers: list[str]            # 107 哪些通道被写就触发我(Day 21 _triggers 用的就是它)
    mapper: Callable | None        # 111 读到输入后先转换一下(可选)
    writers: list[Runnable]        # 114 跑完后执行的写器们(把输出写回通道)
    bound: Runnable[Any, Any]      # 118 节点的主逻辑 = 你的函数
    retry_policy: Sequence[RetryPolicy] | None   # 122 重试策略(Day 25)
    cache_policy: CachePolicy | None             # 125 缓存策略(Day 51)
channels声明"我要读哪些通道"。是 str 就把那个通道的值直接当输入;是 list 就把这几个通道拼成一个 dict 当输入。StateGraph 编译时会把它设成"整个 state 的所有字段通道"。
triggers声明"哪些通道更新就唤醒我"。注意 channels 和 triggers 是两回事:triggers 决定"要不要跑我",channels 决定"跑我时读什么"。二者常常但不总是一样。
bound你的函数。默认是 DEFAULT_BOUND_read.py:94,恒等函数 lambda x: x),说明"没主逻辑、纯做读写转发"也是合法的节点。
writers写器列表。StateGraph 编译时会往里塞一个 ChannelWrite,负责把你的返回值写回 state 通道。
🎯 设计取舍①:为什么把 channels(读)和 triggers(触发)拆成两个字段? 朴素想法:一个节点"读什么"和"被什么触发"应该是同一批通道吧?多数时候是。但拆开能表达更细的语义:比如一个节点被通道 A 触发,但运行时需要同时读 A 和 B 的当前值(B 没变也要读它的最新值参与计算)。如果只有一个字段,就没法区分"因谁而醒"和"醒了读谁"。LangGraph 把两者解耦,让"触发条件"和"输入范围"各自独立配置——这也是条件边、Send、subgraph 等高级控制流能灵活组合的基础。
L03

node:把"读—算—写"组装成一根链

node 是个 cached_property,把 bound 和 writers 拼成一个可执行的 Runnable(_read.py:221-234):

# _read.py:221
@cached_property
def node(self) -> Runnable[Any, Any] | None:
    """Get a runnable that combines `bound` and `writers`."""
    writers = self.flat_writers                 # 224 先做写器合并优化(L04)
    if self.bound is DEFAULT_BOUND and not writers:
        return None                             # 226 既没逻辑也没写器 → 空节点,不跑
    elif self.bound is DEFAULT_BOUND and len(writers) == 1:
        return writers[0]                       # 228 纯写:就一个写器,直接用
    elif self.bound is DEFAULT_BOUND:
        return RunnableSeq(*writers)            # 230 纯写:多个写器串起来
    elif writers:
        return RunnableSeq(self.bound, *writers)  # 232 常态:你的函数 + 写器们,串成链
    else:
        return self.bound                       # 234 只有逻辑没写器:就跑你的函数
RunnableSeq(self.bound, *writers)最常见的一支:把你的函数 bound 和写器 writersRunnableSeq(顺序执行链)串起来。执行时先跑你的函数得到返回值,再把返回值喂给写器写回通道。这就是"算→写"两段。
return None空节点(无逻辑无写器)返回 None。Day 21 prepare_single_taskif node := proc.node 拿到 None 就不造可执行任务——省掉无意义的空跑。
cached_property这根链组装一次就缓存住,同一个节点被反复触发(循环图里 agent 跑很多次)不用每次重组。copy() 时会清掉这个缓存(_read.py:200)保证更新后重算。
数据结构:PregelNode 组装出的"读—算—写"链 ChannelRead 读 channels→输入 bound(你) 业务逻辑 ChannelWrite 写回 channels channels 状态通道 RunnableSeq(bound, *writers) —— node 属性把它们串成一根链 (Day 21 造 task 时读)
图注:ChannelRead 读入 → bound 计算 → ChannelWrite 写出,三段被 RunnableSeq 串成 node,供执行器调用。
L04

flat_writers:把相邻写器合并成一个

组装前有个小优化 flat_writers_read.py:204-219):把连续的 ChannelWrite 合并:

# _read.py:204
@cached_property
def flat_writers(self) -> list[Runnable]:
    """Get writers with optimizations applied. Dedupes consecutive ChannelWrites."""
    writers = self.writers.copy()
    while (len(writers) > 1
           and isinstance(writers[-1], ChannelWrite)
           and isinstance(writers[-2], ChannelWrite)):   # 211 末尾两个都是 ChannelWrite
        writers[-2] = ChannelWrite(                      # 215 合并成一个
            writes=writers[-2].writes + writers[-1].writes,   # 216 写入项拼起来
        )
        writers.pop()                                    # 218 去掉被合并的那个
    return writers
while 末尾两个都是 ChannelWrite从后往前看,只要末尾连着两个写器,就把它们的 writes(写入项列表)拼成一个新的 ChannelWrite。反复合并直到不再连续。
writers[-2].writes + writers[-1].writes关键:合并只是把两个写器要写的条目列表拼接,语义完全等价——本来要执行两次写入调用,现在一次搞定。
🎯 设计取舍②:合并写器省的那点开销,值得专门写段代码吗? 单看一次节点执行,省一次 ChannelWrite.invoke 调用确实微不足道。但想想循环图里的 agent 节点:一次对话可能跑几十个超步,每步都走一遍这根链。而且 LangGraph 内部为条件边、Command、Send 等会往 writers 里塞好几个 ChannelWrite——不合并的话每个节点每步都多好几层 Runnable 调用栈。热路径上的常数级优化,乘以运行次数就不小了。加上 cached_property 缓存,合并本身只算一次。这是"框架代码愿意为极致省一点点、因为它会被跑无数遍"的典型心态。
L05

ChannelRead:从 config 里的读函数取值

读这一段由 ChannelRead_read.py:25)完成。它自己不直接碰 channel,而是从 config 里取一个"读函数"(_read.py:73-91):

# _read.py:73
@staticmethod
def do_read(config, *, select, fresh=False, mapper=None) -> Any:
    try:
        read: READ_TYPE = config[CONF][CONFIG_KEY_READ]   # 82 从 config 取"读函数"
    except KeyError:
        raise RuntimeError("Not configured with a read function ...")  # 84 不在 Pregel 上下文里
    if mapper:
        return mapper(read(select, fresh))                # 89 读完再转换一下
    else:
        return read(select, fresh)                        # 91 直接读
config[CONF][CONFIG_KEY_READ]读通道不是硬编码的,而是运行时从 config 注入的一个函数。谁注入的?PregelLoop 在造任务时把"当前这套 channels 的读函数"塞进 config。这样同一个节点在不同 run / 不同 checkpoint 下读到的是各自对应的状态。
read(select, fresh)select 是要读的通道名(单个或多个),fresh 表示是否要读"最新提交后"的值。返回值就是喂给你函数的输入。
RuntimeError(...)如果不在 Pregel 运行上下文里(config 没有读函数)就报错——这也是为什么你不能脱离图、单独 node.invoke() 期望它自己读 state。
💡 本质:读/写都靠 config 注入的函数,节点本身"无状态" ChannelRead(读)和马上要看的 ChannelWrite(写)都不持有 channels,而是从 config 里现取 CONFIG_KEY_READ / CONFIG_KEY_SEND 这两个函数。好处是节点定义本身完全无状态、可复用——真正的"读到哪套状态、写到哪套状态"由运行时 config 决定。这让同一张编译好的图能被无数个并发 thread 安全共享(各自 config 不同)。
L06

ChannelWrite:把返回值写回通道

对称地,ChannelWrite_write.py:46)负责写。它的 do_write 先校验再调 config 里的发送函数(_write.py:105-126):

# _write.py:105
@staticmethod
def do_write(config, writes, allow_passthrough=True) -> None:
    for w in writes:                                # 112 逐条校验
        if isinstance(w, ChannelWriteEntry):
            if w.channel == TASKS:                  # 114 禁止直接写保留通道 TASKS
                raise InvalidUpdateError("Cannot write to the reserved channel TASKS")
            if w.value is PASSTHROUGH and not allow_passthrough:
                raise InvalidUpdateError("PASSTHROUGH value must be replaced")  # 119
        ...
    write: TYPE_SEND = config[CONF][CONFIG_KEY_SEND]   # 125 从 config 取"发送函数"
    write(_assemble_writes(writes))                    # 126 组装成 (通道,值) 元组列表后发出
w.channel == TASKS → 报错边界守卫:TASKS 是给 Send 专用的保留通道,你不能直接往它写普通值。要派发动态任务只能用 Send 对象(L07 会看到 Send 才被允许进 TASKS)。
PASSTHROUGH一个哨兵值,表示"这一项要写的值就是节点的输入本身"。比如 add_edge 产生的"透传写"——把上游值原样传给下游通道。校验它必须在有输入可透传时才用。
config[CONF][CONFIG_KEY_SEND]和读对称:写也是取 config 里注入的发送函数。这个函数最终把 (通道名, 值) 攒进当前 task 的 writes——就是 Day 22 commit 时收走、Day 21 apply_writes 时合并的那批写入。
⚠️ 边界/坑:节点返回值里写了 state 里没有的字段,会怎样? 不会立刻报错。ChannelWrite 照样把它组装成 (未知字段, 值) 发出去,一路传到 Day 21 的 apply_writes——那里发现"这个通道不在 channels 里",只 logger.warning(...unknown channel...) 警告并丢弃,图继续跑。所以经典 bug 是:你 return {"mesage": ...}(拼错 messages),运行没报错但状态里那字段永远是空的。排查这类"写了没生效"的问题,第一件事是去日志里搜 unknown channel。
L07

Send 怎么变成 TASKS 通道里的一条写入

写入的最后一步 _assemble_writes_write.py:172)把各种写入项标准化成 (通道, 值) 元组。Send 在这里被特殊处理:

# _write.py:172
def _assemble_writes(writes) -> list[tuple[str, Any]]:
    tuples: list[tuple[str, Any]] = []
    for w in writes:
        if isinstance(w, Send):
            tuples.append((TASKS, w))               # 179 Send → 写进 TASKS 通道!
        elif isinstance(w, ChannelWriteTupleEntry):
            if ww := w.mapper(w.value):             # 181 一个映射产出多条写入
                tuples.extend(ww)
        elif isinstance(w, ChannelWriteEntry):
            value = w.mapper(w.value) if w.mapper is not None else w.value
            if value is SKIP_WRITE:                 # 185 mapper 返回 SKIP_WRITE → 跳过不写
                continue
            if w.skip_none and value is None:       # 187 skip_none 且值为 None → 跳过
                continue
            tuples.append((w.channel, value))       # 189 普通写:(通道名, 值)
        else:
            raise ValueError(f"Invalid write entry: {w}")
    return tuples
isinstance(w, Send) → (TASKS, w)今天的高光时刻:一个 Send 对象被组装成 (TASKS, Send对象) 这条写入。它流进 TASKS 通道后,Day 21 prepare_next_tasks 开头那段就会把它捞出来变成 PUSH 任务。Send 的"派发"本质上就是一次向 TASKS 通道的写入——控制流被统一成了数据流!
SKIP_WRITE / skip_none两个"条件不写"的开关。SKIP_WRITE 让 mapper 能动态决定"这次不写";skip_none 让"值是 None 就别覆盖"。条件边、可选写入靠它们实现"有时写有时不写"。
ChannelWriteTupleEntry一个写入项能通过 mapper 产出多条写入(一对多)。add_messages、Command.update 等把一个返回值拆成多个通道写入时走这条。
💡 本质:LangGraph 把"控制流"也编码成"数据流" 这是整个引擎最精巧的一处统一:你以为 Send 是"跳转/派发"(控制流),可它落地成了往一个特殊通道写数据(数据流)。于是执行引擎只需要会做一件事——"合并写入、按版本触发",就同时实现了普通边、条件边、循环、map-reduce 扇出。没有单独的"调度器"处理 Send,Send 就是数据,和 messages、counter 一样流过同一套 channel 机制。这就是 Pregel 模型"一切皆消息"的威力。
控制流:一个 Send 从"写出"到"被跑"的全程 节点 return Send("b", arg) _assemble_writes → (TASKS, Send) apply_writes 进 TASKS 通道 下一步 prepare_next_tasks PUSH 任务跑 注意:跨越了超步边界——Send 本步写入,下一超步才被捞出执行(BSP) 控制流(派发)被彻底编码成数据流(写通道)
图注:Send 被组装成 (TASKS, Send) 写入通道,下一超步 prepare_next_tasks 捞出变 PUSH 任务执行。
L08

今日小结 + 动手 + 明日预告

🧠 今天你应该能回答

  • 为什么你的函数要被 PregelNode 包一层?(适配"函数世界"和"通道世界",前后夹读器/写器)
  • channels 和 triggers 有何区别?(channels=跑我时读什么;triggers=谁更新唤醒我)
  • node 属性做了什么?(RunnableSeq 把 bound + writers 串成"算→写"链)
  • flat_writers 优化了什么?(合并相邻 ChannelWrite,减少热路径调用层数)
  • ChannelRead/ChannelWrite 自己持有 channels 吗?(不,从 config 取读/发送函数,节点无状态)
  • 写了 state 没有的字段会怎样?(apply_writes 只 warning 丢弃,不报错——排查搜 unknown channel)
  • Send 怎么被执行的?(组装成 (TASKS, Send) 写入 → 下步 prepare 捞成 PUSH 任务,控制流即数据流)

✋ 10 分钟动手

# 1. 读 PregelNode 字段 + node 组装
sed -n '97,234p' libs/langgraph/langgraph/pregel/_read.py

# 2. 读 ChannelWrite.do_write(校验 + 发送)
sed -n '105,126p' libs/langgraph/langgraph/pregel/_write.py

# 3. 读 _assemble_writes(Send→TASKS 那行)
sed -n '172,192p' libs/langgraph/langgraph/pregel/_write.py

# 4. 亲手触发 unknown channel 警告(故意拼错字段)
python -c "
import logging; logging.basicConfig(level=logging.WARNING)
from langgraph.graph import StateGraph, START, END
from typing import TypedDict
class S(TypedDict): count: int
g = StateGraph(S)
g.add_node('a', lambda s: {'cont': 1})   # 拼错 count→cont
g.add_edge(START,'a'); g.add_edge('a',END)
print(g.compile().invoke({'count':0}))    # 观察 warning + count 没变
"
明天预告 · Day 24:今天看的是"节点内部"的读写。明天上升一层看"整张图的边界"——你 invoke({...}) 传进去的输入怎么变成第一批 channel 写入?图跑完又怎么从 channel 读出最终结果返回给你?进 _io.pymap_input / map_output_values / map_output_updates,看清输入输出的映射规则,以及 stream 的 values/updates 两种模式差在哪。
← Day 22 执行器 Day 24 · 输入输出映射 _io.py →