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

任务准备:prepare_next_tasks 与 apply_writes

昨天 tick() 里留了两个"黑盒":prepare_next_tasks(...)(这一步该跑哪些节点?)和 apply_writes(...)(这一步的写入怎么合并、谁又被触发?)。这俩就是整个 Pregel 引擎的心脏——一个负责 Plan,一个负责 Update,中间夹着执行器。今天钻进 _algo.py 把它俩逐行拆开,看清"版本号驱动的触发"这套机制到底怎么运转。

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

备菜与结账:两个函数的分工

🤔 痛点:昨天说 tick 会"计划任务"、after_tick 会"合并写入",可"计划"到底怎么算出来的?凭什么某个节点这步跑、下步不跑? Day 19 说"哪些 channel 更新了,就触发订阅它的节点"。这句话听起来简单,可代码里怎么判断"一个 channel 更新了没有"?总不能每步都把整个 state 存两份来 diff 吧——那太慢也太占内存。LangGraph 的答案是版本号:每个 channel 有一个单调递增的版本号,每个节点记住"我上次看到的版本号"。新版本 > 已见版本 = 有新东西 = 触发。这套机制就藏在今天的两个函数里。

阶段 4 的核心其实就三个动词:Plan(备菜)→ Execute(炒菜)→ Update(结账)。今天讲头尾两个,都在 _algo.py 里:

  • prepare_next_tasks(...)_algo.py:392)= Plan:读当前 checkpoint 和 channels,算出"这一步该激活哪些节点",把每个节点的输入读好、打包成一个个 task。
  • apply_writes(...)_algo.py:232)= Update:把这一步所有 task 攒的写入按 reducer 合并进 channel、给被写的 channel 升版本号,并返回"这一步到底更新了哪些 channel"。
💡 本质:一个把"channel 版本"翻译成"任务列表",另一个把"任务写入"翻译回"channel 版本" 这两个函数是一对互逆的翻译器,靠"版本号"这个中间货币衔接。apply_writes 升高的版本号,正是下一次 prepare_next_tasks 判断触发的依据。Day 20 那个 updated_channels 闭环,物理载体就是这个版本号。
类比:图书馆的"上新登记本" 每本书(channel)有个"第几次上架"的编号(版本号)。每位读者(节点)在借书证上记着"我上次借到第几版"。图书馆员(prepare_next_tasks)只需对比两个数字,就知道"这位读者关注的书上新了没",上新了就通知他来看(触发任务)。读者看完还书、写了读后感(写入),管理员(apply_writes)把读后感归档、把书的上架编号 +1——于是又能通知下一批关注这本书的读者。整座图书馆不需要把每本书复印一份来对比。
L02

prepare_next_tasks 的骨架:PUSH + PULL 两类任务

先看整体结构。函数很长,但骨架就两段:先收 PUSH 任务(Send 动态派发),再收 PULL 任务(被 channel 触发的普通节点)。_algo.py:437-513

# _algo.py:437
def prepare_next_tasks(checkpoint, pending_writes, processes, channels, ...) -> dict[str, ...]:
    tasks: list[PregelTask | PregelExecutableTask] = []
    # ① 消费 PUSH 队列:上一步 Send 出去的动态任务(Day 16 map-reduce)
    tasks_channel = cast(Topic[Send] | None, channels.get(TASKS))     # 442
    if tasks_channel and tasks_channel.is_available():
        for idx, _ in enumerate(tasks_channel.get()):                # 444
            if task := prepare_single_task((PUSH, idx), None, ...):   # 445
                tasks.append(task)
    # ② 计算 PULL 候选节点(见 L03 的优化)
    if updated_channels and trigger_to_nodes:                        # 475
        ...candidate_nodes = 只查被更新 channel 触发的那几个节点...
    elif not checkpoint["channel_versions"]:                         # 483
        candidate_nodes = ()                                         # 全新空图,没人可跑
    else:
        candidate_nodes = processes.keys()                           # 兜底:全图扫一遍
    # ③ 逐个候选节点尝试造 PULL 任务
    for name in candidate_nodes:                                     # 490
        if task := prepare_single_task((PULL, name), None, ...):     # 491
            tasks.append(task)
    return {t.id: t for t in tasks}                                  # 513 按 task_id 建字典
