- 发布日期
第 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(可排序,时间有序) |
ts | ISO8601 时间戳 |
channel_values | 通道名 → 反序列化后的值 |
channel_versions | 通道名 → 版本号(第 10 讲触发判定用) |
versions_seen | 节点 → {通道 → 已见版本}(第 10/11 讲) |
updated_channels | 本 checkpoint 更新了哪些通道 |
注意 channel_versions 和 versions_seen 都被持久化—— 这就是为什么续跑后调度还能正确进行:版本信息一并存档了。
CheckpointMetadata(38–86 行)
| 字段 | 含义 |
|---|---|
source | "input" / "loop" / "update" / "fork"(这条 checkpoint 怎么来的) |
step | 超步编号(input 首个为 -1) |
parents | checkpoint 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_checkpoint | input | 输入后状态 |
| 每超步结束 | after_tick → _put_checkpoint | loop | 第 11 讲,主力 |
| 执行中 | put_writes | — | 增量 pending writes |
| 退出/中断 | _put_exit_delta_writes + put | — | exit 模式 / 中断收尾 |
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__.py的BaseCheckpointSaver(176–318 行)与数据结构(38–146 行)。 - 浏览
WRITES_IDX_MAP(795 行)与get_next_version(692–711 行)。
小结
Checkpoint= 通道值快照 + 版本信息(channel_versions / versions_seen),后者保证续跑后调度正确。CheckpointTuple含pending_writes;特殊 channel(ERROR/INTERRUPT/RESUME)用负索引。- 写入主力在每超步
after_tick;加载在 run 开始。
下一讲:这些值是怎么序列化进存储的——serde 机制。