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

执行器:PregelRunner 怎么并行跑一步的任务

Day 20 那句 for _ in runner.tick(...) 就是"炒菜"环节。tick 把 task 备好了,真正把它们并行跑起来、跑完把写入攒回去、有一个崩了就取消其他的——全靠 PregelRunner_runner.py:135)和它底下的线程池 BackgroundExecutor_executor.py:40)。今天看清一个超步内的并发是怎么被安全地组织起来的。

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

谁来炒菜:执行器要解决的四个问题

🤔 痛点:一个超步里可能有好几个节点要跑,它们能并行吗?其中一个报错了,其他的怎么办? Day 19 说"同一超步的节点互不依赖、可以并行"。可"可以并行"是理论,落到代码:谁开线程?开几个?某个任务抛异常了,是等其他跑完还是立刻掐掉?跑完的写入先放哪,什么时候交给 apply_writes?流式输出又怎么做到"每完成一个就吐一个"?这四个问题,就是 PregelRunner 存在的理由。

执行器夹在 Day 20 的 tick 和 after_tick 中间,职责被官方注释说得很清楚(_runner.py:136-138):

  • 并发执行一组 Pregel 任务;
  • 把它们的写入 commit 回去(攒到 checkpoint 的 pending_writes);
  • 有输出可吐时把控制权 yield 给外层(流式);
  • 必要时 中断/取消其他任务(一个失败,全部叫停)。
💡 本质:执行器 = "并发编排器",节点逻辑本身它不碰 PregelRunner 不关心节点里写了什么业务代码——那是节点自己的事。它只管"把这批任务丢进线程池、盯着谁先完成、把结果收好、出事就止损"。这种"编排与业务分离"让引擎能对任何节点一视同仁地做并发、超时、重试。
类比:后厨的传菜长 tick 是配菜台(把每道菜的料备齐),PregelRunner 是传菜长:把这一轮的菜同时派给几个炒锅(线程)一起炒(并行),哪个锅先出菜就先端上桌(yield 流式),某个锅打翻着火了就喊停其他锅别炒了(失败即止),所有菜炒完把账单汇总交给收银台(commit → after_tick)。传菜长不亲自炒菜,只负责调度。
L02

PregelRunner:拿两个"弱引用"就够干活

看它的 __init___runner.py:140-169)。最关键是两个参数都用 weakref.ref 包着:

# _runner.py:135
class PregelRunner:
    """Responsible for executing a set of Pregel tasks concurrently, committing
    their writes, yielding control to caller when there is output to emit, and
    interrupting other tasks if appropriate."""
    def __init__(self, *,
        submit: weakref.ref[Submit],                                   # 143 拿到线程池的提交入口
        put_writes: weakref.ref[Callable[[str, Sequence[...]], None]], # 144 把写入交回 loop
        use_astream: bool = False,
        node_finished: Callable[[str], None] | None = None,
        node_error_handler_map: Mapping[str, str] | None = None,       # 147 节点→错误处理节点
        ...
    ) -> None:
        self.submit = submit
        self.put_writes = put_writes
        self._handled_exception_ids: set[int] = set()                  # 169 已被错误处理器接管的异常
submit: weakref.ref[Submit]提交任务到线程池的函数(来自 BackgroundExecutor,L07)。用时 self.submit()(fn, ...)——先解弱引用拿到真函数,再调用。
put_writes: weakref.ref[...]把某个 task 的写入交回 PregelLoop 的回调。commit 时(L06)调它,写入就进了 checkpoint 的 pending_writes,等 after_tick 的 apply_writes 来合并(Day 21)。
node_error_handler_map节点级容错:某节点失败时路由到指定的"错误处理节点",而不是直接崩掉整个图(Day 52 深入)。这里先记住有这么个机制。
🎯 设计取舍①:为什么 submit / put_writes 要用 weakref 弱引用? PregelRunnerPregelLoop 是同一次运行里彼此持有的两个对象:loop 持有 runner,runner 又要回调 loop 的方法。如果 runner 用强引用抓着 loop 的方法,就形成"你抓我、我抓你"的引用环,Python 的引用计数无法回收,得等 GC 兜底——在长期运行的服务里容易攒出内存。用 weakref 让 runner"能用但不占有"这些回调,运行一结束 loop 就能被立即回收。代价是每次用都要 self.submit() 先解引用、多一层调用;换来的是干净的生命周期。
L03

