发布日期

第 12 讲 · PregelRunner 与并发执行模型

学习目标

  • 看懂 PregelRunner.tick/atick 如何并行调度任务、commit 如何回收结果。
  • 理解 BackgroundExecutor / AsyncBackgroundExecutor 的并发模型与 max_concurrency
  • 理解 generator 式 tick 如何与流式输出 interleave。

PregelRunner:Execute 阶段的执行者

libs/langgraph/langgraph/pregel/_runner.py,类 PregelRunner(135 行起)。 第 9 讲主循环里的 runner.tick([未完成任务]) 就是它。它负责:把任务丢给执行器并行跑、 收集每个任务的写、把异常路由到错误处理。

它由 PregelLoop 创建,注入两个回调:

  • submit = loop.submit:提交到后台执行器(见下文)。
  • put_writes = loop.put_writes:任务完成后把写交回 loop 缓冲 + 持久化。

tick:并行调度

tick(176–358 行,async 版 atick 360–572 行)。核心策略:

1. 只跑还没有写的任务

主循环已经过滤过 not t.writes,runner 内部再次确保—— 已有写的任务 = 已完成或 replay 恢复的,不重复执行(第 18 讲 replay 用到)。

2. 单任务快路径

flowchart TD
    A[runner.tick tasks] --> B{任务数 & 配置}
    B -->|单任务,无 timeout/waiter| C[inline 执行<br/>不开线程, 省开销]
    B -->|多任务| D[submit 到 BackgroundExecutor<br/>并行]
    C --> E[run_with_retry 第14讲]
    D --> E
    E --> F[future done callback]
    F --> G[commit: 回收写/异常]

只有一个任务且无超时/无 waiter 时,直接 inline 跑,避免线程开销——常见的线性图就走这条。

3. 多任务并行 + FuturesDict

多任务时,每个任务 submit(run_with_retry, task, ...) 提交到执行器。 FuturesDict(75–132 行)维护 future↔task 映射,并给每个 future 挂 done callback → 触发 commit

4. generator 式 tick:与流式 interleave

注意第 9 讲那段 for _ in runner.tick(...): yield from _output(...)—— tick 是个 generator,每当有任务产出可流式的内容,它就 yield 一次,把控制权还给 Pregel.stream,让后者及时把 chunk 吐给用户。这就是"边执行边流式输出"的实现。

commit:回收任务结果

commit(574–613 行)。任务完成(成功或失败)时被 done callback 调用:

  • 成功:调 put_writes(task.id, task.writes),把缓冲的写交回 loop(第 11 讲 after_tick 再合并)。
  • 失败:根据异常类型写特殊 channel:
    • GraphInterrupt/GraphBubbleUp → 写 INTERRUPT(中断冒泡,第 18 讲)。
    • 普通异常 → 写 ERROR,并可能路由到错误处理节点(第 14 讲)。
    • 无写 → 写 NO_WRITES 标记。

并发模型:BackgroundExecutor

libs/langgraph/langgraph/pregel/_executor.py

同步:BackgroundExecutor(40–119 行)

  • 基于 ThreadPoolExecutor
  • 关键:用 copy_context() 把当前上下文(含 config)拷给子线程, 保证节点里能读到正确的 RunnableConfig / contextvars。

异步:AsyncBackgroundExecutor(122–211 行)

  • 基于 asyncio,用 run_coroutine_threadsafe 调度。
  • 支持 max_concurrency:内部用 Semaphore 限流(gated,214–217 行)—— 这就是你在 config 里设 max_concurrency 限制并行度的落点(防止把下游 API 打爆)。
flowchart LR
    subgraph 同一超步内并行
        T1[task1] --> EX[Executor]
        T2[task2] --> EX
        T3[task3] --> EX
    end
    EX -->|Semaphore max_concurrency| RUN[实际并行 N 个]

durability 与并发重叠

loop.submit 还承担 checkpoint 持久化的后台化。durability 三模式影响并发:

模式checkpoint 时机与执行的关系
"async"(默认)每步后台异步存checkpoint 与下一超步执行重叠,吞吐高
"sync"每步等存完才继续慢但每步落盘,更安全
"exit"仅 run 结束/中断时存最快,但中途崩溃丢失

stream 里 sync 模式会 loop._put_checkpoint_fut.result() 阻塞等待。

使用场景

  • 控制并行度:fan-out 大量子任务时设 config={"max_concurrency": 5}, 避免同时打几百个 LLM/工具请求被限流或 OOM。
  • 性能调优:默认 durability="async" 已让 checkpoint 与执行重叠; 对一致性要求极高用 "sync",对性能极致且能容忍崩溃丢失用 "exit"
  • 排查上下文丢失:节点子线程里读不到 config,往往是没走 copy_context 的路径, 检查是否在非标准线程里手动跑了东西。

动手实验

  1. 做一个 fan-out(一个节点用 Send 扇出 5 个子任务),分别设 max_concurrency=1 和不限, 在子任务里 time.sleep + 打印时间戳,观察并行 vs 串行。
  2. commit 打断点,制造一个节点抛异常,观察它写 ERROR channel 的过程。
  3. 用三种 durability 跑同一个图,对比耗时(可在节点里加延时放大差异)。

阅读作业

  • 精读 _runner.pytick(176–358 行)与 commit(574–613 行)。
  • 精读 _executor.pyBackgroundExecutorAsyncBackgroundExecutor

小结

  • PregelRunner 执行 Execute 阶段:单任务 inline 快路径,多任务并行 submit。
  • commit 回收写(成功)或写 ERROR/INTERRUPT(失败),是错误路由的入口。
  • 并发由 BackgroundExecutor 提供,max_concurrency 限流;durability 控制 checkpoint 与执行的重叠。

下一讲:节点内部的读写机制——ChannelRead / ChannelWrite / local_read。