- 发布日期
第 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的路径, 检查是否在非标准线程里手动跑了东西。
动手实验
- 做一个 fan-out(一个节点用
Send扇出 5 个子任务),分别设max_concurrency=1和不限, 在子任务里time.sleep+ 打印时间戳,观察并行 vs 串行。 - 在
commit打断点,制造一个节点抛异常,观察它写ERRORchannel 的过程。 - 用三种 durability 跑同一个图,对比耗时(可在节点里加延时放大差异)。
阅读作业
- 精读
_runner.py的tick(176–358 行)与commit(574–613 行)。 - 精读
_executor.py的BackgroundExecutor与AsyncBackgroundExecutor。
小结
PregelRunner执行 Execute 阶段:单任务 inline 快路径,多任务并行 submit。commit回收写(成功)或写 ERROR/INTERRUPT(失败),是错误路由的入口。- 并发由
BackgroundExecutor提供,max_concurrency限流;durability 控制 checkpoint 与执行的重叠。
下一讲:节点内部的读写机制——ChannelRead / ChannelWrite / local_read。