BSP / Pregel 模型:图是怎么"跑起来"的
前 18 天我们学会了怎么描述一张图(节点、边、状态、控制流)。可这张图只是"图纸"——它自己不会动。从今天起进入全课最硬核的一段:执行引擎。第一站先建立心智模型——LangGraph 的引擎叫 Pregel,它把图的运行组织成一波一波的"超步"(BSP,Bulk Synchronous Parallel)。读 pregel/main.py:450 的 Pregel 类。
为什么图需要一套"执行引擎"
messages,谁的改动先生效?一个节点还没跑完,另一个节点能不能看到它的中间结果?——这些问题必须有一个统一的裁判来回答。这个裁判就是执行引擎 Pregel。把前 18 天学的东西回顾一下:StateGraph 描述"有哪些节点、边怎么连、状态长什么样",但它只是一份声明。真正让图"动起来"——决定每一轮跑哪些节点、怎么并行、写入怎么合并、什么时候停——是另一套完全独立的代码,位于 libs/langgraph/langgraph/pregel/。
graph/ 包管"画图",pregel/ 包管"跑图"。你写的 StateGraph 最后会被 .compile() 编译成一个 Pregel 实例——就像源代码被编译成可执行文件。之后所有的 invoke/stream,其实都是在驱动这个 Pregel 引擎。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)。三阶段拆解:Plan / Execution / Update
把 docstring 里那三个阶段翻译成"这一步引擎具体在忙什么",并预告后面几天分别讲哪块:
| 阶段 | 做什么 | 源码位置 | 哪天细讲 |
|---|---|---|---|
| Plan 计划 | 看上一步哪些 channel 被更新了,据此决定这一步激活哪些节点,并为每个节点读好输入打包成"任务" | prepare_next_tasks _algo.py:392 | Day 21 |
| Execution 执行 | 把任务丢进线程池/事件循环并行跑;每个任务把写入攒在自己身上,不直接改 channel | _runner.py / _executor.py | Day 22 |
| Update 更新 | 所有任务跑完,把它们攒的写入统一合并进 channel(用各自的 reducer),算出"下一步谁被触发" | apply_writes _algo.py:232 | Day 21 |
而把这三阶段串成循环、反复跑的那个"主发条",是 _loop.py 的 PregelLoop.tick()(Day 20 主讲)。今天只需记住这张表——它是整个阶段 4 的地图。
👶 小白:那如果一个节点执行时间特别长,会不会拖慢整个超步?
👨🏫 老师:会。BSP 的代价就是木桶效应——一个超步必须等它里面最慢的节点跑完才能进入 Update。这是"同步"换来的确定性所付的税。所以 LangGraph 才配了 step_timeout(Day 20)和节点级 TimeoutPolicy(Day 25):给慢节点上闹钟,别让一个卡死的节点拖垮整波。
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 写结果。你在 StateGraph 里 add_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 求和)。
callB();A 只是往某个 channel 写值,而 B 恰好订阅了这个 channel,于是下一个超步 B 就被触发了。节点彼此解耦、只认 channel——这让"加一条边"变成"让某节点多订阅一个 channel",也让并行、扇出(Send)、条件路由都能统一用"谁写了哪个 channel、谁订阅了哪个 channel"来表达。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 引擎逻辑一行不用改。tick() 里,每种组合都得改引擎。做成可插拔字段(依赖注入)后,引擎只依赖抽象接口(BaseCheckpointSaver 等),具体实现运行时注入。代价是字段多、构造函数长(main.py:758 的 __init__ 参数一大串),但换来了引擎内核的稳定——阶段 6 换四种存档器,这段循环代码纹丝不动。脱掉语法糖:直接用 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.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 / START | input_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。今日小结 + 动手 + 明日预告
🧠 今天你应该能回答
- 为什么图需要独立的执行引擎?(图允许环/并行/多写,"谁先跑"需要统一裁判;描述与执行分离)
- 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))
"
_loop.py 的 PregelLoop.tick()——看主循环怎么判断"该不该再来一步"、step 和 stop 怎么算、递归上限从哪来,以及 invoke 时那个 while loop.tick(): ... loop.after_tick() 驱动循环的真身。