tasks_channel = channels.get(TASKS)PUSH 任务藏在一个叫 TASKS 的特殊 Topic 通道里。上一步谁调了 Send("node", arg),就往这个通道塞了一条。这里把它们取出来一个个变成任务。
(PUSH, idx) / (PULL, name)每个任务用一个"路径元组"标识类型。PUSH = 动态扇出(Send,按下标 idx),PULL = 被边/channel 触发(按节点名 name)。prepare_single_task 根据元组第一项分派处理(L04/L05)。
return {t.id: t for t in tasks}最后按 task_id 建成字典返回。task_id 是根据 checkpoint_id + 节点名 + 步号 + 触发通道算出的确定性哈希——同样的输入必得同样的 id,这是"断点续跑不重复执行"的地基(Day 46)。
💡 本质:PULL 是"被数据拉起来",PUSH 是"被别人推过来" LangGraph 的两种触发方式在这里第一次同框:PULL 对应静态图结构(A 写完 channel,订阅它的 B 被"拉"起来);PUSH 对应运行时动态(节点主动 Send 出一批子任务,被"推"进队列)。map-reduce 就靠 PUSH(Day 16)。两类任务在同一个超步里可以并存、并行跑。
L03

增量触发:候选集从"全图"缩到"极小"

L02 里那段"算 candidate_nodes"是 Pregel 能扛大图的关键优化。放大看(_algo.py:468-486):

# _algo.py:468 —— 注释原文:an optimization that allows which nodes will be active
if updated_channels and trigger_to_nodes:          # 475 上一步更新了哪些 channel + channel→节点反查表
    triggered_nodes: set[str] = set()
    for channel in updated_channels:               # 478 只遍历"变了的" channel
        if node_ids := trigger_to_nodes.get(channel):
            triggered_nodes.update(node_ids)        # 480 把订阅它的节点加进候选
    candidate_nodes: Iterable[str] = sorted(triggered_nodes)  # 482 排序保证确定性
elif not checkpoint["channel_versions"]:           # 483 从没写过任何 channel
    candidate_nodes = ()                            # 484 空图,没有候选
else:
    candidate_nodes = processes.keys()              # 486 兜底:老实扫全部节点
updated_channels来自上一步 apply_writes 的返回值(Day 20 的闭环)。它只包含"这一步真正被写、且升了版本"的 channel——通常就一两个。
trigger_to_nodes编译期就建好的反查表:{channel名: [订阅它的节点名]}。有了它,"哪些 channel 变了"能 O(1) 换算成"哪些节点该醒"。
candidate_nodes = processes.keys()兜底分支:第一步、或没有反查表信息时,退化成遍历全部节点。安全但慢——所以正常运行几乎总走上面的快路径。
🎯 设计取舍①:为什么留一条"扫全图"的兜底路,而不是永远只走增量? 增量路径依赖两个前提:① 上一步告诉了我 updated_channels;② 反查表存在。可第一个超步还没有"上一步",从 checkpoint 恢复时也可能拿不到 updated_channels。这时若强行走增量会漏掉本该触发的节点。源码的选择是:正常态走增量(快),边缘态退化成全扫(慢但对)。用一个 elif/else 把"正确性"兜底,把"性能"留给常态——这是工程里典型的"快路径 + 慢路径"设计。
📝 真实值走一遍 假设图有 30 个节点,上一步只有 messages 通道被更新:updated_channels={"messages"}trigger_to_nodes={"messages": ["agent"]}。增量路径直接得到 candidate_nodes=["agent"]——只需检查 1 个节点,而不是 30 个。图越大,这个优化省得越多。
L04

_triggers:一个节点到底该不该醒

候选节点只是"可能要跑",真正的判定在 prepare_single_task 的 PULL 分支里调 _triggers_algo.py:597-612):

# _algo.py:597
elif task_path[0] == PULL:
    name = cast(str, task_path[1])
    proc = processes[name]
    if _triggers(                                    # 606 真正判定"该不该跑"
        channels,
        checkpoint["channel_versions"],              # 608 每个 channel 现在第几版
        checkpoint["versions_seen"].get(name),       # 609 这个节点上次看到第几版
        checkpoint_null_version, proc,
    ):
        triggers = tuple(sorted(proc.triggers))      # 613 命中触发的通道,排序
        ...造 task...

