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

BSP / Pregel 模型:图是怎么"跑起来"的

前 18 天我们学会了怎么描述一张图(节点、边、状态、控制流)。可这张图只是"图纸"——它自己不会动。从今天起进入全课最硬核的一段:执行引擎。第一站先建立心智模型——LangGraph 的引擎叫 Pregel,它把图的运行组织成一波一波的"超步"(BSP,Bulk Synchronous Parallel)。读 pregel/main.py:450Pregel 类。

📍 你在 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

为什么图需要一套"执行引擎"

🤔 痛点:我画了一张有环、有分支、有并行的图,谁来决定"先跑谁、后跑谁"? 普通函数调用是"你调我、我调你",调用栈自然决定顺序。可 LangGraph 的图允许环(loop)、允许一个节点同时触发多个下游、允许多个节点并行往同一个状态字段写。这时候"执行顺序"就不再是显而易见的了:两个并行节点都改 messages,谁的改动先生效?一个节点还没跑完,另一个节点能不能看到它的中间结果?——这些问题必须有一个统一的裁判来回答。这个裁判就是执行引擎 Pregel。

把前 18 天学的东西回顾一下:StateGraph 描述"有哪些节点、边怎么连、状态长什么样",但它只是一份声明。真正让图"动起来"——决定每一轮跑哪些节点、怎么并行、写入怎么合并、什么时候停——是另一套完全独立的代码,位于 libs/langgraph/langgraph/pregel/

💡 本质:把"图长什么样"和"图怎么执行"彻底分开 这是 LangGraph 最重要的架构决定。graph/ 包管"画图",pregel/ 包管"跑图"。你写的 StateGraph 最后会被 .compile() 编译成一个 Pregel 实例——就像源代码被编译成可执行文件。之后所有的 invoke/stream,其实都是在驱动这个 Pregel 引擎。
用一个类比兜住整个阶段 4 Pregel 就像一个剧组导演:演员(节点)各演各的戏,但导演喊"这一场,A、B、C 三位一起上!"——三人同时表演(并行),演完导演喊"卡!",把这一场所有人的表演统一剪进片子(合并写入),然后才安排下一场。演员在同一场里看不到彼此的最新表演,只能看到上一场剪好的成片。这个"一场一场、场内并行、场末统一收工"的节奏,就是 BSP 超步。
L02

BSP 超步模型:一波一波地跑

BSP 全称 Bulk Synchronous Parallel(批量同步并行),源自 Google 处理超大图(比如全网页链接图)的 Pregel 论文。LangGraph 直接借用了这个名字和模型。核心一句话:计算被切成一个个"超步(super-step)",每个超步内部所有活跃节点并行跑,跑完统一同步一次,再进下一步。

Pregel 类的 docstring 把这个模型写得清清楚楚(main.py:461-477):

# main.py:461 —— Pregel 类文档,逐字翻译在下面
# Pregel organizes the execution of the application into multiple steps,
# following the **Pregel Algorithm**/**Bulk Synchronous Parallel** model.
#
# Each step consists of three phases:
#  - Plan:  决定这一步要跑哪些 actor(节点)
#  - Execution: 并行执行所有选中的 actor,直到全部完成 / 其一失败 / 超时
#               —— 期间,channel 更新对 actor 不可见,要等到下一步
#  - Update: 用本步 actor 写出的值,更新 channel
#
# Repeat until no actors are selected for execution,
# or a maximum number of steps is reached.
multiple steps执行不是"一条路走到黑",而是分成很多个 step(超步)。每个 step 是一次"计划→执行→更新"的完整循环。
in parallel一个超步里被选中的节点同时并行跑,不是一个跑完再跑下一个。这就是为什么两个并行分支能真正并发。
invisible until next step最关键的一句:本步内某节点写的值,本步内其它节点看不到,要等到下一步才可见。这消除了"谁先谁后"的竞态。
Repeat until no actors停止条件:某一步"计划"阶段发现没有节点需要跑了,或达到步数上限(recursion_limit)。
💡 本质:用"步与步之间才同步"换来"步内无锁并行" 传统并发要用锁保护共享变量,容易死锁、难调试。BSP 换了个思路:步内大家读的都是上一步的快照、彼此隔离,写出去的改动先攒着不生效;等这一步所有人都跑完了,引擎在一个明确的时刻把所有写入统一合并。于是"并行"和"确定性"同时拿到了——这正是 LangGraph 敢让你随便并行、还能配合 checkpoint 精确恢复的底气。
控制流:三个超步,步内并行、步末同步 超步 0 A 同步·合并写入 超步 1(并行) B C B、C 看不到彼此的写入 同步·统一 apply_writes 超步 2 D D 读到的,是超步 1 结束时"剪好的成片",绝不会读到 B/C 的半成品 Plan → Execution(并行) → Update(同步合并) ↺ 直到没有节点可跑
图注:一个超步 = Plan + Execution + Update。步内并行、彼此隔离;步末统一同步。
L03