tick 开场:FuturesDict 与"先让出控制权"

tick()_runner.py:176)是执行器的主入口。开头先建一个特殊的 FuturesDict,然后 yield 一次:

# _runner.py:176
def tick(self, tasks, *, reraise=True, timeout=None, retry_policy=None,
         get_waiter=None, schedule_task) -> Iterator[None]:
    tasks = tuple(tasks)
    futures = FuturesDict(                          # 190 {future: task} 的特殊字典
        callback=weakref.WeakMethod(self.commit),   # 191 每个 future 完成自动回调 commit
        event=threading.Event(),
        should_stop=partial(_should_stop_others, ...),  # 193 判断"要不要停其他任务"
        future_type=concurrent.futures.Future,
    )
    yield                                           # 199 先把控制权还给调用方
    if len(tasks) == 0:
        return                                      # 201 没任务,直接结束
Iterator[None]tick 是个生成器!所以 Day 20 用 for _ in runner.tick(...) 驱动它。每 yield 一次,外层就得到一次"该看看有没有新输出可以流式吐出了"的机会。
FuturesDict(callback=WeakMethod(self.commit))一个把 future 映射到 task 的字典,且绑了完成回调:任何 future 一完成,自动调 commit 把该 task 的写入收好。不用手动轮询"谁好了"。
yield(第一次)还没开始跑就先让出控制权。这是流式的礼貌——让外层先处理上一步遗留的输出、注册好监听,再回来真正调度任务。
💡 本质:tick 是"可暂停的调度器",靠 yield 和外层协作流式 普通函数是"进去跑完出来";生成器 tick 是"跑一段、yield 出去让外层吐输出、外层再驱动我跑下一段"。这种"协程式"结构,正是 LangGraph 能做到"节点一完成就实时流式"而不是"整步跑完才一次性返回"的底层原因。
数据结构:PregelRunner 的记账(弱引用 + FuturesDict) PregelRunner submit(weakref) put_writes(weakref) FuturesDict {future → task} callback=commit(完成即回调) BackgroundExecutor ThreadPoolExecutor 真正开线程跑 submit() runner 用弱引用拿到线程池与回调 → 提交任务得 future → FuturesDict 记账、完成自动 commit 弱引用避免 loop↔runner 引用环,运行结束即可回收
图注:PregelRunner 靠两个弱引用与 FuturesDict 记账,把任务提交给线程池并在完成时自动 commit。
L04

单任务快路径:不开线程,当场跑

如果这一步只有一个任务、没超时、没 waiter,开线程池纯属浪费。tick 有一条快路径_runner.py:203-221):

# _runner.py:203
elif len(tasks) == 1 and timeout is None and get_waiter is None:
    t = tasks[0]
    try:
        run_with_retry(                             # 207 直接在当前线程跑(含重试,Day 25)
            t, retry_policy,
            configurable={CONFIG_KEY_CALL: partial(_call, weakref.ref(t), ...)},  # 210
        )
        self.commit(t, None)                        # 221 成功:commit(无异常)
    except Exception as exc:
        self.commit(t, exc)                         # 223 失败:也 commit(带异常,存错误)
        ...错误处理器路由 / reraise...