核心比较就在 _triggers_algo.py:1260):

# _algo.py:1260
def _triggers(channels, versions, seen, null_version, proc) -> bool:
    if seen is None:                                 # 1267 从没跑过:只要订阅的通道有值就触发
        for chan in proc.triggers:
            if channels[chan].is_available():
                return True
    else:
        for chan in proc.triggers:                   # 1272 跑过:比版本号
            if channels[chan].is_available() and versions.get(
                chan, null_version
            ) > seen.get(chan, null_version):        # 1273 通道现版本 > 我已见版本 → 触发
                return True
    return False                                     # 1277 都没新版本 → 不跑
seen is None节点从没跑过(versions_seen 里没它)。规则简单:只要它订阅的任一通道"有值"(is_available)就触发。这让入口节点在第一步能被 START 通道拉起来。
versions.get(chan) > seen.get(chan)整个 Pregel 的触发核心一行:通道当前版本号 > 该节点"已见"版本号,说明这个通道在该节点上次运行后又被写过——有新数据,触发。反之(相等)说明没新东西,跳过。
is_available()额外守卫:通道当前得"有值"才算数。空通道即便版本变了也不触发——避免拿空值去跑节点。
数据结构:版本号如何驱动"触发判定" checkpoint["channel_versions"] 通道当前版本(apply_writes 升) messages → 5 counter → 3 versions_seen["agent"] agent 节点已见版本 messages → 4 counter → 3 _triggers: 逐通道比较 messages: 5 > 4 触发! counter: 3 = 3 不动 任一通道命中即返回 True → agent 这步要跑
图注:channel_versions(通道现版本)对比 versions_seen(节点已见版本),任一通道"现版本 > 已见"即触发该节点。
🎯 设计取舍②:为什么用"版本号比较"而不是"内容 diff"? 朴素做法:存下上一步的 state,这一步跑完再逐字段比较"变没变"。问题是——大 state 存两份费内存,逐字段深比较费时间,还分不清"写了个相同的值"算不算变(有些 reducer 就算写相同值也要触发下游)。版本号的做法把"变化"抽象成一个单调递增的整数比较:谁写过谁就升版本,比较 O(1),且天然记录"发生过写入"这一事实而非"值变没变"。代价是要额外维护 channel_versions / versions_seen 两张表——但这点记账换来的是可扩展、可持久化(版本号能存进 checkpoint,恢复后照样比)的触发机制。
L05

PUSH 任务:Send 动态派发的落点

回到 L02 那段 PUSH 消费(_algo.py:442-466)。它处理的是"上一步某节点调了 Send"产生的动态任务:

# _algo.py:442
tasks_channel = cast(Topic[Send] | None, channels.get(TASKS))
if tasks_channel and tasks_channel.is_available():
    for idx, _ in enumerate(tasks_channel.get()):   # 444 TASKS 通道里每个 Send 一条
        if task := prepare_single_task(
            (PUSH, idx), None,                        # 446 用下标 idx 标识第几个 Send
            checkpoint=checkpoint, channels=channels, step=step, stop=stop,
            for_execution=for_execution, ...
        ):
            tasks.append(task)
channels.get(TASKS)TASKS 是个保留通道,类型是 Topic[Send](发布订阅式,能累积多个值)。节点里写的每个 Send(node, arg),最终都被 ChannelWrite(Day 23)塞进这里。
enumerate(tasks_channel.get())把通道里攒的 Send 一条条取出,用下标 idx 区分。map-reduce 里 fan-out 出 5 个子任务,这里就循环 5 次,各造一个 PUSH 任务。
(PUSH, idx)路径元组标记为 PUSH。prepare_single_task 看到 PUSH 会走另一套逻辑:不做版本比较(PUSH 天生就是"要跑"的),直接把 Send 携带的参数当输入造任务。
💡 本质:PUSH 不看版本号——它是"命令式"的,来了就跑 PULL 是"声明式"的(订阅了谁、谁更新就醒);PUSH 是"命令式"的(我直接点名让你跑,带上参数)。所以 PUSH 分支里根本没有 _triggers 那套版本比较——通道里有几个 Send 就造几个任务。这也是为什么同一个节点能在一个超步里被 Send 触发跑很多次(每个 Send 一个独立 task、独立输入),这正是 map-reduce 扇出的实现根基。
⚠️ 边界/坑:Send 出去的任务在"下一个超步"才跑,不是立刻跑 很多人以为 Send 像函数调用一样当场执行。其实不是——Send 只是往 TASKS 通道写了一条,本超步结束、apply_writes 后,下一次 prepare_next_tasks 才把它捞出来变成任务。这完全符合 BSP:"本步的写入下步才可见"。所以如果你在节点里 Send 然后指望马上拿到结果,会扑空——结果要等下一步。理解这一点,map-reduce 的"扇出一步、汇总下一步"节奏就顺理成章了。
L06