三阶段拆解:Plan / Execution / Update

把 docstring 里那三个阶段翻译成"这一步引擎具体在忙什么",并预告后面几天分别讲哪块:

阶段做什么源码位置哪天细讲
Plan 计划看上一步哪些 channel 被更新了,据此决定这一步激活哪些节点,并为每个节点读好输入打包成"任务"prepare_next_tasks _algo.py:392Day 21
Execution 执行把任务丢进线程池/事件循环并行跑;每个任务把写入攒在自己身上,不直接改 channel_runner.py / _executor.pyDay 22
Update 更新所有任务跑完,把它们攒的写入统一合并进 channel(用各自的 reducer),算出"下一步谁被触发"apply_writes _algo.py:232Day 21

而把这三阶段串成循环、反复跑的那个"主发条",是 _loop.pyPregelLoop.tick()(Day 20 主讲)。今天只需记住这张表——它是整个阶段 4 的地图。

🎯 设计取舍①:为什么"读输入"放在 Plan、"合并写入"放在 Update,中间隔着 Execution? 因为这正是 BSP "隔离"的实现方式。Plan 阶段一次性把每个节点的输入从 channel 里读好、冻结——所以节点执行期间即使别的节点在写,也影响不到它读到的输入。写入全攒到 Update 一起合并——所以并行节点的写入有一个统一的、确定的合并时刻。如果读和写都在执行时零散进行,就退回到"需要加锁的共享内存"模型了。把 IO 挤到执行的前后两端,是 BSP 能做到"无锁并行 + 确定性"的关键工程手法。

👶 小白:那如果一个节点执行时间特别长,会不会拖慢整个超步?

👨‍🏫 老师:会。BSP 的代价就是木桶效应——一个超步必须等它里面最慢的节点跑完才能进入 Update。这是"同步"换来的确定性所付的税。所以 LangGraph 才配了 step_timeout(Day 20)和节点级 TimeoutPolicy(Day 25):给慢节点上闹钟,别让一个卡死的节点拖垮整波。

L04

Actors + Channels:Pregel 的两块积木

Pregel 世界只有两种东西,docstring 开头就点明了(main.py:458-460):

# main.py:458
# Pregel combines **actors** and **channels** into a single application.
# **Actors** read data from channels and write data to channels.
Actor(演员)就是一个 PregelNode(Day 23 细讲)。它订阅若干 channel(读输入),干活,再往若干 channel 写结果。你在 StateGraphadd_node 加的每个节点,编译后都变成一个 PregelNode。
Channel(通道)节点之间传数据的"信箱"(阶段 5 D27-32 专门讲)。每个 channel 有自己的"合并规则":LastValue 后写覆盖前写、BinaryOperatorAggregate 累加、Topic 追加成列表……你在状态里写的 Annotated[list, add_messages],那个 reducer 就是 channel 的合并规则。

docstring 还列了内置 channel 类型(main.py:496-512):LastValue(默认,存最后一个值)、Topic(发布订阅/累积列表)、Context(管理外部资源生命周期)、BinaryOperatorAggregate(用二元运算累积,如 operator.add 求和)。

💡 本质:节点之间不直接互相调用,只通过 channel 间接通信 这是 Pregel 和"普通函数互调"最大的不同。节点 A 不会 callB();A 只是往某个 channel 写值,而 B 恰好订阅了这个 channel,于是下一个超步 B 就被触发了。节点彼此解耦、只认 channel——这让"加一条边"变成"让某节点多订阅一个 channel",也让并行、扇出(Send)、条件路由都能统一用"谁写了哪个 channel、谁订阅了哪个 channel"来表达。
数据结构:Actor 读/写 Channel,彼此不直接调用 Actor A channel: msgs channel: count Actor B 订阅→读 A 写 msgs → B(订阅 msgs)在下个超步被触发;每个 channel 自带合并规则
图注:Actor 只和 Channel 打交道;"边"的本质是"订阅关系"。
L05