len(tasks) == 1 and ...命中条件:单任务 + 无超时 + 无 waiter。绝大多数线性图(A→B→C,每步就一个节点)都走这条路。
run_with_retry(t, ...)直接在当前线程同步跑,不提交线程池。run_with_retry 是 Day 25 的主角,负责按 retry_policy 重试。
self.commit(t, None) / commit(t, exc)无论成功失败都调 commit——成功传 None,失败传异常。commit 决定"把正常写入存起来,还是把错误/中断存起来"(L06)。
🎯 设计取舍②:为什么专门为"单任务"写一条不开线程的快路径? 线程池的 submit + Future + concurrent.futures.wait 有实打实的开销:上下文拷贝、线程调度、future 状态机。可现实里大量图是线性的——每一步就一个节点。为这种最常见的情况开线程池,等于"为了端一杯水叫一辆卡车"。快路径检测到"就一个任务"时当场同步跑,省掉全部并发开销;只有真需要并行(多任务)或需要超时/等待时才动用线程池(L05)。常见情况优化到极致,复杂情况才付复杂的代价。
L05

并行调度循环:提交 + 等最先完成的那个

多任务时走通用路径:全部提交线程池,然后循环"等最先完成的一个、处理它、再等下一个"(_runner.py:259-289):

# _runner.py:259
for t in tasks:
    fut = self.submit()(                            # 260 提交到线程池,拿到 future
        run_with_retry, t, retry_policy,
        configurable={CONFIG_KEY_CALL: partial(_call, weakref.ref(t), ...)},
        __reraise_on_exit__=reraise,
    )
    futures[fut] = t                                # 276 记账:这个 future 对应哪个 task
# 每有一个任务完成就吐一次输出
end_time = timeout + time.monotonic() if timeout else None   # 280 墙钟截止时刻
while len(futures) > (1 if get_waiter is not None else 0):    # 282
    done, inflight = concurrent.futures.wait(
        futures, return_when=concurrent.futures.FIRST_COMPLETED,   # 285 等"第一个完成"
        timeout=(max(0, end_time - time.monotonic()) if end_time else None),
    )
    if not done:
        break                                       # 289 一个都没完成=超时了,跳出
self.submit()(run_with_retry, t, ...)把每个任务提交进线程池,立刻返回一个 Future(占位凭据,代表"将来会有结果")。N 个任务=N 个 future,它们同时在不同线程跑。
FIRST_COMPLETED不等所有任务、只等最先完成的一个就醒来。这样某个任务先跑完,就能马上处理它、yield 流式输出,而不用干等最慢的那个。
timeout=end_time - nowDay 20 传下来的 step_timeout 在这里生效:整个超步的墙钟上限。每轮 wait 都算"还剩多少时间"。
if not done: breakwait 返回空 done 集,说明超时到了还没人完成——直接跳出循环,后续 _panic_or_proceed 会按超时处理。
控制流:多任务并行 + FIRST_COMPLETED 逐个收割 submit 全部任务 线程1: node A 线程2: node B 线程3: node C wait(FIRST_ COMPLETED) commit + yield 输出 循环回去等下一个完成,直到 futures 空 任一任务抛错 → _should_stop_others 为真 → break 停掉其余
图注:多任务全部提交线程池,主循环用 FIRST_COMPLETED 逐个收割、边收边流式;出错则止损其余。
L06

失败即止 + commit:写入/错误/中断三条路

每轮收割后判断要不要停其他任务(_runner.py:330-335 + _runner.py:616):

# _runner.py:330
if _should_stop_others(done_for_stop, handled_exception_ids=self._handled_exception_ids):
    break                                           # 333 有任务真失败 → 停掉其余
yield                                               # 335 否则让出控制权,流式吐输出

# _runner.py:616
def _should_stop_others(done, *, handled_exception_ids=None) -> bool:
    """Check if any task failed, if so, cancel all other tasks.
    GraphInterrupts are not considered failures."""
    for fut in done:
        if fut.cancelled():
            continue
        elif exc := fut.exception():                # 626 这个 future 抛异常了?
            if (id(exc) not in (handled_exception_ids or set())
                and not isinstance(exc, GraphBubbleUp)):   # 629 中断/续跑信号不算失败
                return True
    return False

