Day 26 / 共 60 天 · 阶段 4 Pregel 执行引擎(核心深水区·收官)

调试与画图:给 Pregel 引擎装上仪表盘

阶段 4 收官。前 7 天把引擎里外拆了个遍,可跑起来时怎么看清它在干什么?今天两件事:stream_mode="debug" 每步吐出 task / checkpoint 事件(debug.py),让你像看回放一样审视每个超步;get_graph() 把图画成 Mermaid(_draw.py)——而它发现边的方式极妙:拿一张空图、空跑一遍 Pregel 主循环,看每步谁触发了谁。学会观测,也顺带为阶段 6 的时间旅行铺路。

📍 你在 60 天里的位置
①入门 D1-6· ②状态 D7-12· ③控制流 D13-18· ④Pregel D19-26· ⑤通道 D27-32· ⑥持久化 D33-40· ⑦中断 D41-46· ⑧函数式 D47-52· ⑨预制件 D53-58· ⑩收官 D59-60
D19 BSP模型 D20 主循环 D21 任务准备 D22 执行器 D23 读写 D24 IO映射 D25 重试超时 D26 调试画图
L01

给引擎装仪表盘:两种"看清"的手段

🤔 痛点:图跑出来结果不对,可它一共跑了几步、每步激活了谁、状态怎么演变——我完全看不到。 Day 24 的 values/updates 只给你"状态/更新",不告诉你引擎的内部节奏:这一步计划了哪些 task?它们的触发通道是什么?存了哪个 checkpoint?下一步 next 是谁?调试复杂图(尤其有环、有 Send 扇出)时,这些"引擎视角"的信息才是关键。stream_mode="debug" 就是把引擎内部状态逐步暴露出来。

两种手段互补:

  • debug 流debug.py):运行时观测。每步产出两类事件——task(这步要跑谁)和 task_result/checkpoint(这步跑完的结果和存档)。
  • 画图_draw.py):静态观测。graph.get_graph().draw_mermaid() 把节点和边画出来,一眼看清结构。
💡 本质:debug 流复用的就是前面所有 map_* 函数 今天你会发现一个惊喜——debug 模式并没有另起炉灶,它就是把 Day 24 的 map_output_* 换成 map_debug_tasks / map_debug_checkpoint,塞进 Day 20 那个 _emit 机制。观测能力是"复用流式管道 + 换个映射函数"实现的,不是额外挂的一套东西。理解了前面的 _emit 和 map_output,debug 就是水到渠成。
类比:行车记录仪 + 车辆结构图 debug 流像行车记录仪:逐帧记录"这一秒踩了油门、那一秒打了方向",出事了能回放查因;画图像车辆结构图:静态地告诉你发动机连着变速箱、变速箱连着车轮。一个看"动态怎么跑",一个看"静态怎么连"。
L02

map_debug_tasks:这一步要跑谁

debug 的第一类事件由 map_debug_tasksdebug.py:41)产出——在任务执行之前发(Day 20 tick 里 _emit("tasks", ...)):

# debug.py:41
def map_debug_tasks(tasks) -> Iterator[TaskPayload]:
    """Produce "task" events for stream_mode=debug."""
    for task in tasks:
        if task.config is not None and TAG_HIDDEN in task.config.get("tags", []):
            continue                                   # 44 隐藏任务不吐
        payload: TaskPayload = {
            "id": task.id,                             # 48 任务唯一 id
            "name": task.name,                         # 49 节点名
            "input": task.input,                       # 50 喂给它的输入(读通道读出来的)
            "triggers": task.triggers,                 # 51 是被哪些通道触发的
        }
        if task.config is not None:                    # 60 附上用户级 metadata(过滤内部键)
            md = {k: v for k, v in (task.config.get("metadata") or {}).items()
                  if k not in EXCLUDED_METADATA_KEYS}
            ...
            if md:
                payload["metadata"] = md
        yield payload