apply_writes:把一步的写入合并进 channel

换到结账端。执行器把本步任务跑完,after_tick 调 apply_writes_algo.py:232)。前半段:排序、更新"已见版本"、按通道分组写入:

# _algo.py:232
def apply_writes(checkpoint, channels, tasks, get_next_version, trigger_to_nodes) -> set[str]:
    tasks = sorted(tasks, key=lambda t: task_path_str(t.path[:3]))   # 256 排序=确定性
    bump_step = any(t.triggers for t in tasks)                       # 259 有真任务才推进步

    # ① 更新每个节点的"已见版本":它触发时看到的那些通道的当前版本
    for task in tasks:                                               # 262
        checkpoint["versions_seen"].setdefault(task.name, {}).update({
            chan: checkpoint["channel_versions"][chan]
            for chan in task.triggers
            if chan in checkpoint["channel_versions"]
        })

    # ② 把所有 task 的写入按通道分组
    pending_writes_by_channel: dict[str, list[Any]] = defaultdict(list)   # 295
    for task in tasks:
        for chan, val in task.writes:                               # 297
            if chan in (NO_WRITES, PUSH, RESUME, INTERRUPT, ...):   # 298 保留通道跳过
                pass
            elif chan in channels:
                pending_writes_by_channel[chan].append(val)         # 309 归到对应通道
            else:
                logger.warning(f"...wrote to unknown channel {chan}, ignoring it.")
sorted(tasks, key=...path[:3])先按任务路径排序。为什么?多个任务写同一个通道时,reducer 的合并顺序必须确定,否则同样输入可能得到不同结果(不可复现)。排序保证每次跑顺序一致。
versions_seen[...].update(...)关键一步:把"这个节点触发时,它订阅的通道各是第几版"记下来。等下面这些通道被写、升了版本,下一步 _triggers 一比"新版 > 已见版",同一个节点就又能被触发——循环图(agent↔tools)就是这么转起来的。
if chan in (NO_WRITES, PUSH, ...)过滤保留通道。NO_WRITES 是"我啥也没写"的占位标记(Day 22 会看到无写入的任务会补一个),RESUME/INTERRUPT 是人在环用的(阶段 7),它们不是普通状态通道,不参与 reducer 合并。
logger.warning(...unknown channel)写到不存在的通道只警告不报错——容错设计。多半是拼错了 state 字段名,图还能继续跑,但你会在日志里看到提醒。
L07

升版本与收尾:updated_channels 怎么算出来

apply_writes 后半段:真正把值 update 进通道、给被写通道升版本、并处理"步末通知"(_algo.py:315-345):

# _algo.py:315
updated_channels: set[str] = set()
for chan, vals in pending_writes_by_channel.items():        # 317
    if channels[chan].update(vals) and next_version is not None:  # 319 交给通道 reducer 合并
        checkpoint["channel_versions"][chan] = next_version # 320 被写→升到新版本
        if channels[chan].is_available():
            updated_channels.add(chan)                      # 323 记入"本步更新集"

# 这步没被写的通道,通知它们"新的一步开始了"(某些通道靠这个清空瞬时值)
if bump_step:                                               # 326
    for chan in channels:
        if channels[chan].is_available() and chan not in updated_channels:
            if channels[chan].update(EMPTY_SEQ) and ...:    # 329 空更新
                ...

# 若本步更新的通道都触发不了任何节点 = 可能是最后一步,通知所有通道"收尾"
if bump_step and updated_channels.isdisjoint(trigger_to_nodes):  # 336
    for chan in channels:
        if channels[chan].finish() and next_version is not None:  # 338
            ...