而每个任务完成后(无论成败)由 commit 归档写入(_runner.py:574):

# _runner.py:574
def commit(self, task, exception) -> None:
    if isinstance(exception, asyncio.CancelledError):
        task.writes.append((ERROR, exception)); self.put_writes()(task.id, task.writes)
    elif exception:
        if isinstance(exception, GraphInterrupt):   # 585 是中断 → 存 INTERRUPT 写入
            ... self.put_writes()(task.id, [(INTERRUPT, exception.args[0]), *resumes])
        elif isinstance(exception, GraphBubbleUp):  # 592 冒泡信号 → 交给上层处理
            pass
        else:                                       # 595 真错误 → 存 ERROR 写入
            task.writes.append((ERROR, exception)); self.put_writes()(task.id, task.writes)
    else:                                           # 604 成功
        if not task.writes:
            task.writes.append((NO_WRITES, None))   # 611 没写任何东西 → 补一个占位标记
        self.put_writes()(task.id, task.writes)     # 613 把写入交回 loop
GraphInterrupt 不算失败_should_stop_others 明确排除 GraphBubbleUp(GraphInterrupt 的父类)。人在环的中断是"计划内暂停"不是"出错",所以不该连累其他并行任务被取消。
commit 里三条分支正常→存 writes;中断→存 INTERRUPT 写入(供恢复用,阶段 7);真错误→存 ERROR 写入。无论哪种,写入都进 checkpoint——这样即便崩了,"崩在哪、崩之前写了啥"也都存了档,能恢复。
NO_WRITES 占位任务成功但啥也没写(比如纯副作用节点),也要补一个 NO_WRITES 标记。为什么?Day 21 apply_writes 里 bump_step 靠"有没有 task"判断要不要推进超步——这个占位保证"我确实跑过了"这一事实被记录。
⚠️ 边界/坑:一个并行节点崩了,同批其他节点的写入会不会白跑? 不会白跑,但要看时机。_should_stop_others 返回 True 只 break还没完成的任务;已经完成的任务,它们的 commit 已经把写入存进 checkpoint 了。所以从崩溃点恢复时,那些已成功的节点不必重跑(Day 46 幂等恢复)。真正的坑是反过来:如果你的节点有外部副作用(比如已经发了一封邮件),而同批另一个节点崩了导致整步回滚重来,副作用会重复发生——引擎只能保证"状态写入"的幂等,管不了你的外部 IO。所以副作用节点要自己做幂等键。
L07

BackgroundExecutor:线程池与"退出时结账"

上面的 self.submit() 从哪来?来自 BackgroundExecutor_executor.py:40)——一个包着线程池的上下文管理器:

# _executor.py:54
def submit(self, fn, *args,
    __cancel_on_exit__=False,     # 59 退出时若还没启动,可取消
    __reraise_on_exit__=True,     # 60 退出时把任务里的异常重新抛出
    __next_tick__=False, **kwargs,
) -> concurrent.futures.Future[T]:
    ctx = copy_context()                            # 64 拷贝上下文(ContextVar 隔离)
    task = self.executor.submit(ctx.run, fn, *args, **kwargs)   # 71 丢进线程池
    self.tasks[task] = (__cancel_on_exit__, __reraise_on_exit__)
    task.add_done_callback(self.done)               # 74 完成即从 tasks 里摘除
    return task

# _executor.py:93 —— 退出上下文时(一步跑完)
def __exit__(self, exc_type, exc_value, traceback):
    tasks = self.tasks.copy()
    for task, (cancel, _) in tasks.items():
        if cancel:
            task.cancel()                           # 104 取消标了 cancel_on_exit 的
    if pending := {t for t in tasks if not t.done()}:
        concurrent.futures.wait(pending)            # 107 等所有未完成的跑完
    self.stack.__exit__(exc_type, exc_value, traceback)   # 109 关线程池
    if exc_type is None:                            # 111 没有正在传播的异常时
        for task, (_, reraise) in tasks.items():    # 113 才把任务内的异常重新抛出
