发布日期

第 15 讲 · 检查点抽象与数据结构

学习目标

  • 理解 BaseCheckpointSaver 的接口契约(put / put_writes / get_tuple / list)。
  • 看懂 Checkpoint / CheckpointMetadata / CheckpointTuple / PendingWrite 数据结构。
  • 理解 checkpoint 在执行循环中的写入时机。

检查点:把通道快照成可恢复的存档

回顾第 1 讲:因为状态全在通道里,把通道值快照下来就是一个 Checkpoint。 它是续跑、时间旅行、人审中断(第 18 讲)的共同基础。

libs/checkpoint/langgraph/checkpoint/base/__init__.py,核心类 BaseCheckpointSaver(176 行起)。 所有存储后端(InMemory / SQLite / Postgres,第 17 讲)都实现这套接口。

数据结构

Checkpoint:某一时刻的状态快照(92–123 行)

字段含义
v格式版本
id单调递增的 UUID6(可排序,时间有序)
tsISO8601 时间戳
channel_values通道名 → 反序列化后的值
channel_versions通道名 → 版本号(第 10 讲触发判定用)
versions_seen节点 → {通道 → 已见版本}(第 10/11 讲)
updated_channels本 checkpoint 更新了哪些通道

注意 channel_versionsversions_seen 都被持久化—— 这就是为什么续跑后调度还能正确进行:版本信息一并存档了。

CheckpointMetadata(38–86 行)

字段含义
source"input" / "loop" / "update" / "fork"(这条 checkpoint 怎么来的)
step超步编号(input 首个为 -1)
parentscheckpoint namespace → checkpoint_id 的父链(子图用,第 18 讲)

CheckpointTuple:读取时的完整包(139–146 行)

CheckpointTuple(
    config,          # 含 thread_id / checkpoint_ns / checkpoint_id
    checkpoint,
    metadata,
    parent_config,   # 父 checkpoint(None=根)
    pending_writes,  # list[PendingWrite] = (task_id, channel, value)
)

PendingWrite 与特殊 channel 索引

PendingWrite = (task_id, channel, value)。它是"任务执行中产生、但还没合并成正式 checkpoint 的写"——第 13 讲 task.writes 持久化后就成了 pending writes。续跑时(第 18 讲) 靠它恢复"已经成功的任务的写",避免重复执行。

特殊 channel 用负索引区分(795 行):

WRITES_IDX_MAP = {ERROR: -1, SCHEDULED: -2, INTERRUPT: -3, RESUME: -4}
  • ERROR:任务失败信息(第 14 讲)。
  • INTERRUPT:中断点(第 18 讲人审)。
  • RESUME:恢复值(第 18 讲 Command(resume=...))。

这几个负索引会在第 18 讲反复出现,先记住它们是"控制流的特殊写"。

核心方法

方法行号职责
get_tuple(config)239–251主读取接口,返回完整 tuple
list(config, filter, before, limit)253–275列举历史 checkpoint(时间旅行/历史)
put(config, checkpoint, metadata, new_versions)277–298写入 checkpoint,返回含新 id 的 config
put_writes(config, writes, task_id, task_path)300–318写入中间 pending writes(含 interrupt/resume)
get_next_version(current, channel)692–711生成下一个版本号
a* 系列417–580全套异步镜像

get_next_version 决定版本号形态:默认 int+1, 但生产实现(Postgres/SQLite/InMemory)重写成 {032d}.{016f} 字符串以保证全局可排序。

写入时机:与执行循环的交互

时机调用source说明
run 开始加载get_tuple + channels_from_checkpoint恢复通道(第 17 讲)
收到输入_put_checkpointinput输入后状态
每超步结束after_tick_put_checkpointloop第 11 讲,主力
执行中put_writes增量 pending writes
退出/中断_put_exit_delta_writes + putexit 模式 / 中断收尾
flowchart TD
    A[run 开始] --> B[get_tuple 加载最新/指定 checkpoint]
    B --> C[channels_from_checkpoint 重建通道]
    C --> D[每超步 after_tick]
    D --> E[put_writes 增量写]
    D --> F[put: source=loop 存 checkpoint]
    F --> D
    D --> G[退出: 最终 checkpoint]

使用场景

  • 多轮对话记忆:编译时传 checkpointer,调用时带 config={"configurable": {"thread_id": "user-123"}}, 同一 thread 的多次 invoke 自动累积状态——这就是"有记忆的 Agent"。
  • 审计/复盘:用 list() 拉某 thread 的全部 checkpoint,重建任意时刻的状态。
  • 不想要记忆:编译时 checkpointer=False,每次 run 独立。

动手实验

from langgraph.checkpoint.memory import InMemorySaver
app = g.compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
app.invoke({"messages": [("user", "hi")]}, cfg)
app.invoke({"messages": [("user", "again")]}, cfg)  # 记得上一轮

# 看历史
for t in app.get_state_history(cfg):
    print(t.metadata["step"], t.metadata["source"])

观察:第二次 invoke 能看到第一轮的消息;get_state_history 列出每步 checkpoint。

阅读作业

  • 精读 base/__init__.pyBaseCheckpointSaver(176–318 行)与数据结构(38–146 行)。
  • 浏览 WRITES_IDX_MAP(795 行)与 get_next_version(692–711 行)。

小结

  • Checkpoint = 通道值快照 + 版本信息(channel_versions / versions_seen),后者保证续跑后调度正确。
  • CheckpointTuplepending_writes;特殊 channel(ERROR/INTERRUPT/RESUME)用负索引。
  • 写入主力在每超步 after_tick;加载在 run 开始。

下一讲:这些值是怎么序列化进存储的——serde 机制。