Skip to content

PregelRunner:超步内的任务执行器

源码版本1.2.9

职责

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_mapschedule_error_handler/aschedule_error_handler 回调,以及 _handled_exception_ids 这个跨 tick 的异常去重集合。
  • FuturesDict:75-134 — 自定义 dict 子类,on_done 回调在每个 future 完成时触发 commit;should_stop 字段持有 _should_stop_others partial。
  • 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-323concurrent.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_WRITES marker。
  • _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-174task.name in self.node_error_handler_map 时返回 True;error handler 节点本身不能被路由(防递归)。

数据流

tick 的开头先构造 FuturesDict,这是个增强版 dict,每个 future 完成时自动调 commit:

python
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:

python
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-613put_writesPregelLoop.put_writes 的 weakref(put_writes:415),它把 writes 塞进 checkpoint_pending_writes 并起 checkpointer 的 future 持久化。注意几种特殊 writes:

  • CancelledError → 写 (ERROR, exception),让 after_tick 收尾时把 task 算成「失败但已记录」,不 panic。
  • GraphInterrupt → 写 (INTERRUPT, ...) + 任何 RESUME writes,中断的恢复值要和 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 不算失败:commitGraphBubbleUp 显式 pass,不写任何 writes(GraphBubbleUp pass:592-594)。这类异常(如 ParentCommand)是给父图发信号的,不是真错误。
  • NO_WRITES 防止 task 重复执行:task 正常完成但没写任何东西,也要补个 (NO_WRITES, None)(NO_WRITES:609-611)。否则 after_tickprepare_next_tasks 看到这个 task 没 writes 会以为没跑过,续跑时再跑一次。
  • schedule_error_handler 可能返回 None:handler 节点自己也可能没配置或无法构造,这时 task 就走普通 panic 路径(handler optional:230-248)。commit 里的 ERROR_SOURCE_NODE write 会和这条路径并存,_resume_error_handlers_if_applicable 续跑时再尝试调度。
  • timeout 撞了不抛 GraphInterrupt:_panic_or_proceedif 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