copy_context() + ctx.run每个任务在独立拷贝的上下文里跑。这样一个节点里改了 ContextVar(比如 config、tracing 上下文)不会污染并行的其他节点——线程安全的关键。
add_done_callback(self.done)任务一完成自动从 self.tasks 字典摘掉,顺手把 GraphBubbleUp(中断信号)吞掉不当错误。自动清账,不用手动管理。
__exit__ 里 wait(pending)退出上下文(一步结束)时阻塞等所有未完成任务跑完。保证不会有"半跑着的僵尸线程"泄漏到下一步——超步的边界是硬的。
if exc_type is None: reraise只有当前没有别的异常在传播时,才把任务里攒的异常重新抛出。避免"异常覆盖异常"丢失原始错误。

👶 小白:async 图也是用这个线程池吗?

👨‍🏫 老师:不是。这是同步BackgroundExecutor(走线程池 ThreadPoolExecutor)。异步图用的是紧挨着的 AsyncBackgroundExecutor_executor.py:122),它不开线程,而是把每个任务变成 asyncio 协程任务,跑在同一个事件循环上,还支持 max_concurrency 用信号量限流。两者接口一样(都提供 submit),所以上层 PregelRunner 的 tick/atick 逻辑几乎对称。同步靠多线程并行、异步靠事件循环并发——这是 Python 并发的两条路,LangGraph 两条都铺了。

L08

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

🧠 今天你应该能回答

  • PregelRunner 解决哪四个问题?(并发执行 / commit 写入 / yield 流式 / 失败即止)
  • submit / put_writes 为什么用 weakref?(避免 loop↔runner 引用环,及时回收内存)
  • 单任务快路径省了什么?(不开线程池,当前线程同步跑,省并发开销)
  • 多任务怎么调度?(全提交线程池,FIRST_COMPLETED 逐个收割、边收边 yield)
  • 为什么 GraphInterrupt 不触发"停其他任务"?(中断是计划内暂停,不是失败)
  • commit 的三条分支?(正常存 writes / 中断存 INTERRUPT / 错误存 ERROR,都进 checkpoint)
  • 同步和异步执行器分别靠什么并发?(线程池 vs 事件循环协程)

✋ 10 分钟动手

# 1. 读 PregelRunner.tick 主体(快路径 + 并行循环)
sed -n '176,290p' libs/langgraph/langgraph/pregel/_runner.py

# 2. 读 commit 三条分支
sed -n '574,613p' libs/langgraph/langgraph/pregel/_runner.py

# 3. 读线程池 BackgroundExecutor
sed -n '40,120p' libs/langgraph/langgraph/pregel/_executor.py

# 4. 亲眼看并行:两个节点同一步跑,观察它们的时间戳交错
python -c "
import time
from langgraph.graph import StateGraph, START, END
from typing import Annotated; import operator
def make(name):
    def f(s):
        print(name,'start',time.time()); time.sleep(0.5); print(name,'end',time.time())
        return {'log': [name]}
    return f
class S(dict): pass
g = StateGraph(dict)
g.add_node('a', make('A')); g.add_node('b', make('B'))
g.add_edge(START,'a'); g.add_edge(START,'b'); g.add_edge('a',END); g.add_edge('b',END)
g.compile().invoke({})   # A、B 同一超步并行,start 几乎同时
"
明天预告 · Day 23:执行器跑的其实不是"你的函数"本身,而是包了一层的 PregelNode——它负责"跑之前读好输入、跑之后把输出写进 channel"。明天进 _read.pyPregelNode_write.pyChannelWrite,看清一个节点的"读—算—写"三段式是怎么被组装出来的,以及 Send 到底怎么变成 TASKS 通道里的一条写入。
← Day 21 任务准备 Day 23 · PregelNode 读写 →