执行器:PregelRunner 怎么并行跑一步的任务
Day 20 那句 for _ in runner.tick(...) 就是"炒菜"环节。tick 把 task 备好了,真正把它们并行跑起来、跑完把写入攒回去、有一个崩了就取消其他的——全靠 PregelRunner(_runner.py:135)和它底下的线程池 BackgroundExecutor(_executor.py:40)。今天看清一个超步内的并发是怎么被安全地组织起来的。
谁来炒菜:执行器要解决的四个问题
PregelRunner 存在的理由。执行器夹在 Day 20 的 tick 和 after_tick 中间,职责被官方注释说得很清楚(_runner.py:136-138):
- 并发执行一组 Pregel 任务;
- 把它们的写入 commit 回去(攒到 checkpoint 的 pending_writes);
- 有输出可吐时把控制权 yield 给外层(流式);
- 必要时 中断/取消其他任务(一个失败,全部叫停)。
PregelRunner 不关心节点里写了什么业务代码——那是节点自己的事。它只管"把这批任务丢进线程池、盯着谁先完成、把结果收好、出事就止损"。这种"编排与业务分离"让引擎能对任何节点一视同仁地做并发、超时、重试。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 深入)。这里先记住有这么个机制。PregelRunner 和 PregelLoop 是同一次运行里彼此持有的两个对象:loop 持有 runner,runner 又要回调 loop 的方法。如果 runner 用强引用抓着 loop 的方法,就形成"你抓我、我抓你"的引用环,Python 的引用计数无法回收,得等 GC 兜底——在长期运行的服务里容易攒出内存。用 weakref 让 runner"能用但不占有"这些回调,运行一结束 loop 就能被立即回收。代价是每次用都要 self.submit() 先解引用、多一层调用;换来的是干净的生命周期。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(第一次)还没开始跑就先让出控制权。这是流式的礼貌——让外层先处理上一步遗留的输出、注册好监听,再回来真正调度任务。单任务快路径:不开线程,当场跑
如果这一步只有一个任务、没超时、没 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)。常见情况优化到极致,复杂情况才付复杂的代价。并行调度循环:提交 + 等最先完成的那个
多任务时走通用路径:全部提交线程池,然后循环"等最先完成的一个、处理它、再等下一个"(_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 会按超时处理。失败即止 + 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。所以副作用节点要自己做幂等键。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 两条都铺了。
今日小结 + 动手 + 明日预告
🧠 今天你应该能回答
- 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 几乎同时
"
PregelNode——它负责"跑之前读好输入、跑之后把输出写进 channel"。明天进 _read.py 的 PregelNode 和 _write.py 的 ChannelWrite,看清一个节点的"读—算—写"三段式是怎么被组装出来的,以及 Send 到底怎么变成 TASKS 通道里的一条写入。