Pregel 类字段走读:一个引擎实例装了什么

现在打开 Pregel 类本体,看它到底持有哪些字段(main.py:705-756,裁剪核心):

class Pregel(...):                              # main.py:450
    nodes: dict[str, PregelNode]               # 705 所有节点(actor),name -> PregelNode
    channels: dict[str, BaseChannel | ...]     # 707 所有通道,name -> Channel
    stream_mode: StreamMode = "values"         # 709 默认流式模式
    input_channels: str | Sequence[str]        # 725 输入写进哪些 channel(如 __start__)
    output_channels: str | Sequence[str]       # 716 从哪些 channel 读最终输出
    stream_channels: str | Sequence[str] | None = None  # 718 流式时读哪些 channel
    interrupt_before_nodes: All | Sequence[str]         # 723 静态断点(Day 43)
    interrupt_after_nodes: All | Sequence[str]          # 721
    step_timeout: float | None = None          # 727 单个超步的墙钟超时(Day 20)
    checkpointer: Checkpointer = None          # 733 存档器(阶段 6)
    store: BaseStore | None = None             # 736 长期记忆(Day 40)
    retry_policy: Sequence[RetryPolicy] = ()   # 742 全局重试策略(Day 25)
    trigger_to_nodes: Mapping[str, Sequence[str]]       # 755 反查表:channel -> 订阅它的节点
nodes / channels引擎的两块核心积木,正好对应 L04 的 Actor 和 Channel。整个执行就是"在 channels 上、按订阅关系反复激活 nodes"。
input_channels / output_channels引擎的"进出口"。invoke 时你传的输入写进 input_channels;跑完从 output_channels 读结果。Day 24 的 map_input/map_output 就是干这个的。
trigger_to_nodes性能关键的反查表:给定"哪个 channel 被更新了",O(1) 查出"哪些节点订阅了它、该被触发"。没有它,Plan 阶段每步都要遍历所有节点。Day 21 会看到它怎么用。
checkpointer / store可插拔的持久化与记忆。留意:它们是 Pregel 的字段——执行引擎和持久化是解耦的,你换个 checkpointer 引擎逻辑一行不用改。
🎯 设计取舍②:为什么把 checkpointer / store / cache / retry_policy 都做成 Pregel 的字段,而不是写死在循环里? 因为执行引擎要面对千差万别的部署:本地跑用内存存档、生产用 Postgres、要不要长期记忆、要不要缓存、要不要重试……如果写死在 tick() 里,每种组合都得改引擎。做成可插拔字段(依赖注入)后,引擎只依赖抽象接口(BaseCheckpointSaver 等),具体实现运行时注入。代价是字段多、构造函数长(main.py:758__init__ 参数一大串),但换来了引擎内核的稳定——阶段 6 换四种存档器,这段循环代码纹丝不动。
L06

脱掉语法糖:直接用 Pregel 写一个最小应用

平时你都用 StateGraph,看不到 Pregel 长什么样。docstring 里给了一个不用 StateGraph、直接手搓 Pregel的例子(main.py:527-553),最能暴露引擎的本来面目:

# main.py:527
from langgraph.channels import EphemeralValue
from langgraph.pregel import Pregel, NodeBuilder

node1 = (
    NodeBuilder().subscribe_only("a")   # 订阅 channel "a"(读输入)
    .do(lambda x: x + x)                # 干活:把输入翻倍
    .write_to("b")                      # 把结果写到 channel "b"
)

app = Pregel(
    nodes={"node1": node1},
    channels={
        "a": EphemeralValue(str),       # 两个通道
        "b": EphemeralValue(str),
    },
    input_channels=["a"],               # 输入进 "a"
    output_channels=["b"],              # 从 "b" 取输出
)

app.invoke({"a": "foo"})                # -> {'b': 'foofoo'}
subscribe_only("a")声明"我这个 actor 订阅 channel a"。等价于 StateGraph 里"这个节点接在 a 后面"。trigger 就是这么来的。
.do(...)节点真正的业务逻辑(bound)。这里是把字符串翻倍。
.write_to("b")声明"我写 channel b"。它编译成一个 ChannelWrite(Day 23)。
invoke({"a":"foo"})引擎流程:把 foo 写进 a → 超步 0 触发 node1(订阅 a)→ node1 输出 foofoo 写进 b → 下步没有节点订阅 b → 停止 → 从 output_channels 读 b。
💡 本质:StateGraph 只是 Pregel 的"高层语法糖" docstring 明说了(main.py:520-523):"如果你不确定要不要直接用 Pregel,那答案多半是不用——用 Graph API 或 Functional API,它们最终都会 compile 成底层的 Pregel。"换句话说:你学的 add_node/add_edge,本质就是在往一个 Pregel 里塞 nodes 和配置 channel 订阅关系。看懂这个最小例子,就看懂了整个引擎的输入输出边界。
L07

