发布日期

第 09 讲 · Pregel 类全景与 Runnable API

Pregel 类的 invoke/stream 调用链、节点/通道/触发器索引全景解析

学习目标

  • 理解 Pregel 类的职责边界:它编排,但不亲自做超步细节。
  • 看懂 invoke / stream(及 async 版)的调用链路。
  • 认识三个协作者:PregelLoopPregelRunner_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 行)还会注入一个内部的 TASKSTopic)通道,用来承载 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}
messagesLLM token 流
custom节点内 stream_writer 主动写的内容
debug每步 task/checkpoint 详细事件(调试神器)

invoke 内部用 values 模式取最终结果。

使用场景

  • 理解一次请求的生命周期:从 invoke 入口到 __enter__ 加载、循环、__exit__ 收尾, 这是排查"卡住/慢/结果不对"的总地图。
  • 选对 API:要中间过程用 stream,只要结果用 invoke;要并发处理多输入考虑自己管理多 thread。

动手实验

  1. Pregel.streamwhile loop.tick() 处打断点,用第 1 讲最小图 invoke, 单步观察"tick → runner.tick → after_tick"循环几轮。
  2. 同一个图用 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 = 消费 streamstream 主循环 = tick(Plan) → runner.tick(Execute) → after_tick(Update)。
  • 这条主循环是后面 5 讲的展开纲领。

下一讲:深入 tick() 的 Plan 阶段——谁该跑,怎么决定。