Checkpoint 概念与数据结构
前 5 个阶段我们把图跑起来了:状态、控制流、Pregel 超步、通道。但那些状态全在内存里——进程一退就没了。从今天起进入阶段 6:怎么把每一步的状态存下来,让图能中断、能恢复、能时间旅行、能跨会话记忆。今天先认识存档的最小单位——Checkpoint。
为什么要存档:内存态的三个致命问题
到 Day 32 为止,图的所有状态都活在通道对象里(内存)。这有三个躲不掉的问题:
- 进程一退就全没:服务重启、崩溃、扩缩容,正在进行的对话/工作流直接蒸发。
- 没法"人在环":想让图停在某一步等人审批,再从那步继续——没有存档就无从"继续"。
- 没法回溯调试:出了错想看"第 3 步时状态长啥样",内存态跑完就覆盖了。
json.dumps(state),你丢掉了版本信息,恢复后 Pregel 引擎就不知道"接下来该唤醒谁"。所以 Checkpoint 必须连同调度所需的元信息一起存。checkpoint 库的分层:一个接口,多种后端
LangGraph 把持久化单独拆成了几个包,核心接口和具体存储分开:
| 包 | 路径 | 装什么 |
|---|---|---|
langgraph-checkpoint | libs/checkpoint/ | 抽象基类 + 内存实现 + 序列化 + Store |
langgraph-checkpoint-sqlite | libs/checkpoint-sqlite/ | SqliteSaver(本地/小项目,Day 36) |
langgraph-checkpoint-postgres | libs/checkpoint-postgres/ | PostgresSaver(生产,Day 37) |
核心包里,一切从 BaseCheckpointSaver 这个抽象基类出发。它定义"存档器长什么样",Sqlite/Postgres 只是不同的"实现"。今天先看它定义的数据结构,明天(D34)看它定义的方法。
Checkpoint 的七个字段(真源码逐行)
存档卡本身是一个 TypedDict,定义在 base/__init__.py:92。这是全阶段最该背下来的结构:
# base/__init__.py:92
class Checkpoint(TypedDict):
"""State snapshot at a given point in time."""
v: int
"""The version of the checkpoint format. Currently `1`."""
id: str
"""The ID of the checkpoint.
This is both unique and monotonically increasing, so can be used for sorting
checkpoints from first to last."""
ts: str
"""The timestamp of the checkpoint in ISO 8601 format."""
channel_values: dict[str, Any]
"""The values of the channels at the time of the checkpoint."""
channel_versions: ChannelVersions
"""The versions of the channels at the time of the checkpoint."""
versions_seen: dict[str, ChannelVersions]
"""Map from node ID to map from channel name to version seen."""
updated_channels: list[str] | None
"""The channels that were updated in this checkpoint."""
v存档卡的格式版本号。注释写 "Currently 1",但 pregel 侧其实已经用到 LATEST_VERSION = 4(D38 讲)——格式演进时靠它做兼容迁移。id存档的唯一 ID,而且单调递增(用 UUID6 生成,D38 详解)。因为递增,直接按字符串大小排序就是"从早到晚",max(keys) 就是最新档。tsISO 8601 时间戳,给人看的,排序不靠它(靠 id)。channel_values核心数据:通道名 → 该通道当时的值。你的 state 里那些字段(messages、counter…)就装在这。channel_versions通道名 → 版本号。每次通道被写,版本号涨一格。用来判断"数据新旧"。versions_seen节点名 → (通道名→已看版本)。记录每个节点"上次看到通道的什么版本",是调度的关键(L04 展开)。updated_channels本次超步更新了哪些通道。可为 None。是个优化缓存,避免恢复时全量对比版本。{"v":4, "id":"1ef4f797-8335-6428-8001-...", "ts":"2024-05-04T06:32:42+00:00", "channel_values":{"messages":[HumanMessage("你好")], "counter":3}, "channel_versions":{"messages":"00000000000000000000000000000002.0.51", "counter":"...0003..."}, "versions_seen":{"chatbot":{"messages":"...0001..."}}, "updated_channels":["counter"]}
最难懂的两个字段:channel_versions 与 versions_seen
新手最容易被这俩绕晕。先看它们的类型定义 base/__init__.py:89:
# base/__init__.py:89
ChannelVersions = dict[str, str | int | float]
# channel_versions: dict[str, ChannelVersions] → 通道名 -> 版本
# versions_seen: dict[str, ChannelVersions] → 节点名 -> {通道名 -> 已看版本}
为什么要两套版本?因为 Pregel 的调度规则是(Day 21 讲过):一个节点该不该在下一步跑,取决于"它订阅的通道有没有出现它没看过的新版本"。
channel_versions = 公众号"最新发到第几期";versions_seen = "这个读者读到第几期"。引擎每步做的事就是:对每个读者,比较"公众号最新期" > "他读过的期"?大于,就把他叫醒去读(执行)。读完把 versions_seen 更新成最新期。channel_values 而不动 channel_versions/versions_seen,那么恢复后引擎会认为"没有新版本"→ 没有任何节点会被唤醒,图看似加载成功却一步都不跑。这正是 D44「时间旅行 / update_state」必须走官方 API(它会正确地 bump 版本)而不能手改存档的原因。版本三件套是"数据"和"调度"之间的咬合齿轮,缺一不可。CheckpointMetadata:存档卡的"标签栏"
光有数据还不够,还要知道"这份档是怎么来的、是第几步"。这就是 CheckpointMetadata,base/__init__.py:38:
# base/__init__.py:38
class CheckpointMetadata(TypedDict, total=False): # total=False:字段都可选,便于扩展
source: Literal["input", "loop", "update", "fork"]
"""- "input": 来自 invoke/stream 的输入
- "loop": pregel 主循环内产生
- "update": 手动 update_state 产生
- "fork": 从另一个 checkpoint 复制而来"""
step: int
"""步号:-1 是第一个 input 档,0 是第一个 loop 档,之后依次 +1。"""
parents: dict[str, str]
"""父 checkpoint 的 ID,按 checkpoint 命名空间映射(子图用)。"""
run_id: str
"""产生此档的那次 run 的 ID。"""
total=False整个 TypedDict 所有字段都可选。注释明说 "to allow for future expansion"——未来加字段不会破坏老代码。source这份档的"来源"。调试时超有用:一眼看出它是用户输入、引擎跑出来的、还是人手动改的。step步号。-1 是最初的输入档,0 是第一次循环,之后递增。时间旅行"回到第 N 步"就靠它。parents父档 ID 映射。有子图时,不同命名空间各有自己的父,所以是个 dict 而不是单个 id。total=False 而 Checkpoint 用普通 TypedDict(全必填)?因为二者稳定性诉求不同。Checkpoint 是引擎恢复必须依赖的核心结构,缺一个字段就没法调度,所以字段全必填、变动极慎重(靠 v 做版本迁移)。CheckpointMetadata 是给人和可观测系统看的旁注,业务/平台经常想往里加自定义键(如 run_id、后来加的 delta 计数器 counters_since_delta_snapshot),做成 total=False 让它"能长肉不破坏兼容"。核心结构求稳、周边结构求活——这是 schema 演进的经典分治。langgraph_* 内部键,写元数据时会被过滤掉,避免内部实现细节污染用户可见的 metadata。D38 会再遇到它。CheckpointTuple:读档时拿到的"整套材料"
存的时候分开存(数据/元数据/writes),读的时候打包成一个 CheckpointTuple 一次给你。base/__init__.py:139:
# base/__init__.py:139
class CheckpointTuple(NamedTuple):
"""A tuple containing a checkpoint and its associated data."""
config: RunnableConfig # 定位这份档的坐标(thread_id/ns/checkpoint_id)
checkpoint: Checkpoint # 存档数据本体(L03 的七字段)
metadata: CheckpointMetadata # 标签栏(L05)
parent_config: RunnableConfig | None = None # 父档的坐标,用来往回走(时间旅行)
pending_writes: list[PendingWrite] | None = None # 尚未落到下一档的待写入
其中 PendingWrite 的定义在 base/__init__.py:31:PendingWrite = tuple[str, str, Any],即 (task_id, channel, value)——"哪个任务、往哪个通道、写了什么"。
config 告诉你"我是谁",parent_config 告诉你"我爸是谁"(顺着它一路往回就是完整历史链 = 时间旅行的骨架),pending_writes 是"上一步任务已经算完、但还没合并进正式通道值"的中间结果(用于崩溃恢复时不丢已完成的任务)。👶 小白:为什么要单独存 pending_writes?跑完直接更新通道不就行了?
👨🏫 老师:因为一个超步里可能有多个任务并行。假设 3 个任务,跑完第 2 个时进程崩了。如果没有 pending_writes,重启后这 3 个任务得全部重跑(可能重复调用 LLM、重复扣费)。有了它,已经跑完的任务结果先记为"待写入",恢复时能识别"这俩已经做完了,只补第 3 个"。这是 D46「replay 与幂等」的基础。
BaseCheckpointSaver 骨架 + 今日小结
最后瞄一眼存档器基类的开头 base/__init__.py:176-217,明天 D34 会逐个方法拆:
# base/__init__.py:176
class BaseCheckpointSaver(Generic[V]):
"""Base class for creating a graph checkpointer."""
serde: SerializerProtocol = JsonPlusSerializer() # 默认序列化器(D39 详讲)
def __init__(self, *, serde: SerializerProtocol | None = None) -> None:
self.serde = maybe_add_typed_methods(serde or self.serde)
Generic[V]基类带泛型参数 V,代表"版本号的类型"。V = TypeVar("V", int, float, str)——版本可以是整数、浮点或字符串。InMemory 用 str,默认基类用 int(L 下方)。serde每个存档器都自带一个序列化器,负责把 Python 对象(含 LangChain Message)转成字节存库、再读回来。默认是 JsonPlusSerializer(D39)。类文档里还给了最重要的一句用法约定:用 checkpointer 时,必须在 config 里传 thread_id——它是存取档案的主键(D38 细讲):
config = {"configurable": {"thread_id": "my-thread"}}
graph.invoke(inputs, config) # 没有 thread_id,checkpointer 无法保存/恢复/时间旅行
🧠 今天你应该能回答
- 为什么内存态不够、必须要 Checkpoint?(重启即丢、没法人在环、没法回溯)
- Checkpoint 的七个字段各是什么?(v/id/ts/channel_values/channel_versions/versions_seen/updated_channels)
- channel_versions 和 versions_seen 有什么区别?(通道最新版 vs 节点看过的版;相减决定谁下一步跑)
- 为什么手改 channel_values 会导致图"读档后不跑"?(没 bump 版本,引擎认为无新版本)
- CheckpointTuple 的 parent_config 有什么用?(串成家谱链,时间旅行沿它回溯)
- pending_writes 解决了什么问题?(超步内并行任务的崩溃恢复/幂等)
✋ 10 分钟动手
# 1. 打开 Checkpoint 的定义,对照今天七字段
sed -n '38,124p' libs/checkpoint/langgraph/checkpoint/base/__init__.py
# 2. 看 CheckpointTuple / PendingWrite
sed -n '31,147p' libs/checkpoint/langgraph/checkpoint/base/__init__.py
# 3. 跑个带 checkpointer 的最小图,打印真实 checkpoint
python - <<'PY'
from langgraph.graph import StateGraph, START
from langgraph.checkpoint.memory import InMemorySaver
from typing import Annotated
import operator
g = StateGraph(dict); g.add_node("inc", lambda s: {"n": s.get("n",0)+1})
g.add_edge(START, "inc"); app = g.compile(checkpointer=InMemorySaver())
cfg = {"configurable": {"thread_id": "t1"}}
app.invoke({"n": 0}, cfg)
print(app.get_state(cfg)) # 观察 values / config / metadata
PY
get_tuple / put / list / put_writes 四大接口方法的签名与契约,以及为什么 put 要额外收一个 new_versions 参数。