return updated_channels                                     # 345 交回给 Day 20 的闭环
channels[chan].update(vals)把这个通道收到的多个写入交给它的 reducer 合并(LastValue 取最后一个、BinaryOp 累加、add_messages 智能合并……阶段 5 细讲)。返回 True 表示"确实变了"。
channel_versions[chan] = next_version变了就升版本号。next_version 是"当前所有通道最大版本 +1"(同一步所有被写通道共用一个新版本号),保证版本单调递增。
updated_channels.add(chan)把确实更新的通道收进集合——这就是返回值,回到 Day 20 的 self.updated_channels = apply_writes(...),喂给下一次 tick 做 L03 的增量触发。闭环合拢!
isdisjoint(trigger_to_nodes)"本步更新的通道,和'能触发节点的通道'没有交集"——意味着没人会被下一步触发,图大概率要停了。此时调 finish() 让通道做收尾(比如把汇总值定型),保证最终输出正确。
控制流:apply_writes 的四步 + 回喂闭环 ①记已见版本 versions_seen ②按通道分组 writes_by_channel ③update+升版本 reducer 合并 ④收尾/finish 最后一步定型 return updated_channels → 下一次 prepare_next_tasks 的增量触发(L03)
图注:apply_writes 四步走,返回 updated_channels 回喂给下一步 Plan,与 Day 20 主循环合成闭环。

👶 小白:同一步里两个并行节点都写了 counter,会不会互相覆盖丢数据?

👨‍🏫 老师:不会。看 ② 那步——两个节点的写入都被 appendpending_writes_by_channel["counter"],变成一个列表,③ 再整体交给 counter 通道的 reducer 一次性合并。如果 counter 用的是 BinaryOperatorAggregate(operator.add),两个值就相加;如果是 LastValue,那确实"后来者赢"(这时并行写同一 LastValue 字段本身就是设计错误,Day 10 讲过)。并行写不丢,是因为写入先攒成列表再统一交给 reducer,而不是各写各的直接盖。

L08

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

🧠 今天你应该能回答

  • prepare_next_tasks 和 apply_writes 各干什么?(前者 Plan:算任务;后者 Update:合并写入+升版本)
  • PULL 和 PUSH 任务的区别?(PULL 声明式、看版本触发;PUSH 命令式、Send 点名就跑)
  • 一个节点"该不该跑"怎么判定?(_triggers:通道现版本 > 该节点已见版本)
  • 增量触发靠什么加速?(updated_channels + trigger_to_nodes,候选集从全图缩到极小)
  • 版本号方案比"内容 diff"好在哪?(O(1) 比较、可持久化、天然记录"发生过写入")
  • 并行写同一通道为何不丢?(先 append 成列表,再整体交 reducer 合并)
  • updated_channels 怎么闭环回去?(apply_writes 返回它 → 下步 Plan 的增量触发用它)

✋ 10 分钟动手

# 1. 读 apply_writes 全貌(合并 + 升版本 + 收尾)
sed -n '232,345p' libs/langgraph/langgraph/pregel/_algo.py

# 2. 读 prepare_next_tasks 骨架(PUSH + PULL + 增量优化)
sed -n '437,513p' libs/langgraph/langgraph/pregel/_algo.py

# 3. 读触发核心 _triggers(版本比较那一行)
sed -n '1260,1277p' libs/langgraph/langgraph/pregel/_algo.py

# 4. 打印看看版本号怎么涨(每步 counter 加 1)
python -c "
from langgraph.graph import StateGraph, START, END
from typing import Annotated; import operator
class S(dict): pass
g = StateGraph(dict)
g.add_node('a', lambda s: {'n': s.get('n',0)+1})
g.add_edge(START,'a'); g.add_edge('a', END)
app = g.compile()
for c in app.stream({'n':0}, stream_mode='debug'):
    print(c['type'], c.get('payload',{}).get('name',''))
"
明天预告 · Day 22:今天 tick 把任务"备"好了,可到底是谁把这些 task 真正"跑"起来的?明天进 _runner.pyPregelRunner.tick()_executor.pyBackgroundExecutor——看它怎么用线程池并行跑多个任务、怎么在一个任务失败时取消其他任务、怎么把每个任务的写入 commit 回去。执行引擎的"炒菜"环节。
← Day 20 超步主循环 Day 22 · 执行器 PregelRunner →