PregelRunner:超步内的任务执行器
职责
PregelRunner 干一件事:把 PregelLoop.tick 准备好的那批 PregelExecutableTask 真正跑起来,把它们的 writes 收齐交回 PregelLoop。它不关心图长什么样,不关心 checkpoint,不关心通道怎么 reduce——它只负责「并发执行 + 失败处理 + writes 落盘」这三件事。
它的位置非常窄:Pregel.astream 的主循环里(sync main loop:2964-2984),while loop.tick() 返回 True 之后,紧跟着就是 for _ in runner.tick([t for t in loop.tasks.values() if not t.writes], ...):,跑完再 loop.after_tick()。也就是说,runner 的输入是「tick 算出来的、还没写过 writes 的 task」,输出是「这些 task 的 writes 已经进了 checkpoint_pending_writes」,中间没有第二层缓冲。
PregelRunner 自身非常薄:两个公开方法 tick(tick:176)和 commit(commit:574),外加一个 atick 异步版。重活全在 _call / _acall(_call:700)/(_acall:789)这两个底层调用器里——它们负责把节点函数包成可调度的 future,处理 retry policy、子 task Send 扇出、streaming chunk 回调,最终把 writes 写回 task 对象。
设计动机
为什么 runner 要做成 generator 而不是普通方法?
- 流式 yield:runner 是
for _ in runner.tick(...)这种 generator 形态(yield return type:188),每跑完一个 task 就 yield 一次,把控制权交还上层——上层就能立刻把 streaming chunk 发给客户端,而不是等整个超步跑完。这是 LangGraph 流式 (streaming) 的底层基础。 - 单 task 快路径:
len(tasks) == 1 and timeout is None and get_waiter is None时走同步直接调用,不起线程池 (fast path:203-254)。绝大多数 agent 单分支场景命中这条路径,避免线程池开销。 commit作为 Future 回调:FuturesDict.on_done(on_done:116)在每个 future 完成时被触发,内部就是weakref.WeakMethod(self.commit)——commit是 per-task 的 writes 落盘钩子,而不是 batch 收尾。这样 writes 进pending_writes的顺序天然按完成顺序,checkpoint 也能尽早写。- error handler 路由:
_should_route_to_error_handler(_should_route_to_error_handler:171-174)判断失败节点是否配了error_handler节点。命中就把 exception id 加进_handled_exception_ids,调度一个 handler task 替代原 task——这样失败不会立即 panic,而是走另一个节点收尾。 _should_stop_others取消机制:任何一个 task 抛非GraphBubbleUp异常时,runner 取消同批其它 in-flight task(_should_stop_others:616-634)——保证同超步内「要么全成,要么进 handler」的语义。GraphInterrupt被显式排除在外,因为中断不算失败。
关键文件
class PregelRunner:135-138— 类定义,docstring 一句话点明职责:执行 task、commit writes、yield 控制权、必要时中断其它 task。__init__:140-169— 持有submit/put_writes的 weakref、node_error_handler_map、schedule_error_handler/aschedule_error_handler回调,以及_handled_exception_ids这个跨 tick 的异常去重集合。FuturesDict:75-134— 自定义 dict 子类,on_done回调在每个 future 完成时触发commit;should_stop字段持有_should_stop_otherspartial。tick signature:176-188— 输入是Iterable[PregelExecutableTask],返回Iterator[None],signature 暴露 generator 语义。fast path:203-254— 单 task + 无 timeout + 无 waiter 时直接同步跑,失败时按需调度 error handler。schedule tasks:259-276— 多 task 路径:每个 task 起一个 future,通过self.submit()(实际是PregelLoop.submit)调度。concurrent wait loop:282-323—concurrent.futures.wait(FIRST_COMPLETED)循环,每完成一个 task 就commit、emit 输出、必要时起 handler task。commit method:574-613— per-task 收尾:cancelled 写 ERROR、GraphInterrupt写 INTERRUPT、普通 exception 写ERROR+ERROR_SOURCE_NODE、正常写task.writes+NO_WRITESmarker。_should_stop_others:616-634— 判断是否要取消其它 task,显式排除GraphBubbleUp和已 handled 的异常。_panic_or_proceed:650-697— 收尾函数,取消所有 in-flight future,合并多个GraphInterrupt成一个,timeout 时抛TimeoutError。_call:700-787— 节点调用的底层包装:retry、stream chunk emit、Send扇出回调schedule_task,把 task.writes 累积起来。_should_route_to_error_handler:171-174—task.name in self.node_error_handler_map时返回 True;error handler 节点本身不能被路由(防递归)。
数据流
tick 的开头先构造 FuturesDict,这是个增强版 dict,每个 future 完成时自动调 commit:
def tick(
self,
tasks: Iterable[PregelExecutableTask],
*,
reraise: bool = True,
timeout: float | None = None,
retry_policy: Sequence[RetryPolicy] | None = None,
get_waiter: Callable[[], concurrent.futures.Future[None]] | None = None,
schedule_task: Callable[
[PregelExecutableTask, int, Call | None],
PregelExecutableTask | None,
],
) -> Iterator[None]:
tasks = tuple(tasks)
futures = FuturesDict(
callback=weakref.WeakMethod(self.commit),
event=threading.Event(),
should_stop=partial(
_should_stop_others, handled_exception_ids=self._handled_exception_ids
),
future_type=concurrent.futures.Future,
)
# give control back to the caller
yield这段来自 tick 头部:176-199。注意第一个 yield 立刻把控制权交还——上层 for _ in runner.tick(...) 第一次迭代只是「启动」,真正的 task 执行发生在后续迭代中。这是 generator 协程的一种简化形态,避免引入 async 关键字到同步路径。
commit 是 per-task 的收尾函数,把 task 的 writes 落进 pending_writes:
def commit(
self,
task: PregelExecutableTask,
exception: BaseException | None,
) -> None:
if isinstance(exception, asyncio.CancelledError):
task.writes.append((ERROR, exception))
self.put_writes()(task.id, task.writes)
elif exception:
if isinstance(exception, GraphInterrupt):
if exception.args[0]:
writes = [(INTERRUPT, exception.args[0])]
if resumes := [w for w in task.writes if w[0] == RESUME]:
writes.extend(resumes)
self.put_writes()(task.id, writes)
elif isinstance(exception, GraphBubbleUp):
pass
else:
task.writes.append((ERROR, exception))
if self._should_route_to_error_handler(task) and not isinstance(
exception, GraphBubbleUp
):
task.writes.append((ERROR_SOURCE_NODE, task.name))
self._handled_exception_ids.add(id(exception))
self.put_writes()(task.id, task.writes)
else:
if self.node_finished and (
task.config is None or TAG_HIDDEN not in task.config.get("tags", [])
):
self.node_finished(task.name)
if not task.writes:
task.writes.append((NO_WRITES, None))
self.put_writes()(task.id, task.writes)这段来自 commit:574-613。put_writes 是 PregelLoop.put_writes 的 weakref(put_writes:415),它把 writes 塞进 checkpoint_pending_writes 并起 checkpointer 的 future 持久化。注意几种特殊 writes:
CancelledError→ 写(ERROR, exception),让after_tick收尾时把 task 算成「失败但已记录」,不 panic。GraphInterrupt→ 写(INTERRUPT, ...)+ 任何RESUMEwrites,中断的恢复值要和 interrupt 一起落盘,下次 resume 时一起读出来。- 普通 exception + 配了 error handler → 多写一条
(ERROR_SOURCE_NODE, task.name),这条标记会被PregelLoop._resume_error_handlers_if_applicable扫到(ERROR_SOURCE_NODE scan:775),用来在续跑时调度对应的 handler 节点。 - 正常完成但没写任何东西 → 补一条
(NO_WRITES, None),防止prepare_next_tasks把它当成「还没跑过」(NO_WRITES marker:609-611)。
整个 runner 的执行流可以这样画:
边界与失败
_handled_exception_ids跨 tick 持久:__init__里这个集合是 instance-level,不是 per-tick(_handled_exception_ids:169)。同一个异常对象 id 进了集合就不会再被_should_stop_others当成失败,避免 error handler 跑完后又把原 exception 重抛。GraphBubbleUp不算失败:commit对GraphBubbleUp显式pass,不写任何 writes(GraphBubbleUp pass:592-594)。这类异常(如ParentCommand)是给父图发信号的,不是真错误。NO_WRITES防止 task 重复执行:task 正常完成但没写任何东西,也要补个(NO_WRITES, None)(NO_WRITES:609-611)。否则after_tick后prepare_next_tasks看到这个 task 没 writes 会以为没跑过,续跑时再跑一次。schedule_error_handler可能返回 None:handler 节点自己也可能没配置或无法构造,这时 task 就走普通 panic 路径(handler optional:230-248)。commit里的ERROR_SOURCE_NODEwrite 会和这条路径并存,_resume_error_handlers_if_applicable续跑时再尝试调度。- timeout 撞了不抛 GraphInterrupt:
_panic_or_proceed的if inflight: raise timeout_exc_cls("Timed out")(timeout:691-697)——timeout 是真异常,会冒泡出 Pregel,不像 GraphInterrupt 那样被吞。 - traceback 修剪:reraise 之前会跳过
EXCLUDED_FRAME_FNAMES列表里的帧(tb trim:241-247),把 langgraph 内部栈帧过滤掉,让用户看到的 traceback 直接指向他们的节点代码。
小结
PregelRunner 是个薄壳:generator tick 调度并发执行 + per-task commit 把 writes 落进 pending_writes,中间塞了 error handler 路由和取消机制。理解了它,就理解了「一个超步里 task 是怎么真跑起来的」。上一层 PregelLoop 的循环驱动见 /pregel/loop;prepare_next_tasks / apply_writes / should_interrupt 这三个算法函数见 /pregel/algo;整体组装见 /pregel/pregel;流式输出怎么在 runner yield 之间被加工见 /stream/run-stream。
对照官方资料:LangGraph 文档 · README。