StateGraph.compile() 到底产出了什么

把 L06 反过来想:你写的 StateGraph.compile() 之后其实就返回了一个 Pregel 的子类实例(CompiledStateGraph)。对照一下你熟悉的写法与引擎字段的关系:

你在 StateGraph 里写的编译后落到 Pregel 的哪里
add_node("agent", fn)nodes["agent"] = PregelNode(...),bound 包着你的 fn
状态里 Annotated[list, add_messages]channels["messages"] = 带该 reducer 的 channel
add_edge("a", "b")让节点 b 订阅"a 完成"这个 channel(trigger_to_nodes 里登记)
add_conditional_edges(...)编译成一个分支 writer,运行时决定写哪个 branch:to:* channel
set_entry_point / STARTinput_channels + 从 START 到入口节点的订阅关系
checkpointer=Saver()Pregel.checkpointer 字段

验证一下也很简单——编译产物就是个 Pregel:

from langgraph.graph import StateGraph
g = StateGraph(dict)
g.add_node("double", lambda s: {"x": s["x"] * 2})
g.set_entry_point("double")
app = g.compile()

type(app).__mro__      # 里面能看到 Pregel —— 编译产物是 Pregel 子类
app.nodes.keys()       # dict_keys(['__start__', 'double'])  ← 多了个内部起点节点
app.channels.keys()    # 你的状态字段都变成了 channel
⚠️ 边界/易错:为什么 app.nodes 里会冒出 __start__ 这种你没加的节点? 这是编译期自动注入的内部节点/通道(constants.py 里的 START="__start__")。引擎需要一个"起点"来把 invoke 的输入喂进图里——它把你的输入写进 __start__ 通道,触发真正的入口节点。所以你直接打印 nodes/channels 会看到一批带下划线的"系统件"。初学者常误以为是 bug,其实是引擎的脚手架。Day 18 讲过 START/END,这里能看到它们在引擎层的真身:就是普通的 channel/node。
L08

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

🧠 今天你应该能回答

  • 为什么图需要独立的执行引擎?(图允许环/并行/多写,"谁先跑"需要统一裁判;描述与执行分离)
  • BSP 超步模型一句话是什么?(切成一个个超步,步内并行、彼此隔离,步末统一同步)
  • 一个超步的三个阶段?(Plan 计划并读输入 / Execution 并行执行攒写入 / Update 合并写入算触发)
  • Actor 和 Channel 是什么、怎么通信?(Actor=PregelNode,只通过订阅 channel 间接通信,不直接互调)
  • Pregel 类持有哪些关键字段?(nodes/channels/input_output_channels/trigger_to_nodes/checkpointer…)
  • StateGraph 和 Pregel 什么关系?(StateGraph 是语法糖,compile 后就是一个 Pregel 实例)

✋ 10 分钟动手

# 1. 读 Pregel 类的 BSP 文档(三阶段就在这里)
sed -n '450,477p' libs/langgraph/langgraph/pregel/main.py

# 2. 看 Pregel 到底有哪些字段
sed -n '705,756p' libs/langgraph/langgraph/pregel/main.py

# 3. Python 里验证"编译产物就是 Pregel"
python -c "
from langgraph.graph import StateGraph
g = StateGraph(dict); g.add_node('d', lambda s: {'x': s['x']*2})
g.set_entry_point('d'); app = g.compile()
print([c.__name__ for c in type(app).__mro__])
print('nodes:', list(app.nodes)); print('channels:', list(app.channels))
"
明天预告 · Day 20:知道了"一波一波跑",那到底是谁在摇动这个发条?明天钻进 _loop.pyPregelLoop.tick()——看主循环怎么判断"该不该再来一步"、stepstop 怎么算、递归上限从哪来,以及 invoke 时那个 while loop.tick(): ... loop.after_tick() 驱动循环的真身。
← Day 18 START/END 入口 Day 20 · 超步主循环 →