Day 23 / 共 60 天 · 阶段 4 Pregel 执行引擎(核心深水区)
PregelNode:一个节点的"读—算—写"三段式
Day 22 说执行器跑的不是"你的函数",而是包了一层的 PregelNode。这一层负责:跑之前从 channel 读好输入,跑完把输出 写回 channel。今天拆 _read.py 的 PregelNode/ChannelRead 和 _write.py 的 ChannelWrite——看清你写的那个普通函数,是怎么被"读器 + 你 + 写器"串成一根 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 和写器 writers 用 RunnableSeq(顺序执行链)串起来。执行时先跑你的函数得到返回值,再把返回值喂给写器写回通道。这就是"算→写"两段。return None空节点(无逻辑无写器)返回 None。Day 21 prepare_single_task 里 if node := proc.node 拿到 None 就不造可执行任务——省掉无意义的空跑。cached_property这根链组装一次就缓存住,同一个节点被反复触发(循环图里 agent 跑很多次)不用每次重组。copy() 时会清掉这个缓存(_read.py:200)保证更新后重算。图注: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 被组装成 (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.py 的 map_input / map_output_values / map_output_updates,看清输入输出的映射规则,以及 stream 的 values/updates 两种模式差在哪。