"triggers": task.triggers调试利器:直接告诉你"这个任务是被哪个通道触发的"。图为什么在这步激活了 agent?看 triggers 就知道——是 messages 通道更新了。省去猜。
"input": task.input任务的实际输入(Day 23 ChannelRead 读出来的那份)。能看到"喂给节点的到底是什么",排查"节点行为怪异"时先看它吃进去的是不是对的。
EXCLUDED_METADATA_KEYS 过滤只保留用户关心的 metadata,丢掉 langgraph 内部键(step/node/checkpoint_ns 等,这些和 task 自身字段重复)。让输出干净。
📝 真实值走一遍 stream(mode="debug") 在某步吐出:{"type":"task", "payload":{"id":"...","name":"agent","input":{"messages":[...]},"triggers":["messages"]}}。读它你立刻知道:"第 N 步,agent 因 messages 通道更新被触发,输入是当前消息列表。"这就是引擎视角的实况。
L03

map_debug_checkpoint:这一步存了什么档

第二类事件由 map_debug_checkpointdebug.py:144)产出——在 after_tick 存档时发,携带这步的完整快照:

# debug.py:144
def map_debug_checkpoint(config, channels, stream_channels, metadata, tasks,
                         pending_writes, parent_config, output_keys) -> Iterator[CheckpointPayload]:
    ...
    yield {
        "config": rm_pregel_keys(patch_checkpoint_map(config, metadata)),   # 177 本步 checkpoint 定位
        "parent_config": ...,                          # 178 上一个 checkpoint(用于时间旅行链)
        "values": read_channels(channels, stream_channels),   # 179 存档时的完整状态
        "metadata": metadata,                          # 180 step/source 等元信息
        "next": [t.name for t in tasks],               # 181 下一步会跑谁
        "tasks": [                                     # 182 每个任务的结果/错误/中断
            {"id": t.id, "name": t.name, "error": t.error, "state": t.state}
            if t.error else
            {"id": t.id, "name": t.name, "result": t.result,
             "interrupts": tuple(asdict(i) for i in t.interrupts), "state": t.state}
            if t.result else
            {"id": t.id, "name": t.name, "interrupts": ..., "state": t.state}
            for t in tasks_w_writes(tasks, pending_writes, task_states, output_keys)
        ],
    }
"values": read_channels(...)存档瞬间的完整状态(复用 Day 24 的 read_channels!)。每个 checkpoint 都带一份全量——这正是时间旅行能回到任意步的数据基础(阶段 6)。
"next": [t.name for t in tasks]下一步待跑的节点名。中断后你 get_state() 看到的 next 就来自这里——它告诉你"恢复后会从哪继续"。
"parent_config"指向上一个 checkpoint。一串 checkpoint 靠 parent 链成一条历史,get_state_history() 就是顺着这条链回溯(Day 44)。
error / result / interrupts 三态任务结果按三种情况分别打包:出错了带 error;正常带 result;被中断带 interrupts。调试时一眼看出"这个任务是成了、崩了、还是停在人在环那了"。
数据结构:debug 流的两类事件长什么样 task 事件(执行前) id: 任务唯一 id name: 节点名 triggers: 被哪些通道触发 ★ input: 喂进去的输入 ★ metadata: 用户级元信息 checkpoint 事件(存档时) values: 存档时完整状态 ★ next: 下一步会跑谁 ★ config / parent_config: 时间旅行链 metadata: step / source tasks: 每任务 result/error/interrupt
图注:task 事件回答"这步跑谁、被谁触发、吃什么";checkpoint 事件回答"现在状态、下一步谁、结果如何"。★ 为 values/updates 拿不到的引擎视角信息。
L04

debug 模式怎么发出:checkpoints/tasks 的重映射

debug 流不是独立通道,而是 _emit(Day 20,_loop.py:1380)把 checkpoints/tasks 事件重新包装成 debug:

