- 发布日期
第 09 讲 · Pregel 类全景与 Runnable API
Pregel 类的 invoke/stream 调用链、节点/通道/触发器索引全景解析
学习目标
- 理解
Pregel类的职责边界:它编排,但不亲自做超步细节。 - 看懂
invoke/stream(及 async 版)的调用链路。 - 认识三个协作者:
PregelLoop、PregelRunner、_algo。
Pregel:执行引擎的"门面"
libs/langgraph/langgraph/pregel/main.py,类 Pregel(449 行起)。 它持有图的全部静态结构,并对外暴露 Runnable 接口。
它持有什么:
nodes:{节点名: PregelNode}(第 8 讲编译产物)。channels:{通道名: BaseChannel}。trigger_to_nodes:通道 → 订阅它的节点 的反向索引(Plan 阶段优化用,第 10 讲)。checkpointer/store/stream_mode等运行配置。
构造时(__init__,757–835 行)还会注入一个内部的 TASKS(Topic)通道,用来承载 Send(第 19 讲)。
它不做什么:超步循环、任务调度、写合并——这些全委托给 _loop.py / _runner.py / _algo.py。 记住这个分工,读源码时就知道该去哪个文件找答案。
对外 API:invoke 是 stream 的薄封装
Pregel 实现了 Runnable 的 invoke/stream/ainvoke/astream。关键事实: invoke 本质是消费 stream 的所有 chunk,取最后的值。
flowchart TD
INV[invoke / ainvoke] --> STR[stream / astream]
STR --> ENTER[SyncPregelLoop.__enter__<br/>加载 checkpoint, 重建 channels]
ENTER --> RUNNER[创建 PregelRunner]
RUNNER --> LOOP{while loop.tick}
LOOP -->|有任务| EXEC[runner.tick 执行节点]
EXEC --> AFTER[loop.after_tick 合并写 + checkpoint]
AFTER --> LOOP
LOOP -->|无任务| OUT[输出最终 state]
源码里 stream 的主循环骨架(同步版)长这样:
2999:libs/langgraph/langgraph/pregel/main.py
while loop.tick():
for task in loop.match_cached_writes():
loop.output_writes(task.id, task.writes, cached=True)
for _ in runner.tick(
[t for t in loop.tasks.values() if not t.writes],
timeout=self.step_timeout,
get_waiter=get_waiter,
schedule_task=loop.accept_push,
):
# emit output
yield from _output(...)
loop.after_tick()
逐行读这段:
while loop.tick():Plan——规划本超步任务,没任务就退出(第 10 讲)。loop.match_cached_writes():命中缓存的任务直接出结果(CachePolicy)。runner.tick([未完成任务]):Execute——并行跑节点(第 12 讲)。 注意只传not t.writes的任务(已有写的说明已完成/replay)。yield from _output(...):把流式 chunk 吐给调用方。loop.after_tick():Update——合并写、bump 版本、存 checkpoint(第 11 讲)。
这就是第 1 讲 BSP 三阶段在代码里的精确落点:tick / runner.tick / after_tick。
调用链路速查
Pregel.invoke(input, config)
└─ Pregel.stream(input, config, stream_mode=...)
├─ with SyncPregelLoop(...) as loop: # __enter__ 加载状态
│ └─ runner = PregelRunner(submit=loop.submit, put_writes=loop.put_writes)
│ └─ while loop.tick(): # Plan
│ runner.tick(tasks) # Execute
│ loop.after_tick() # Update
└─ 返回最后一个 values chunk
async 版 ainvoke → astream → AsyncPregelLoop + runner.atick 完全对称。
stream_mode:一次 run 能产出哪些"视角"
stream 支持多种 stream_mode(第 20 讲细讲),决定它 yield 什么:
| mode | 产出 |
|---|---|
values | 每步结束后的完整 state |
updates | 每个节点本步写了啥 {node: update} |
messages | LLM token 流 |
custom | 节点内 stream_writer 主动写的内容 |
debug | 每步 task/checkpoint 详细事件(调试神器) |
invoke 内部用 values 模式取最终结果。
使用场景
- 理解一次请求的生命周期:从
invoke入口到__enter__加载、循环、__exit__收尾, 这是排查"卡住/慢/结果不对"的总地图。 - 选对 API:要中间过程用
stream,只要结果用invoke;要并发处理多输入考虑自己管理多 thread。
动手实验
- 在
Pregel.stream的while loop.tick()处打断点,用第 1 讲最小图invoke, 单步观察"tick → runner.tick → after_tick"循环几轮。 - 同一个图用
stream_mode=["updates", "values"]跑,打印每个 chunk, 建立"updates=增量、values=全量"的直觉。
阅读作业
- 通读
Pregel类 docstring(449–476 行)与stream主体(2670–3036 行)。 - 找到
Pregel.invoke(3851 行)确认它如何消费 stream 取最终值。
小结
Pregel是门面:持有静态图结构 + 对外 Runnable API,把执行委托给 Loop/Runner/_algo。invoke= 消费stream;stream主循环 =tick(Plan) →runner.tick(Execute) →after_tick(Update)。- 这条主循环是后面 5 讲的展开纲领。
下一讲:深入 tick() 的 Plan 阶段——谁该跑,怎么决定。