# _loop.py:1380
def _emit(self, mode, values, *args, **kwargs) -> None:
    if self.stream is None:
        return
    debug_remap = mode in ("checkpoints", "tasks") and "debug" in self.stream.modes  # 1389
    if mode not in self.stream.modes and not debug_remap:
        return                                         # 1390 既不是订阅的模式、也不需重映射 → 不发
    for v in values(*args, **kwargs):                  # 1392 调 map_debug_* 生成事件
        if mode in self.stream.modes:
            self.stream((self.checkpoint_ns, mode, v))  # 1394 原模式直接发
        if debug_remap:                                # 1396 若订阅了 debug,再包一层发一份
            self.stream((self.checkpoint_ns, "debug", {
                "step": self.step - 1 if mode == "checkpoints" else self.step,  # 1402
                "timestamp": datetime.now(timezone.utc).isoformat(),
                "type": "checkpoint" if mode == "checkpoints" else "task_result" ...
            }))
debug_remap核心开关:当你订阅了 "debug",引擎把内部的 checkpointstasks 事件额外包装成带 step/timestamp/type 的 debug 事件发一份。debug = "checkpoints + tasks 的加壳版"。
"step" / "timestamp" / "type"debug 事件比原始事件多了这三样"元信息包装"——步号、时间戳、类型。有了它们你才能把一串事件按时间/步号排成清晰的回放。
同时发原模式 + debug如果你同时订阅了 tasksdebug,一个任务事件会发两份(各自格式)。互不干扰。
🎯 设计取舍①:debug 为什么是"重映射已有事件"而不是独立埋点? 本可以在引擎各处单独插 debug 埋点。但那意味着两套并行的观测代码,容易信息漂移——debug 显示的和实际发生的不一致(埋点忘了更新)。LangGraph 的做法是让 debug 寄生在既有的 checkpoints/tasks 事件上,只加一层 step/timestamp/type 外壳。这样 debug 看到的就是引擎真正发出的事件,天然同步、零漂移。代价是 debug 事件的结构受限于底层事件(不能任意定制),但换来"所见即真实"的可靠性——对调试工具而言这是对的取舍。
L05

写入聚合:同通道多次写怎么呈现

一个小而精的细节:debug 展示任务写入时,同一通道被写多次要特殊呈现(map_task_result_writesdebug.py:83):

# debug.py:83
def map_task_result_writes(writes) -> dict[str, Any]:
    """单次写记为 {channel: write};多次写记为 {channel: {'$writes': [w1, w2, ...]}}"""
    result: dict[str, Any] = {}
    for channel, value in writes:
        existing = result.get(channel)
        if existing is not None:                       # 93 这个通道之前已经写过
            channel_writes = (existing["$writes"]      # 94 已是聚合形式就取列表
                              if is_multiple_channel_write(existing) else [existing])
            channel_writes.append(value)               # 99 追加本次写入
            result[channel] = {"$writes": channel_writes}   # 100 聚合成 $writes 列表
        else:
            result[channel] = value                    # 102 首次写:直接记值
单次 vs 多次通道只被写一次 → 直接 {通道: 值};被写多次(比如 Send 扇出多个都写同字段)→ {通道: {"$writes": [值1, 值2, ...]}}。用 $writes 包裹区分"这是多条写入"。
为什么要区分?因为 debug 要展示"reducer 合并之前的原始写入"。如果只展示合并后的值,你就看不到"其实有 3 个节点都往这写了"这个事实——而这恰恰是排查"并行写冲突"的关键信息。

而任务跑完后的 task_result 事件(debug.py:106)就用它来打包结果里的写入:

# debug.py:106
def map_debug_task_results(task_tup, stream_keys) -> Iterator[TaskResultPayload]:
    """Produce "task_result" events for stream_mode=debug."""
    task, writes = task_tup
    yield {
        "id": task.id, "name": task.name,
        "error": next((w[1] for w in writes if w[0] == ERROR), None),   # 118 有错取错
        "result": map_task_result_writes(                                # 119 用 L05 那个聚合
            [w for w in writes if w[0] in stream_channels_list or w[0] == RETURN]),
        "interrupts": [asdict(v) for w in writes if w[0] == INTERRUPT    # 122 中断信息
                       for v in (w[1] if isinstance(w[1], Sequence) else [w[1]])],
    }
task_result 事件里 result 字段调的正是 L05 的 map_task_result_writes——所以你在 debug 里看到的任务结果,写入部分就是上面那套"单写记值、多写记 $writes"的规则。
⚠️ 边界/坑:debug 里看到的 $writes 列表 ≠ 最终状态值 新手看到 debug 输出里某字段是 {"$writes": [1, 2, 3]} 会困惑"我的 count 怎么变成列表了?"。其实这是合并前的原始写入清单——3 个并行任务各写了 1、2、3。最终 count 是什么,取决于该通道的 reducer:LastValue 取最后(3),BinaryOp 累加(6)。debug 展示的是"过程",最终值要看 values 模式或 reducer 语义。别把 $writes 当成实际状态。
L06

draw_graph:靠"空跑一遍"发现所有边

换到画图。draw_graph_draw.py:42)最惊艳的地方:它不去读图的静态结构,而是拿空 checkpoint 真的跑一遍 Pregel 循环,看每步谁触发谁,反推出边(_draw.py:88-118):

# _draw.py:88
input_writes = list(map_input(input_channels, {}))    # 89 造一批"空输入"写入(复用 Day 24!)
updated_channels = apply_writes(                       # 90 空跑 apply_writes(复用 Day 21!)
    checkpoint, channels,
    [PregelTaskWrites((), INPUT, input_writes, [])],
    get_next_version, trigger_to_nodes,
)
tasks = prepare_next_tasks(                            # 100 空跑 prepare_next_tasks(复用 Day 21!)
    checkpoint, [], nodes, channels, managed, config,
    step, -1, for_execution=True, ...,
    updated_channels=updated_channels,
)
start_tasks = tasks
for step in range(step, limit):                        # 118 真的循环模拟超步
    if not tasks:
        break                                          # 119 没任务了=图跑完=画完了
    ...记录这步 tasks 是被哪些通道触发的 → 连边...
map_input(input_channels, {})空字典造输入写入——不关心真实数据,只想触发第一批节点。复用的正是 Day 24 的 map_input。
apply_writes / prepare_next_tasks直接调 Day 21 的两个核心函数!画图不是另写一套遍历逻辑,而是真的驱动引擎跑。每步 prepare_next_tasks 算出哪些节点被触发,就等于发现了"上一步的节点 → 这一步的节点"这条边。
for step in range(...): if not tasks: break和 Day 20 主循环一模一样的停止条件!"没任务了就停",此时所有能到达的边都已走过、记录完毕。
控制流:draw_graph 靠"空跑 Pregel 循环"逐步发现边 空输入 map_input 触发入口 prepare_next_tasks 这步谁被触发 记录触发关系→连边 apply_writes 推进一步 循环直到 prepare_next_tasks 返回空 → 所有边发现完毕 全程复用 Day 21/24 的真实引擎函数,不另写遍历
图注:画图=空跑一遍 Pregel 主循环,每步的触发关系就是一条边,直到没任务可跑。
L07

为什么"空跑"而不是"读静态结构"

🎯 设计取舍②:画图为什么不直接读 add_edge 记录的边表,非要模拟运行? 因为 LangGraph 的"边"不全是静态声明的。普通 add_edge 确实能直接读;但条件边(走哪条取决于运行时函数返回)、Send 扇出(运行时动态派发)、Command.goto(节点内决定跳哪)——这些"边"在静态结构里根本不完整。要画出真实的可达关系,唯一可靠的办法是让引擎按真实触发逻辑跑一遍prepare_next_tasks 用的是和真运行一模一样_triggers 版本比较(Day 21),所以它"发现"的边就是真实会发生的边。用静态边表会漏掉动态控制流;用空跑则和实际执行天然一致。这就是为什么画图代码里赫然出现 apply_writes / prepare_next_tasks——它借用了引擎本身作为"事实来源"。

👶 小白:空跑会不会真的执行我的节点函数、产生副作用?

👨‍🏫 老师:不会。空跑只调 prepare_next_tasks(算"谁该被触发")和 apply_writes(推进版本号),并不调 runner 去真正 invoke 你的节点——注意 L06 代码里没有 runner.tick。它关心的只是"触发关系"(谁连谁),不关心节点算出什么值,所以喂的是空输入、也不跑业务逻辑。你的函数一次都不会被调用,放心画图。

💡 本质:整个阶段 4 的函数在这里"闭环复用"了一次 今天画图这段是对阶段 4 最好的复习:它把 Day 24 的 map_input(造输入)、Day 21 的 apply_writes(推进)、Day 21 的 prepare_next_tasks + _triggers(发现触发)、Day 20 的循环骨架(while not tasks break)全串了一遍——只是把"执行器 runner"这一环抽掉了。看懂这段,说明你真的理解了 Pregel 引擎:它的每个零件都能被单独拿来复用。
L08

阶段 4 小结 + 动手 + 明日预告

🧠 今天你应该能回答

  • debug 流给你哪些 values/updates 没有的信息?(task 的 triggers/input、checkpoint 的 next、每任务结果三态)
  • debug 模式是怎么实现的?(_emit 把 checkpoints/tasks 事件加 step/timestamp/type 外壳重映射)
  • 为什么 debug 用重映射而非独立埋点?(寄生既有事件,天然同步、零信息漂移)
  • debug 里 $writes 列表是什么?(合并前的原始多写入,不是最终状态值)
  • draw_graph 怎么发现边?(空跑一遍 Pregel 主循环,每步触发关系即一条边)
  • 为什么画图要模拟运行而非读静态边表?(条件边/Send/Command 是动态的,只有真跑才准)
  • 空跑会执行我的节点吗?(不会,只调 prepare/apply,不调 runner,无副作用)

🏔️ 阶段 4 全景回顾(D19-26)

D19 BSP 模型 → D20 主循环 tick/after_tick → D21 Plan(prepare)/Update(apply_writes) → D22 执行器并行 → D23 节点读写三段式 → D24 图边界 IO → D25 重试超时 → D26 观测。一句话串起来:引擎一波波跑超步(D19/20),每步先算任务(D21)、并行执行(D22),任务本身是"读通道→跑函数→写通道"(D23),图的进出靠 IO 映射(D24),失败有重试兜底(D25),全程可观测可画图(D26)。Pregel 深水区到此结束,你已经拿到了这台引擎的完整心智模型。

✋ 10 分钟动手

# 1. 读 debug 事件映射
sed -n '41,71p' libs/langgraph/langgraph/pregel/debug.py
sed -n '144,206p' libs/langgraph/langgraph/pregel/debug.py

# 2. 读画图"空跑发现边"那段
sed -n '88,120p' libs/langgraph/langgraph/pregel/_draw.py

# 3. 看 debug 流的真实事件 + 画 Mermaid
python -c "
from langgraph.graph import StateGraph, START, END
from typing import TypedDict
class S(TypedDict): n: int
g = StateGraph(S)
g.add_node('a', lambda s:{'n':s['n']+1}); g.add_node('b', lambda s:{'n':s['n']*2})
g.add_edge(START,'a'); g.add_edge('a','b'); g.add_edge('b',END)
app = g.compile()
for ev in app.stream({'n':1}, stream_mode='debug'):
    print(ev['type'], '| step', ev.get('step'), '|', ev['payload'].get('name') or ev['payload'].get('next'))
print(app.get_graph().draw_mermaid())   # 打印 Mermaid 源码
"
明天预告 · Day 27(进入阶段 5):这 8 天我们一直在说"通道 channel"——它有版本、能 update、能 get、能存进 checkpoint,但它内部到底长啥样?下一阶段专攻 Channels。Day 27 先看抽象基类 BaseChannelchannels/base.py)的四个核心方法 update/get/checkpoint/from_checkpoint,看清所有通道的共同契约——然后逐个拆 LastValue、BinaryOp、Topic、屏障通道。
← Day 25 重试与超时 Day 27 · 通道抽象 BaseChannel →