LangGraph 源码深潜:Checkpoint、Thread 与 Durable Execution 机制解剖
核心问题: LangGraph 如何通过
checkpointer + thread_id + checkpoint把一次图执行变成可读取、可回放、可恢复、可中断续跑的生产级工作流?源码主线:
StateGraph.compile(checkpointer=...) → Pregel.checkpointer → graph.invoke(input, config) → checkpoint put / put_writes → get_state / get_state_history → StateSnapshot前置文章: 第 6 篇
StateGraph源码解剖、第 7 篇Conditional Edge 与 Router源码解剖、第 8 篇Reducer 与并行状态合并源码解剖、第 9 篇Plan-and-Execute源码解剖、第 10 篇Reflection Loop源码解剖依赖基线:
langgraph==1.2.7、langgraph-checkpoint==4.1.1源码基线:
langchain-ai/langgraphtag1.2.7;PyPI source distribution SHA256dcdf5b441bf8c7c7c154e603b302c9dbfbfd2d11e1b7ae7d93a5aba979dc87bd阅读边界: 本文覆盖 checkpoint、thread、StateSnapshot、super-step、writes、metadata、time travel 与故障恢复的源码机制;不展开外部数据库 checkpointer 的部署细节,不深入 LangGraph Platform 的托管持久化实现。
0. 本篇在源码学习主线中的位置
前面几篇已经解释了 LangGraph 如何把 Agent 范式落到图结构中:
StateGraph ↓Node / Edge / Conditional Edge ↓Reducer / Parallel Merge ↓Plan-and-Execute / Reflection Loop这些能力解决的是“图如何构建、如何路由、如何循环、如何合并状态”。但是生产级 Agent 还需要回答另一个问题:
如果一次图执行不是几百毫秒内完成,而是持续数分钟、数小时,甚至中途需要人工确认、进程重启、失败恢复,该怎么办?
这就是本篇要解决的核心:
普通 graph.invoke ↓带 checkpointer 的 graph.invoke ↓thread-scoped checkpoint ↓StateSnapshot ↓恢复、回放、分叉、HITL本篇只解决:
checkpointer如何在构建期注入CompiledStateGraph / Pregel。thread_id如何把不同会话或任务隔离成不同 checkpoint 序列。checkpoint保存哪些核心信息。get_state()与get_state_history()如何从 checkpointer 中恢复StateSnapshot。- checkpoint 如何支持 time travel、故障恢复和 HITL。
本篇不展开:
- Postgres、SQLite、Redis 等持久化 checkpointer 的具体建表、索引、连接池优化。
- LangGraph Platform / Agent Server 的托管 checkpointer 行为。
- 长期记忆
Store的语义,这与 thread-scoped checkpointer 不同。
1. 本篇问题、学习目标与能力边界
1.1 核心问题
LangGraph 如何在每个执行步骤后保存 graph state,并通过
thread_id组织成可恢复、可回放、可分叉的执行历史?
1.2 学习目标
完成本篇后,读者必须能够:
- 解释
StateGraph.compile(checkpointer=...)在构建期如何把持久化能力注入可执行图。 - 解释
thread_id、checkpoint_id、checkpoint_ns分别解决什么定位问题。 - 画出一次
graph.invoke(input, config)如何在 super-step 边界产生 checkpoint。 - 解释
checkpoint、CheckpointTuple、StateSnapshot的字段差异。 - 说明
get_state()和get_state_history()如何从 checkpointer 读取最新状态和历史状态。 - 说明
put()与put_writes()的分工:一个保存完整 checkpoint,一个保存中间 writes。 - 解释 time travel 为什么不是“回滚原线程”,而是从旧 checkpoint replay 或 fork。
- 判断生产中何时用
InMemorySaver,何时必须换成持久化 checkpointer。
1.3 能力边界
| 能力 | 本篇是否覆盖 | 说明 |
|---|---|---|
| thread-scoped 短期记忆 | 是 | 解释 thread_id 如何定位同一执行线的 checkpoint 序列 |
| durable execution | 是 | 解释 super-step 后保存 checkpoint、故障后恢复执行的机制 |
| time travel / replay / fork | 是 | 解释从历史 checkpoint 继续执行与 update_state 分叉的语义 |
| HITL interrupt resume | 部分覆盖 | 解释 checkpoint 是前提,但 interrupt 细节放到下一篇 |
| long-term memory store | 否 | Store 面向跨 thread 应用数据,不属于本篇核心 |
| 数据库 checkpointer 部署 | 否 | 本篇只讲接口契约与源码主链,不讲运维配置 |
2. 核心概念与最小心智模型
2.1 一句话定义
Checkpoint 是 LangGraph 在一次图执行的 super-step 边界保存的 graph state 快照,负责让 thread 内状态可读取、可恢复、可回放,不负责跨线程长期知识管理。
2.2 最小心智模型
graph.invoke(input, config={"configurable": {"thread_id": "travel-thread-001"}}) ↓Pregel runtime 按 super-step 执行节点 ↓每个 super-step 产生 channel writes 和 state 更新 ↓checkpointer.put / put_writes 保存 checkpoint 与中间写入 ↓graph.get_state(config) 读取最新 StateSnapshot ↓graph.get_state_history(config) 读取历史 StateSnapshot 序列2.3 核心术语
| 术语 | 源码对象 | 语义 | 不要误解为 |
|---|---|---|---|
| Checkpointer | BaseCheckpointSaver / InMemorySaver | checkpoint 存取接口与具体存储实现 | 普通 dict 缓存 |
| Thread | thread_id | 一条会话或任务执行线的 checkpoint 命名空间 | 操作系统线程 |
| Checkpoint | Checkpoint TypedDict | 某一时刻 channel values、versions、versions_seen 等底层快照 | 用户可直接读写的业务 state |
| CheckpointTuple | CheckpointTuple | checkpoint + config + metadata + parent_config + pending_writes | 只有 checkpoint 本体 |
| StateSnapshot | StateSnapshot | get_state() 面向用户返回的可读状态视图 | 原始持久化记录 |
| Super-step | Pregel step | 一轮图调度边界,同一轮可有多个节点并行 | 单个节点执行 |
| Writes | put_writes() 保存的数据 | 节点任务产生的中间写入、错误、interrupt、resume 等 | 最终完整 state |
| Metadata | CheckpointMetadata | step、source、writes、parents 等运行元数据 | 用户业务字段 |
2.4 与相邻抽象的边界
| 对象 | 负责什么 | 不负责什么 | 与本篇对象的关系 |
|---|---|---|---|
StateGraph | 声明 state schema、nodes、edges | 不直接执行,不直接持久化 | compile(checkpointer=...) 把 checkpointer 传给可执行图 |
Pregel | 运行图、调度节点、应用 writes | 不决定业务语义 | 运行时在 step 边界与 checkpointer 协作 |
BaseCheckpointSaver | 定义 checkpoint 存取接口 | 不知道业务节点语义 | Pregel 通过它保存和读取执行快照 |
Store | 跨 thread 的长期键值存储 | 不保存每一步 graph state | 与 checkpointer 可同时使用,但语义不同 |
interrupt() | 暂停图并等待外部 resume | 不负责底层持久化 | 依赖 checkpointer 保存暂停点 |
3. 完整执行链路
3.1 高层链路
最小代码用于观察 checkpointer + thread_id 的主链:
from typing_extensions import TypedDictfrom langgraph.graph import StateGraph, START, ENDfrom langgraph.checkpoint.memory import InMemorySaver
class TravelState(TypedDict): user_request: str plan: str | None
def plan_node(state: TravelState): return {"plan": "生成旅行计划草稿"}
builder = StateGraph(TravelState)builder.add_node("plan_node", plan_node)builder.add_edge(START, "plan_node")builder.add_edge("plan_node", END)
checkpointer = InMemorySaver()graph = builder.compile(checkpointer=checkpointer)
config = {"configurable": {"thread_id": "travel-thread-001"}}
graph.invoke({"user_request": "东京 5 天旅行", "plan": None}, config=config)
latest_state = graph.get_state(config)history = list(graph.get_state_history(config))对应高层链路:
用户输入 TravelState ↓graph.invoke(input, config) ↓configurable.thread_id 定位执行 thread ↓Pregel runtime 创建 input checkpoint / loop checkpoint ↓plan_node 执行并返回 partial state update ↓writes 应用到 channel ↓checkpointer 保存 checkpoint ↓graph.get_state(config) 读取最新 StateSnapshot ↓graph.get_state_history(config) 读取 StateSnapshot 历史3.2 对象流转
| 阶段 | 输入类型 | 核心函数 | 输出类型 | 状态变化 |
|---|---|---|---|---|
| 编译 | StateGraph + checkpointer | StateGraph.compile() | CompiledStateGraph | checkpointer 进入 Pregel runtime |
| 调用 | input + RunnableConfig | CompiledStateGraph.invoke() / Pregel.stream() | dict 或 stream chunks | 初始化或读取 thread checkpoint |
| 节点执行 | StateSnapshot.values | plan_node() | Partial[TravelState] | 产生 writes |
| 持久化 | Checkpoint + metadata + versions | BaseCheckpointSaver.put() | updated config | 保存完整 checkpoint |
| 中间写入 | writes + task_id | BaseCheckpointSaver.put_writes() | None | 保存 task 级 pending writes |
| 最新状态读取 | config(thread_id) | Pregel.get_state() | StateSnapshot | 无写入,只读取 |
| 历史状态读取 | config(thread_id) | Pregel.get_state_history() | Iterator[StateSnapshot] | 无写入,只读取历史 |
3.3 时序链路
Caller │ │ graph.invoke(input, config={thread_id}) ▼CompiledStateGraph / Pregel │ │ 读取或创建 thread checkpoint ▼Pregel Super-step Runtime │ │ 执行 plan_node,收集 writes ▼Checkpoint Saver │ │ put_writes / put checkpoint ▼Pregel │ │ 返回最终输出 ▼Caller │ │ graph.get_state(config) ▼StateSnapshot3.4 正常结束条件
一次带 checkpoint 的图执行正常结束,至少满足三个条件:
1. 图调度到 END,没有剩余 next task。2. 最后一次状态更新已经应用到 channels。3. 最新 checkpoint 已经通过 checkpointer 保存,可由 get_state(config) 读取。最终用户可读状态保存在:
StateSnapshot.values历史执行轨迹可通过:
graph.get_state_history(config)读取为多个 StateSnapshot。
4. 源码地图、关键文件与阅读顺序
4.1 核心目录
langgraph/├── graph/│ └── state.py├── pregel/│ ├── main.py│ ├── loop.py│ ├── algo.py│ └── _checkpoint.py├── checkpoint/│ ├── base/│ │ └── __init__.py│ └── memory/│ └── __init__.py└── types.py4.2 关键文件
| 优先级 | 文件 | 核心对象 | 阅读目的 |
|---|---|---|---|
| 1 | langgraph/graph/state.py | StateGraph.compile()、CompiledStateGraph | 看 checkpointer 如何从 builder 传入 compiled graph |
| 2 | langgraph/pregel/main.py | Pregel、get_state()、get_state_history() | 看运行时如何读取 checkpoint 并构造 StateSnapshot |
| 3 | langgraph/checkpoint/base/__init__.py | Checkpoint、CheckpointTuple、BaseCheckpointSaver | 看 checkpoint 数据结构和 saver 接口契约 |
| 4 | langgraph/checkpoint/memory/__init__.py | InMemorySaver | 看最小 checkpointer 如何按 thread/ns/checkpoint_id 存储 |
| 5 | langgraph/pregel/algo.py | apply_writes()、prepare_next_tasks() | 看 writes 如何应用到 channel 并决定 next tasks |
| 6 | langgraph/types.py | StateSnapshot、PregelTask、Command | 看用户可见状态快照结构 |
4.3 推荐阅读顺序
1. `StateGraph.compile(checkpointer=...)`2. `BaseCheckpointSaver` / `Checkpoint` / `CheckpointTuple`3. `InMemorySaver.get_tuple()` / `put()` / `put_writes()`4. `Pregel.get_state()`5. `Pregel._prepare_state_snapshot()`6. `Pregel.get_state_history()`7. time travel 文档中的 replay / fork / update_state4.4 不建议的阅读顺序
不建议直接从 pregel/loop.py 或 pregel/algo.py 开始。那里是调度核心,包含大量并发、重试、stream、interrupt、managed value、subgraph 逻辑,如果没有先理解 checkpoint 数据结构,很容易把三个层次混在一起:
业务 state ≠ checkpoint ≠ StateSnapshot ≠ pending writes更合理的顺序是:先看 checkpointer 契约,再看 get_state 如何把底层 checkpoint 转成用户可读 StateSnapshot,最后再反查 Pregel 主循环何时写 checkpoint。
5. 对象模型、继承关系与协议边界
5.1 核心对象关系
StateGraph ↓ compile(checkpointer=...)CompiledStateGraph ↓ inherits / wraps Pregel behaviorPregel ↓ usesBaseCheckpointSaver ↓ implemented byInMemorySaver / PostgresSaver / SqliteSaver / custom saver另一个数据视角:
Checkpoint ↓ wrapped with config / metadata / parent_config / pending_writesCheckpointTuple ↓ converted by Pregel._prepare_state_snapshotStateSnapshot5.2 对象职责
| 对象 | 生命周期 | 输入 | 输出 | 核心职责 |
|---|---|---|---|---|
StateGraph | 构建期 | state schema、nodes、edges | CompiledStateGraph | 注册图结构,不执行 |
CompiledStateGraph | 编译后 | invoke/stream 输入 | graph output / chunks | 对外暴露 Runnable 接口 |
Pregel | 运行时 | config、channels、nodes、checkpointer | 状态更新与结果 | 执行图、调度节点、保存/读取 checkpoint |
BaseCheckpointSaver | 运行时依赖 | config、checkpoint、metadata | checkpoint tuple / updated config | 定义 checkpoint 存储契约 |
InMemorySaver | 开发测试期 | thread_id/ns/checkpoint_id | 内存中的 checkpoint tuple | 基于内存保存 checkpoint,不跨进程持久 |
StateSnapshot | 读取期 | checkpoint tuple | 用户可读状态 | get_state / get_state_history 的返回结构 |
5.3 协议边界
CheckpointSaver 协议 负责:保存和读取 checkpoint、pending writes、metadata。 不负责:决定节点怎么执行、state 如何 reducer 合并、业务字段代表什么。
Pregel 协议 负责:调度节点、应用 writes、维护 channel 版本、调用 checkpointer。 不负责:持久化底层介质的事务实现。
StateSnapshot 协议 负责:向用户展示 values、next、metadata、tasks、interrupts。 不负责:作为底层存储格式直接写入数据库。5.4 稳定接口与内部实现
| 类型 | 对象 | 文章中的使用原则 |
|---|---|---|
| 公共 API | builder.compile(checkpointer=...) | 工程代码可以直接使用 |
| 公共 API | graph.get_state(config) | 工程代码可以用于调试、恢复、HITL 页面展示 |
| 公共 API | graph.get_state_history(config) | 工程代码可以用于回放、审计、time travel |
| 扩展接口 | BaseCheckpointSaver | 自定义 checkpointer 应实现该协议 |
| 内部实现 | Pregel._prepare_state_snapshot() | 用于理解,不建议业务代码依赖 |
| 内部实现 | checkpoint["channel_versions"] / versions_seen | 用于解释调度,不建议直接修改 |
6. 源码阅读策略与证据标准
6.1 本篇阅读策略
先找公开入口 compile(checkpointer) ↓确认 checkpointer 进入 Pregel ↓阅读 BaseCheckpointSaver 的数据结构与方法契约 ↓阅读 InMemorySaver 的 thread/ns/checkpoint_id 存储布局 ↓阅读 get_state / get_state_history 如何读取 checkpoint tuple ↓阅读 _prepare_state_snapshot 如何构造 StateSnapshot ↓反查 super-step、writes、time travel 与故障恢复6.2 证据等级
| 标记 | 含义 | 写作要求 |
|---|---|---|
| 源码事实 | 可以由 langgraph==1.2.7 源码直接证明 | 附 tag 1.2.7 源码链接 |
| 官方契约 | 官方文档或 API Reference 明确承诺 | 附官方文档链接 |
| 简化伪代码 | 对真实控制流的压缩表达 | 明确标注“不是源码逐字复制” |
| 作者推断 | 根据调用链得出的设计理解 | 明确使用“从调用关系可以推断” |
| 工程建议 | 面向项目实践的建议 | 说明适用条件 |
6.3 本篇证据清单
| 结论 | 证据类型 | 文件或文档 | 定位 |
|---|---|---|---|
| LangGraph 使用 checkpointer 保存 thread graph state | 官方契约 | Persistence docs | Checkpointers / Stores 区分 |
thread_id 是 checkpoint 的主定位键 | 源码事实 / API Reference | BaseCheckpointSaver | class docstring |
Checkpoint 包含 channel values、versions、versions_seen、updated_channels | 源码事实 | checkpoint/base/__init__.py | Checkpoint TypedDict |
get_state() 调用 checkpointer.get_tuple(config) | 源码事实 | pregel/main.py | Pregel.get_state() |
get_state_history() 调用 checkpointer.list(...) | 源码事实 | pregel/main.py | Pregel.get_state_history() |
InMemorySaver 只适合调试/测试 | 源码事实 | checkpoint/memory/__init__.py | InMemorySaver docstring |
| time travel replay/fork 基于 checkpoint config | 官方契约 | Time travel docs | replay / fork |
7. 构建期源码解剖
本章回答:
用户传入的 checkpointer 如何被接入图编译产物,并在运行时成为 checkpoint 读写能力?
7.1 构建期职责
| 输入 | 归一化动作 | 构建结果 |
|---|---|---|
checkpointer=InMemorySaver() | 校验为可用 checkpointer | 保存到 CompiledStateGraph / Pregel.checkpointer |
checkpointer=None | 不启用持久化 | 图可执行,但 get_state 无 checkpointer 时会失败 |
checkpointer=True | 子图或继承场景的特殊语义 | 运行时从 config 或父图注入 checkpointer |
store=... | 作为长期存储依赖 | 与 checkpointer 一起进入 runtime,但语义不同 |
7.2 构建期总链路
StateGraph(...) ↓add_node / add_edge ↓builder.compile(checkpointer=InMemorySaver()) ↓validate graph structure ↓construct CompiledStateGraph ↓attach nodes / edges / branches ↓store checkpointer on Pregel runtime ↓return executable graph7.3 StateGraph.compile(checkpointer=...) 源码解剖
职责与源码签名
以下签名来自当前正式版源码的语义化压缩,参数保留本文相关部分:
def compile( self, checkpointer: Checkpointer = None, *, store: BaseStore | None = None, interrupt_before: All | list[str] | None = None, interrupt_after: All | list[str] | None = None, debug: bool = False, name: str | None = None, cache: BaseCache | None = None,) -> CompiledStateGraph: ...调用方与被调用方
用户代码 ↓builder.compile(checkpointer=checkpointer) ↓StateGraph.validate() ↓CompiledStateGraph(..., checkpointer=checkpointer, store=store, ...) ↓compiled.attach_node / attach_edge / attach_branch ↓compiled.validate()输入与输出
| 项目 | 类型 | 语义 |
|---|---|---|
| 输入 | `BaseCheckpointSaver | None |
| 输出 | CompiledStateGraph | 可执行图 |
| 副作用 | builder 标记 compiled | 后续继续改 builder 会有警告或不影响已编译图 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def compile(self, checkpointer=None, store=None, **options): # 1. 校验图结构 self.validate( interrupt_before=options.get("interrupt_before"), interrupt_after=options.get("interrupt_after"), )
# 2. 标记 builder 已经编译 self.compiled = True
# 3. 构造可执行图 compiled = CompiledStateGraph( builder=self, nodes={}, channels={**self.channels, START: EphemeralValue(...)} , input_channels=START, output_channels=self.output_channels, stream_channels=self.stream_channels, checkpointer=checkpointer, store=store, interrupt_before=..., interrupt_after=..., debug=..., name=..., cache=..., )
# 4. 把 builder 中的 node/edge/branch 挂到 compiled graph compiled.attach_node(START, None) for name, node_spec in self.nodes.items(): compiled.attach_node(name, node_spec)
for start, end in self.edges: compiled.attach_edge(start, end)
for start, branches in self.branches.items(): for branch_name, branch_spec in branches.items(): compiled.attach_branch(start, branch_name, branch_spec)
# 5. 最终校验,并返回可执行图 return compiled.validate()逐段解释
第一段 validate() 确保图结构在进入运行时之前是自洽的:节点存在、边合法、入口出口明确、interrupt 节点合法。checkpoint 不负责修复图结构错误,因此必须在编译期先完成结构校验。
第二段 self.compiled = True 表明 StateGraph 是 builder,而不是执行器。编译后产物才是真正带运行时依赖的对象,后续修改 builder 不应该被理解为修改已经返回的 compiled graph。
第三段构造 CompiledStateGraph 时,checkpointer 和 store 进入 Pregel runtime。也就是说,checkpoint 不是每个节点自己保存,而是图运行时统一保存。
第四段把节点、普通边、条件边全部挂载到可执行结构。这里的关键是:checkpoint 能保存的是运行时 channel 状态,前提是节点和边已经被编译成 Pregel 可调度结构。
第五段返回的 CompiledStateGraph 实现了 invoke / stream / get_state / get_state_history 等接口。
关键分支
| 条件 | 分支动作 | 构建结果 |
|---|---|---|
checkpointer=None | 不设置持久化器 | 可执行,但不能读取 checkpoint state |
checkpointer=InMemorySaver() | 保存到内存 | 适合调试、测试,不跨进程 |
checkpointer=True | 运行时从父图或 config 获取 | 常见于子图继承场景 |
store 同时传入 | 长期存储也进入 runtime | 节点可访问 store,但不替代 checkpoint |
设计原因
checkpointer 必须在构建期接入,而不是节点内手动保存,原因是:
1. checkpoint 需要覆盖所有节点和边,不应分散在业务节点。2. checkpoint 必须与 Pregel super-step 边界对齐。3. checkpoint 需要保存 versions_seen、channel_versions 等调度信息,业务节点无法可靠维护。4. checkpoint 需要支持 get_state、history、time travel、interrupt resume 等统一能力。源码证据
langgraph/graph/state.py::StateGraph.compilehttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py
7.4 BaseCheckpointSaver 契约源码解剖
职责与所处阶段
BaseCheckpointSaver 是持久化层的抽象基类。它定义了 Pregel runtime 能调用哪些方法来保存、读取、列出、删除 checkpoint。
真实源码签名
以下签名来自当前正式版源码的关键方法:
class BaseCheckpointSaver(Generic[V]): def get(self, config: RunnableConfig) -> Checkpoint | None: ... def get_tuple(self, config: RunnableConfig) -> CheckpointTuple | None: ... def list( self, config: RunnableConfig | None, *, filter: dict[str, Any] | None = None, before: RunnableConfig | None = None, limit: int | None = None, ) -> Iterator[CheckpointTuple]: ... def put( self, config: RunnableConfig, checkpoint: Checkpoint, metadata: CheckpointMetadata, new_versions: ChannelVersions, ) -> RunnableConfig: ... def put_writes( self, config: RunnableConfig, writes: Sequence[tuple[str, Any]], task_id: str, task_path: str = "", ) -> None: ...调用方与被调用方
Pregel runtime ↓BaseCheckpointSaver.get_tuple / list / put / put_writes ↓Concrete saver: InMemorySaver / PostgresSaver / SqliteSaver / custom saver输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | RunnableConfig | 包含 configurable.thread_id、可选 checkpoint_id、checkpoint_ns |
| 输出 | `CheckpointTuple | Iterator[CheckpointTuple] |
| 状态变化 | saver 内部存储 | 存储 checkpoint、metadata、writes、blobs |
| 副作用 | 有 | 写入内存或外部数据库 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
class BaseCheckpointSaver: def get(self, config): # 1. get 是 get_tuple 的简化包装 checkpoint_tuple = self.get_tuple(config) if checkpoint_tuple is None: return None return checkpoint_tuple.checkpoint
def get_tuple(self, config): # 2. 抽象方法:具体 saver 必须实现 raise NotImplementedError
def list(self, config, *, filter=None, before=None, limit=None): # 3. 抽象方法:按 thread / namespace / filter 列出 checkpoint raise NotImplementedError
def put(self, config, checkpoint, metadata, new_versions): # 4. 抽象方法:保存完整 checkpoint raise NotImplementedError
def put_writes(self, config, writes, task_id, task_path=""): # 5. 抽象方法:保存 task 级中间写入 raise NotImplementedError逐段解释
get() 只是便利方法,真正读取通常走 get_tuple(),因为恢复状态不仅需要 checkpoint 本身,还需要 metadata、parent_config 和 pending_writes。
list() 是 get_state_history() 的基础。没有 list(),就只能读取最新 checkpoint,无法支持 history、time travel、审计和回放。
put() 保存完整 checkpoint。它处理的是“某个 super-step 后的完整状态快照”。
put_writes() 保存中间写入。它处理的是节点任务级别的 writes,例如 pending writes、interrupt、error、resume 等。它让系统能在未完成完整 step 时保留关键中间信息。
正常路径
Pregel 执行 step ↓节点产生 writes ↓put_writes 保存 task writes ↓writes 应用到 channels ↓create checkpoint ↓put 保存完整 checkpoint关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
自定义 saver 未实现 get_tuple | 抛 NotImplementedError | 不能恢复 state |
自定义 saver 未实现 list | 抛 NotImplementedError | 不能读取 history |
未传 thread_id | 配置定位失败 | 无法可靠保存/恢复 |
| 存储不可用 | 具体 saver 抛异常 | graph 执行失败或恢复失败 |
设计原因与工程影响
BaseCheckpointSaver 把“图运行时”和“存储介质”解耦。Pregel 不关心 checkpoint 存在内存、Postgres、SQLite 还是其他系统;它只依赖统一接口。
工程影响是:
开发调试:InMemorySaver本地持久:SqliteSaver生产服务:PostgresSaver / 托管 checkpointer / 自定义高可用 saver源码证据
langgraph/checkpoint/base/__init__.py::BaseCheckpointSaverhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/base/__init__.py
7.5 Checkpoint 与 CheckpointTuple 数据结构源码解剖
职责与所处阶段
Checkpoint 是底层快照;CheckpointTuple 是 saver 读取时返回的完整记录,包括快照、配置、元数据、父 checkpoint 和 pending writes。
真实源码签名
class Checkpoint(TypedDict): v: int id: str ts: str channel_values: dict[str, Any] channel_versions: ChannelVersions versions_seen: dict[str, ChannelVersions] updated_channels: list[str] | None
class CheckpointTuple(NamedTuple): config: RunnableConfig checkpoint: Checkpoint metadata: CheckpointMetadata parent_config: RunnableConfig | None = None pending_writes: list[PendingWrite] | None = None调用方与被调用方
Concrete saver ↓returns CheckpointTuple ↓Pregel._prepare_state_snapshot ↓StateSnapshot输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | channel 状态、versions、metadata | Pregel step 产生的底层运行信息 |
| 输出 | Checkpoint / CheckpointTuple | 持久化记录和读取记录 |
| 状态变化 | 无 | 数据结构本身不执行写入 |
| 副作用 | 无 | 副作用由 saver.put 执行 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def create_checkpoint(previous_checkpoint, channels, step): # 1. 生成 checkpoint id 与 timestamp checkpoint_id = generate_monotonic_id(step) timestamp = now_iso8601()
# 2. 从 channel 中提取可 checkpoint 的值 values = {} for channel_name, channel in channels.items(): if channel_name not in previous_checkpoint["channel_versions"]: continue try: values[channel_name] = channel.checkpoint() except EmptyChannelError: pass
# 3. 保存 channel version 与 versions_seen return { "v": LATEST_VERSION, "id": checkpoint_id, "ts": timestamp, "channel_values": values, "channel_versions": previous_checkpoint["channel_versions"], "versions_seen": previous_checkpoint["versions_seen"], "pending_sends": previous_checkpoint.get("pending_sends", []), "updated_channels": None, }逐段解释
第一段生成 checkpoint ID。ID 不是普通随机日志 ID,而是可排序、单调递增的定位信息,用于历史顺序和 checkpoint 定位。
第二段从 channels 中提取快照值。LangGraph 的 state key 在运行时是 channel,所以 checkpoint 保存的是 channel-level state,而不是直接保存用户的 TypedDict 原对象。
第三段保存 channel_versions 与 versions_seen。这两个字段是调度恢复的关键:框架需要知道每个节点看过哪些 channel 版本,才能判断恢复后哪些节点应该继续执行。
正常路径
channels ↓channel.checkpoint() ↓Checkpoint.channel_values ↓CheckpointTuple(checkpoint, metadata, parent_config, pending_writes)关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 某 channel 尚未有值 | 捕获 EmptyChannelError | 不写入该 channel value |
| checkpoint 有 parent | 写入 parent_config | 支持历史链和 fork |
| 有 pending writes | 写入 pending_writes | 恢复时可应用未完成写入 |
设计原因与工程影响
Checkpoint 保存的不只是业务 state,还保存调度状态。这就是 checkpoint 能恢复执行,而普通业务日志只能回看历史的根本差异。
如果你只保存:
{"plan": "生成旅行计划草稿"}你只能知道结果是什么,不能知道:
下一步该执行哪个节点?哪些 channel 版本已经被哪些节点看过?是否有 pending writes?是否处于 interrupt?是否可从某个 checkpoint fork?源码证据
langgraph/checkpoint/base/__init__.py::Checkpointlanggraph/checkpoint/base/__init__.py::CheckpointTuplelanggraph/checkpoint/base/__init__.py::create_checkpointhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/base/__init__.py
7.6 InMemorySaver 存储布局源码解剖
职责与所处阶段
InMemorySaver 是最小可用 checkpointer,适合调试和测试。它把 checkpoint 存在当前进程内存里,进程重启后丢失。
真实源码签名
class InMemorySaver( BaseCheckpointSaver[str], AbstractContextManager, AbstractAsyncContextManager,): storage: defaultdict[str, dict[str, dict[str, tuple[...]]]] writes: defaultdict[tuple[str, str, str], dict[tuple[str, int], tuple[...]]] blobs: dict[tuple[str, str, str, str | int | float], tuple[str, bytes]]
def __init__( self, *, serde: SerializerProtocol | None = None, factory: type[defaultdict] = defaultdict, ) -> None: ...调用方与被调用方
用户代码 ↓InMemorySaver() ↓builder.compile(checkpointer=memory) ↓Pregel runtime ↓InMemorySaver.get_tuple / put / put_writes / list输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | serde、factory | 序列化器与存储容器工厂 |
| 输出 | InMemorySaver | 内存 checkpointer 实例 |
| 状态变化 | 初始化 storage/writes/blobs | 准备保存 checkpoint |
| 副作用 | 仅内存 | 不写磁盘、不写数据库 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
class InMemorySaver(BaseCheckpointSaver): def __init__(self, serde=None, factory=defaultdict): # 1. 初始化序列化器 super().__init__(serde=serde)
# 2. thread_id -> checkpoint_ns -> checkpoint_id -> checkpoint entry self.storage = factory(lambda: defaultdict(dict))
# 3. (thread_id, checkpoint_ns, checkpoint_id) -> task writes self.writes = factory(dict)
# 4. (thread_id, checkpoint_ns, channel, version) -> serialized blob self.blobs = factory()
# 5. 支持 context manager 管理 self.stack = ExitStack()逐段解释
storage 是 checkpoint 主记录表。第一层是 thread_id,说明 checkpoint 天然按 thread 隔离。
writes 保存某个 checkpoint 下的 task writes。它不是完整 checkpoint,而是节点任务产生的中间写入。
blobs 保存 channel-level value。这样可以避免所有 checkpoint 都重复保存全部大对象,并支持按 channel/version 取值。
正常路径
thread_id = "travel-thread-001" ↓checkpoint_ns = "" ↓checkpoint_id = "..." ↓storage[thread_id][checkpoint_ns][checkpoint_id] ↓writes[(thread_id, checkpoint_ns, checkpoint_id)] ↓blobs[(thread_id, checkpoint_ns, channel, version)]关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
同一 thread_id 多次 invoke | 读取同一 thread 下最新 checkpoint | 实现短期会话连续性 |
不同 thread_id | 访问不同 storage 分支 | 状态隔离 |
| 进程重启 | 内存清空 | checkpoint 丢失 |
| 生产使用 | 不推荐 | 应换 Postgres/SQLite/托管 saver |
设计原因与工程影响
InMemorySaver 把 checkpoint 接口行为完整跑通,但不提供真正 durable storage。它适合:
单元测试本地 debug理解 checkpoint 数据结构演示 HITL / time travel不适合:
生产服务多进程部署重启恢复审计留存高可用任务恢复源码证据
langgraph/checkpoint/memory/__init__.py::InMemorySaverhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/memory/__init__.py
7.7 构建期产物
| 产物 | 保存的信息 | 运行时用途 |
|---|---|---|
CompiledStateGraph.checkpointer | 用户传入的 saver 或继承标记 | invoke/get_state/history 读写 checkpoint |
CompiledStateGraph.store | 长期 store | 节点访问跨 thread 数据 |
channels | state key 对应 channel | checkpoint 提取 channel values |
input_channels/output_channels/stream_channels | 输入输出和流式字段 | 构造 StateSnapshot.values 和输出 |
nodes/edges/branches | 编译后的调度结构 | super-step 执行和 next task 计算 |
8. 运行时主链源码解剖
本章回答:
构建完成后,一次
invoke / get_state / get_state_history如何进入 checkpoint 读写逻辑并产生状态快照?
8.1 运行时入口
| 调用方式 | 公开入口 | 核心内部入口 | 返回类型 |
|---|---|---|---|
| 同步调用 | graph.invoke(input, config) | Pregel.stream() / loop runtime | OutputT |
| 异步调用 | graph.ainvoke(input, config) | Pregel.astream() / async loop | OutputT |
| 流式调用 | graph.stream(input, config) | Pregel step loop | Iterator chunks |
| 批量调用 | graph.batch(inputs, configs) | Runnable batch | list[OutputT] |
| 最新状态 | graph.get_state(config) | checkpointer.get_tuple() + _prepare_state_snapshot() | StateSnapshot |
| 历史状态 | graph.get_state_history(config) | checkpointer.list() + _prepare_state_snapshot() | Iterator[StateSnapshot] |
8.2 运行时总链路
graph.invoke(input, config) ↓ensure_config(config) ↓读取 configurable.thread_id ↓读取最新 checkpoint 或创建 input checkpoint ↓按 Pregel super-step 执行节点 ↓节点产生 writes ↓apply_writes 更新 channels ↓create_checkpoint ↓checkpointer.put / put_writes ↓返回最终 output8.3 Config 与 thread_id 归一化源码解剖
职责与所处阶段
运行时必须从 RunnableConfig 中拿到 thread_id,因为 checkpointer 需要它区分不同会话或任务执行线。
真实源码签名
def get_state(self, config: RunnableConfig, *, subgraphs: bool = False) -> StateSnapshot: ...其中 config 需要包含:
{"configurable": {"thread_id": "travel-thread-001"}}调用方与被调用方
用户代码 ↓graph.invoke / graph.get_state ↓ensure_config / merge_configs ↓config["configurable"]["thread_id"] ↓checkpointer.get_tuple / put / list输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | RunnableConfig | 包含 configurable.thread_id |
| 输出 | 归一化后的 config | thread_id 转为字符串,合并默认 config |
| 状态变化 | 无 | 只是准备定位信息 |
| 副作用 | 无 | 后续 saver 调用才有副作用 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def normalize_checkpoint_config(graph, config): # 1. 确保 config 结构存在 config = ensure_config(config)
# 2. 合并 graph 自身默认 config if graph.config: config = merge_configs(graph.config, config)
# 3. 从 configurable 中读取 thread_id thread_id = config["configurable"]["thread_id"]
# 4. thread_id 作为存储 key,统一转成 str if not isinstance(thread_id, str): config["configurable"]["thread_id"] = str(thread_id)
# 5. 保留可选 checkpoint_id / checkpoint_ns checkpoint_id = config["configurable"].get("checkpoint_id") checkpoint_ns = config["configurable"].get("checkpoint_ns", "")
return config逐段解释
ensure_config() 保证 configurable 等字段存在,避免后续直接取嵌套 key 时失败。
合并默认 config 是为了让 graph 编译期或 wrapper 设置的默认配置与当前调用配置同时生效。
thread_id 是 checkpoint 的主定位键。不同 thread_id 对应不同 checkpoint 序列。
checkpoint_id 是可选的历史定位键,用于 time travel、replay、fork 或读取指定 checkpoint。
checkpoint_ns 用于子图或命名空间隔离,避免父图和子图 checkpoint 冲突。
正常路径
config={"configurable": {"thread_id": "travel-thread-001"}} ↓ensure_config ↓thread_id="travel-thread-001" ↓checkpointer 使用 thread_id 存取 checkpoint关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
缺少 thread_id | checkpointer 无法定位 | 运行或状态读取失败 |
thread_id 非字符串 | 转成字符串 | 避免存储 key 类型不一致 |
传入 checkpoint_id | 读取指定 checkpoint | time travel / replay |
传入 checkpoint_ns | 进入指定 namespace | 子图状态读取 |
设计原因与工程影响
thread_id 不只是“会话 ID”,而是 checkpoint 存储主键。生产系统应该把它设计成稳定、可追踪、长度受控的 ID,例如:
user_id + conversation_idworkflow_run_idorder_id + task_idUUID不要随便每次生成新 thread_id,否则就无法恢复同一条执行线。
源码证据
langgraph/pregel/main.py::Pregel.get_statelanggraph/checkpoint/base/__init__.py::BaseCheckpointSaverhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.pyhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/base/__init__.py
8.4 Super-step 与 checkpoint 保存主链源码解剖
职责与所处阶段
Pregel 运行时以 super-step 推进图执行。每个 super-step 会执行一批准备好的 task,收集 writes,应用到 channels,并在边界保存 checkpoint。
真实源码签名
Pregel 主循环分散在 pregel/main.py、pregel/loop.py、pregel/algo.py 中。这里用运行时主链的语义签名表示:
def stream(self, input: InputT | Command | None, config: RunnableConfig | None = None, **kwargs) -> Iterator[Any]: ...调用方与被调用方
graph.invoke ↓graph.stream ↓Pregel loop ↓prepare_next_tasks ↓run tasks ↓apply_writes ↓create_checkpoint ↓checkpointer.put / put_writes输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | input、config | 初始 state 或 resume command |
| 输出 | chunks 或最终 output | 取决于 invoke/stream |
| 状态变化 | channels 更新 | 节点 partial update 通过 reducer 合并到 state |
| 副作用 | checkpointer 写入 | 保存 checkpoint 和 writes |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def pregel_run(input_value, config): # 1. 准备 checkpoint 定位信息 config = normalize_checkpoint_config(config) checkpointer = resolve_checkpointer(config)
# 2. 读取已有 checkpoint;没有则创建空 checkpoint saved = checkpointer.get_tuple(config) if checkpointer else None checkpoint = saved.checkpoint if saved else empty_checkpoint()
# 3. 如果传入 input,则作为 INPUT 写入 if input_value is not None: input_writes = map_input_to_channel_writes(input_value) apply_writes(checkpoint, channels, input_writes) checkpoint = create_checkpoint(checkpoint, channels, step=-1) if checkpointer: config = checkpointer.put(config, checkpoint, metadata={"source": "input"}, new_versions=...)
# 4. 进入 Pregel super-step 循环 while True: tasks = prepare_next_tasks(checkpoint, channels, nodes, config) if not tasks: break
# 5. 执行当前 super-step 中所有 task all_task_writes = [] for task in run_ready_tasks(tasks): try: writes = task.node.invoke(task.input, task.config) all_task_writes.append((task.id, writes)) if checkpointer: checkpointer.put_writes(config, writes, task_id=task.id, task_path=task.path) except Exception as exc: error_write = make_error_write(exc) all_task_writes.append((task.id, error_write)) if checkpointer: checkpointer.put_writes(config, [error_write], task_id=task.id, task_path=task.path) raise
# 6. super-step 边界:应用 writes,更新 channel versions apply_writes(checkpoint, channels, all_task_writes)
# 7. 创建并保存新的 checkpoint checkpoint = create_checkpoint(checkpoint, channels, step=current_step) if checkpointer: config = checkpointer.put(config, checkpoint, metadata={"source": "loop", "step": current_step}, new_versions=...)
# 8. 继续下一 super-step current_step += 1
# 9. 从 channels 读取输出 return read_output_channels(channels)逐段解释
第一段准备 config 和 checkpointer。所有 checkpoint 操作都依赖 thread_id,否则无法定位存储位置。
第二段读取已有 checkpoint。若同一 thread_id 曾经执行过,运行时可以从保存的状态继续;若没有,则创建空 checkpoint。
第三段把用户输入映射成 input writes。LangGraph 的 state 更新统一通过 channel writes 表达,即使初始 input 也要进入这个机制。
第四段进入 Pregel 循环。每一轮 super-step 根据 channel versions 和 versions_seen 计算哪些 task 需要执行。
第五段执行 task,并调用 put_writes() 保存中间写入。这样即使在完整 checkpoint 前发生 interrupt 或异常,系统也能知道 task 已经产生了什么。
第六段在 super-step 边界应用 writes。Reducer、LastValue、BinaryOperatorAggregate 等 channel 合并逻辑在这里发挥作用。
第七段保存完整 checkpoint。这个 checkpoint 代表“本轮 super-step 应用完所有 writes 后”的稳定状态。
第八段进入下一轮,直到没有 next tasks 或到达 END。
正常路径
read checkpoint ↓prepare tasks ↓run tasks ↓put_writes ↓apply_writes ↓put checkpoint ↓next super-step关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 没有历史 checkpoint | 创建 empty checkpoint | 从头执行 |
| 有历史 checkpoint | 从 latest checkpoint 读取 | 继续当前 thread |
| 节点产生 interrupt | 保存 interrupt write | 图暂停,等待 resume |
| 节点异常 | 保存 error write 或传播异常 | 可用于故障诊断 |
| 超过 recursion_limit | 抛 GraphRecursionError | 防止无限循环 |
设计原因与工程影响
super-step 边界是 checkpoint 的合理保存点,因为同一轮可能有多个节点并行执行,必须等这些 writes 合并后,state 才是稳定可恢复状态。
这意味着:
checkpoint 不是“每个节点调用前后随便存一下”。checkpoint 是 Pregel 调度语义下的状态快照。源码证据
langgraph/pregel/main.py::Pregellanggraph/pregel/algo.py::apply_writeslanggraph/checkpoint/base/__init__.py::create_checkpointhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.pyhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/algo.py
8.5 Pregel.get_state() 源码解剖
职责与所处阶段
get_state() 读取某个 thread 的最新 checkpoint 或指定 checkpoint,并转换为用户可读的 StateSnapshot。
真实源码签名
def get_state( self, config: RunnableConfig, *, subgraphs: bool = False,) -> StateSnapshot: ...调用方与被调用方
用户代码 graph.get_state(config) ↓Pregel.get_state ↓ensure_config / merge_configs ↓checkpointer.get_tuple(config) ↓Pregel._prepare_state_snapshot ↓StateSnapshot输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | RunnableConfig | 必须定位 thread,可选 checkpoint_id |
| 输出 | StateSnapshot | 最新或指定 checkpoint 对应的用户可读状态 |
| 状态变化 | 无 | 只读操作 |
| 副作用 | 无 | 不写 checkpoint |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def get_state(self, config, *, subgraphs=False): # 1. 从 config 或 graph 自身拿 checkpointer config = ensure_config(config) checkpointer = config["configurable"].get("checkpointer", self.checkpointer)
# 2. 如果没有 checkpointer,不能读取 state if not checkpointer: raise ValueError("No checkpointer set")
# 3. 处理 subgraph checkpoint namespace checkpoint_ns = config["configurable"].get("checkpoint_ns", "") if checkpoint_ns and "checkpointer" not in config["configurable"]: subgraph = find_subgraph_by_checkpoint_ns(checkpoint_ns) return subgraph.get_state( patch_configurable(config, {"checkpointer": checkpointer}), subgraphs=subgraphs, )
# 4. 合并默认 config,并规范 thread_id 类型 config = merge_configs(self.config, config) if self.config else config thread_id = config["configurable"]["thread_id"] if not isinstance(thread_id, str): config["configurable"]["thread_id"] = str(thread_id)
# 5. 从 checkpointer 读取 checkpoint tuple saved = checkpointer.get_tuple(config)
# 6. 转换为 StateSnapshot return self._prepare_state_snapshot( config=config, saved=saved, recurse=checkpointer if subgraphs else None, apply_pending_writes="checkpoint_id" not in config["configurable"], )逐段解释
第一段解析 checkpointer。运行时可能从 graph 自身获取,也可能从 config 中注入,尤其是子图和 namespace 场景。
第二段如果没有 checkpointer,get_state() 没有数据来源,因此直接失败。这是为什么只 compile() 不传 checkpointer 时不能使用状态历史能力。
第三段处理子图 namespace。如果 checkpoint namespace 指向子图,get_state() 会路由到对应 subgraph 的 get_state()。
第四段规范 thread_id。存储层通常以字符串作为 key,统一类型避免查询不到。
第五段调用 checkpointer.get_tuple(config)。这里读取的不只是 checkpoint,还包括 metadata、parent_config、pending_writes。
第六段调用 _prepare_state_snapshot(),把底层 checkpoint tuple 转成用户可理解的 StateSnapshot。
正常路径
graph.get_state(config) ↓checkpointer.get_tuple(config) ↓_prepare_state_snapshot ↓StateSnapshot(values, next, config, metadata, tasks, interrupts)关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| 没有 checkpointer | 抛 ValueError | 无法读取 state |
| 没有 saved checkpoint | 返回空 StateSnapshot | 表示 thread 尚无状态 |
指定 checkpoint_id | 读取历史 checkpoint | 不自动应用 pending writes |
未指定 checkpoint_id | 读取 latest | 可应用 pending writes 形成当前视图 |
subgraphs=True | 递归读取子图状态 | tasks 中可包含子图 snapshot |
设计原因与工程影响
get_state() 返回的是用户视图,不是底层 checkpoint。它会补充:
values:当前 state 值next:下一步将执行的节点名tasks:下一步 task 详情interrupts:当前中断信息metadata:checkpoint 元数据parent_config:父 checkpoint这使它适合做:
HITL 审批页面debug 面板失败恢复检查运行进度展示time travel 入口源码证据
langgraph/pregel/main.py::Pregel.get_statehttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py
8.6 _prepare_state_snapshot() 源码解剖
职责与所处阶段
_prepare_state_snapshot() 是 CheckpointTuple → StateSnapshot 的核心转换函数。它把底层 checkpoint 恢复成 channels,再计算 next tasks、tasks、interrupts 和用户可读 values。
真实源码签名
def _prepare_state_snapshot( self, config: RunnableConfig, saved: CheckpointTuple | None, recurse: BaseCheckpointSaver | None = None, apply_pending_writes: bool = False,) -> StateSnapshot: ...调用方与被调用方
Pregel.get_state / get_state_history ↓Pregel._prepare_state_snapshot ↓channels_from_checkpoint ↓prepare_next_tasks ↓apply_writes(optional pending writes) ↓read_channels ↓StateSnapshot输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | `CheckpointTuple | None` |
| 输出 | StateSnapshot | 面向用户的状态快照 |
| 状态变化 | 局部 channels | 只在内存中恢复和应用 pending writes,不写 saver |
| 副作用 | 无 | 只读转换 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def _prepare_state_snapshot(config, saved, recurse=None, apply_pending_writes=False): # 1. 没有 checkpoint 时,返回空 snapshot if saved is None: return StateSnapshot( values={}, next=(), config=config, metadata=None, created_at=None, parent_config=None, tasks=(), interrupts=(), )
# 2. 迁移旧格式 checkpoint self._migrate_checkpoint(saved.checkpoint)
# 3. 根据 checkpoint 恢复 channels / managed values step = saved.metadata.get("step", -1) + 1 stop = step + 2 channels, managed = channels_from_checkpoint( self.channels, saved.checkpoint, saver=self.checkpointer, config=saved.config, )
# 4. 计算该 checkpoint 之后的 next tasks next_tasks = prepare_next_tasks( checkpoint=saved.checkpoint, pending_writes=saved.pending_writes or [], nodes=self.nodes, channels=channels, managed=managed, config=saved.config, step=step, stop=stop, for_execution=True, store=self.store, checkpointer=self.checkpointer, )
# 5. 处理子图状态 task_states = {} if recurse: for task in next_tasks.values(): if task.name in self.get_subgraphs(): task_states[task.id] = get_subgraph_state(task, recurse)
# 6. 可选应用 pending writes,得到更贴近当前状态的视图 if apply_pending_writes and saved.pending_writes: for task_id, channel, value in saved.pending_writes: if channel in ("ERROR", "INTERRUPT"): continue if task_id in next_tasks: next_tasks[task_id].writes.append((channel, value))
tasks_with_writes = [task for task in next_tasks.values() if task.writes] if tasks_with_writes: apply_writes(saved.checkpoint, channels, tasks_with_writes, ...)
# 7. 组装 task 可读信息和 interrupts tasks = tasks_w_writes( next_tasks.values(), saved.pending_writes, task_states, self.stream_channels_asis, )
# 8. 从 channels 读取用户可见 values values = read_channels(channels, self.stream_channels_asis)
# 9. 返回 StateSnapshot return StateSnapshot( values=values, next=tuple(t.name for t in next_tasks.values() if not t.writes), config=patch_checkpoint_map(saved.config, saved.metadata), metadata=saved.metadata, created_at=saved.checkpoint["ts"], parent_config=patch_checkpoint_map(saved.parent_config, saved.metadata), tasks=tasks, interrupts=tuple(i for task in tasks for i in task.interrupts), )逐段解释
第一段处理无 checkpoint 情况。这个分支让 get_state() 对新 thread 返回空快照,而不是直接崩溃。
第二段处理 checkpoint 版本迁移。LangGraph 版本演进时,旧 checkpoint 可能需要迁移到当前 channel layout。
第三段从 checkpoint 恢复 channels。用户看到的是 state dict,但运行时恢复的是 channel 对象。
第四段计算 next tasks。checkpoint 不只是保存值,还保存 versions_seen 等调度信息,因此可以根据旧状态算出下一步该执行什么。
第五段处理子图。若 subgraphs=True,StateSnapshot.tasks 中可以携带子图状态或子图 config。
第六段可选应用 pending writes。对于最新状态读取,pending writes 可以让用户看到未完成 task 的写入效果;对于指定历史 checkpoint,通常不应用 pending writes,以保持历史点原貌。
第七段把底层 task 与 writes 组装成可读 task 信息。
第八段读取 channels 形成 values。
第九段返回 StateSnapshot,它是调试、恢复、HITL、time travel 的核心对象。
正常路径
CheckpointTuple ↓channels_from_checkpoint ↓prepare_next_tasks ↓apply_pending_writes(optional) ↓read_channels ↓StateSnapshot关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
saved is None | 返回空 snapshot | 新 thread 无状态 |
subgraphs=True | 递归读取子图 | tasks 中包含子图状态 |
apply_pending_writes=True | 应用 pending writes | 最新状态更完整 |
| 有 interrupt writes | 聚合到 interrupts | 支持 HITL 展示 |
设计原因与工程影响
_prepare_state_snapshot() 是 checkpoint 系统的“解码层”。它把底层可持久化格式转成用户可理解的运行状态。
工程上,不建议你直接解析 checkpoint["channel_values"],而应使用:
graph.get_state(config)graph.get_state_history(config)因为只有 StateSnapshot 才包含 next、tasks、interrupts、metadata、parent_config 等完整运行语义。
源码证据
langgraph/pregel/main.py::Pregel._prepare_state_snapshothttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py
8.7 get_state_history() 源码解剖
职责与所处阶段
get_state_history() 列出某个 thread 的 checkpoint 历史,并把每个 checkpoint tuple 转成 StateSnapshot。
真实源码签名
def get_state_history( self, config: RunnableConfig, *, filter: dict[str, Any] | None = None, before: RunnableConfig | None = None, limit: int | None = None,) -> Iterator[StateSnapshot]: ...调用方与被调用方
用户代码 graph.get_state_history(config) ↓Pregel.get_state_history ↓checkpointer.list(config, before, limit, filter) ↓Pregel._prepare_state_snapshot for each tuple ↓Iterator[StateSnapshot]输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | RunnableConfig | 包含 thread_id,可选 before/filter/limit |
| 输出 | Iterator[StateSnapshot] | 历史快照序列 |
| 状态变化 | 无 | 只读 |
| 副作用 | 无 | 读取存储 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def get_state_history(self, config, *, filter=None, before=None, limit=None): # 1. 准备 config 和 checkpointer config = ensure_config(config) checkpointer = config["configurable"].get("checkpointer", self.checkpointer) if not checkpointer: raise ValueError("No checkpointer set")
# 2. 子图 namespace 转发 checkpoint_ns = config["configurable"].get("checkpoint_ns", "") if checkpoint_ns and "checkpointer" not in config["configurable"]: subgraph = find_subgraph_by_checkpoint_ns(checkpoint_ns) yield from subgraph.get_state_history( patch_configurable(config, {"checkpointer": checkpointer}), filter=filter, before=before, limit=limit, ) return
# 3. 合并 config 并规范 thread_id config = merge_configs( self.config, config, {"configurable": {"thread_id": str(config["configurable"]["thread_id"])}} )
# 4. 从 saver 列出 checkpoint tuple checkpoint_tuples = list( checkpointer.list( config, before=before, limit=limit, filter=filter, ) )
# 5. 逐个转换为 StateSnapshot for checkpoint_tuple in checkpoint_tuples: yield self._prepare_state_snapshot( checkpoint_tuple.config, checkpoint_tuple, )逐段解释
第一段和 get_state() 一样,必须有 checkpointer。
第二段处理子图命名空间,说明 state history 也可以定位到子图。
第三段规范 thread_id,保证查询同一个 thread。
第四段调用 checkpointer.list()。这里是 get_state_history() 与 get_state() 的关键差别:前者列出多个 checkpoint,后者只读取一个 checkpoint tuple。
第五段逐个转成 StateSnapshot,因此用户得到的是可读快照序列,而不是底层存储记录。
正常路径
thread_id ↓checkpointer.list ↓CheckpointTuple A/B/C ↓_prepare_state_snapshot ↓StateSnapshot A/B/C关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
limit 设置 | 限制返回数量 | 控制读取成本 |
before 设置 | 从某个 checkpoint 之前读取 | 支持分页和 time travel 定位 |
filter 设置 | 按 metadata 过滤 | 支持审计查询 |
saver 未实现 list | 抛异常 | 无法读取历史 |
设计原因与工程影响
get_state_history() 是 time travel、审计、可观测性和故障诊断的基础。它不仅告诉你“现在是什么状态”,还能告诉你:
状态是怎么一步步变成现在这样的?哪一步写入了错误字段?哪个 checkpoint 可以作为 replay/fork 起点?是否在某一步进入了 interrupt?源码证据
langgraph/pregel/main.py::Pregel.get_state_historyhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py
8.8 Callback、事件与可观测性
| 时机 | 事件 | 携带数据 | 失败行为 |
|---|---|---|---|
| invoke 开始 | run start / graph start | input、config、thread_id | 配置错误直接失败 |
| node 开始 | task start | node name、task id、state slice | node 异常进入错误路径 |
| task 写入 | put_writes | task_id、writes、task_path | saver 异常可能导致执行失败 |
| step 完成 | put checkpoint | checkpoint、metadata、new_versions | saver 异常影响 durable guarantee |
| get_state | state read | checkpoint tuple、StateSnapshot | 无 checkpointer 失败 |
| get_history | history read | checkpoint tuple list | saver.list 失败则无法读取历史 |
8.9 运行时主链总结
Input + thread_id ↓Pregel runtime ↓Super-step task execution ↓Writes ↓Checkpoint ↓StateSnapshot / History9. 关键分支、异常与边界
本章回答:
当输入、执行模式或运行结果偏离正常主链时,checkpoint 系统如何分流、恢复、终止或失败?
9.1 分支矩阵
| 分支类型 | 触发条件 | 核心函数 | 结果 |
|---|---|---|---|
| 无 checkpointer | compile() 未传 checkpointer,但调用 get_state | Pregel.get_state() | 抛 ValueError |
| 新 thread | thread_id 无历史 checkpoint | checkpointer.get_tuple() | 返回空 snapshot 或从头执行 |
| 已有 thread | thread_id 已有 checkpoint | get_tuple() | 从最新 checkpoint 恢复状态 |
| 指定 checkpoint | config 带 checkpoint_id | get_state() / invoke(None, config) | 读取或 replay 历史点 |
| 子图状态 | config 带 checkpoint_ns | get_state() | 转发到 subgraph |
| pending writes | checkpoint tuple 带 pending writes | _prepare_state_snapshot() | 可选应用 writes |
| interrupt | task writes 中有 interrupt | put_writes() / StateSnapshot.interrupts | 暂停并等待 resume |
| store 失效 | 外部 saver 不可用 | saver 方法 | 读写失败 |
9.2 同步与异步分支
| 维度 | 同步路径 | 异步路径 |
|---|---|---|
| 状态读取 | get_state() | aget_state() |
| 历史读取 | get_state_history() | aget_state_history() |
| saver 方法 | get_tuple/list/put/put_writes | aget_tuple/alist/aput/aput_writes |
| 调度方式 | 阻塞当前线程 | await 异步 IO |
| 生产建议 | 本地或低并发可用 | 数据库/网络 saver 更推荐 async |
9.3 Batch、Stream、Parallel 或路由分支
Checkpoint 本身不是 batch、stream、parallel 的业务 API,但它与这些运行模式协作:
Batch:每个输入通常应有独立 thread_id,否则状态会相互污染。Stream:流式输出过程中仍会在 step 边界保存 checkpoint。Parallel:同一 super-step 多节点并行执行后,合并 writes 再保存 checkpoint。Routing:条件边决定 next tasks,checkpoint 保存 next 所需的 versions_seen 等调度信息。当前对象不直接决定 graph 的分支拓扑。分支由 StateGraph 的 edges / branches 和 Pregel 调度决定;checkpoint 负责把分支执行后的状态与调度信息保存下来。
9.4 异常分类
| 异常类别 | 抛出位置 | 是否可恢复 | 处理策略 | 是否反馈上层 |
|---|---|---|---|---|
| 缺少 checkpointer | get_state() / get_state_history() | 否 | compile 时传 checkpointer | 是 |
| 缺少 thread_id | saver 定位或 config 读取 | 否 | 调用时传 configurable.thread_id | 是 |
| saver 写入失败 | put() / put_writes() | 视存储而定 | 重试或失败恢复 | 是 |
| saver 读取失败 | get_tuple() / list() | 视存储而定 | 重试或降级 | 是 |
| 节点执行异常 | task runtime | 可诊断 | 保存 error writes 后传播或 handler 处理 | 是 |
| recursion limit | Pregel loop | 可调整 | 提高 limit 或修复循环条件 | 是 |
| 历史 checkpoint 不兼容 | _migrate_checkpoint() | 部分可恢复 | 迁移或版本锁定 | 是 |
9.5 异常路径源码解剖
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
def read_state_with_checkpoint(config): try: checkpointer = resolve_checkpointer(config) if not checkpointer: raise ValueError("No checkpointer set")
thread_id = config["configurable"]["thread_id"] if thread_id is None: raise ValueError("thread_id required")
saved = checkpointer.get_tuple(config) return prepare_state_snapshot(config, saved)
except KeyError as exc: # config 结构错误,例如缺少 configurable/thread_id raise ValueError("Invalid checkpoint config") from exc
except StorageTimeout as exc: # 具体 saver 的存储异常,可重试 raise RecoverableCheckpointError from exc
except Exception: # 不吞掉未知错误,避免返回虚假状态 raise逐段解释
缺少 checkpointer 和缺少 thread_id 是配置错误,不能被静默降级。否则系统可能让用户误以为状态已保存,但实际上没有 durable guarantee。
存储超时等具体 saver 异常可以在外层做 retry,但必须注意:checkpoint 写入可能不是幂等的,生产 saver 应根据 checkpoint_id 和 task_id 设计幂等写入。
未知异常不应该被吞掉,因为 checkpoint 是恢复能力的可信来源。返回一个“看似正常但实际不完整”的状态,比直接失败更危险。
9.6 Retry、Fallback 与恢复边界
| 机制 | 适用条件 | 不适用条件 | 幂等要求 |
|---|---|---|---|
| Retry | saver 临时网络故障、数据库短暂不可用 | 节点副作用已经执行但未记录 | put / put_writes 应按 checkpoint_id/task_id 幂等 |
| Fallback | 主 saver 不可用时切只读或降级模式 | 需要严格 durable guarantee 的支付/退款流程 | 必须明确告警 |
| Repair | checkpoint 版本迁移或 metadata 修复 | 业务 state 语义错误 | 需要离线脚本和备份 |
| Replay | 从历史 checkpoint 重新执行后续节点 | 非幂等工具调用未隔离 | 工具层必须幂等或加确认 |
| Fork | 从历史 checkpoint 改 state 后继续 | 原执行必须保持不变的审计场景 | 新 checkpoint 应保留 parent 链 |
9.7 停止条件与保护上限
正常结束:图执行到 END,latest StateSnapshot.next 为空。提前结束:节点或 Command 显式终止后保存 checkpoint。人工中断:interrupt 写入 checkpoint,StateSnapshot.interrupts 非空。框架保护:recursion_limit 限制 super-step 数,超限抛错。异常失败:节点或 saver 异常传播,可能保留 pending writes / error writes。9.8 能力边界
| 容易误判的能力 | 实际提供者 | 本篇对象的真实职责 |
|---|---|---|
| 长期用户偏好记忆 | Store / 数据库 / 向量库 | Checkpointer 只保存 thread graph state |
| 工具幂等控制 | Tool 层 / 业务 API | Checkpoint 只能记录执行状态,不能保证外部副作用安全 |
| 业务规则校验 | Verifier / Policy Engine | Checkpoint 不判断业务对错 |
| 无限循环控制 | recursion_limit / route 逻辑 | Checkpoint 可记录循环过程,但不自动修复循环 |
| 高可用存储 | 具体 saver 实现 | BaseCheckpointSaver 只是接口契约 |
10. 扩展机制与框架协作
本章回答:
checkpoint 允许在哪里插入自定义存储,它如何与 memory、HITL、time travel、subgraph 等能力协作?
10.1 扩展点总览
| 扩展点 | 扩展方式 | 执行时机 | 可修改内容 | 约束 |
|---|---|---|---|---|
| 自定义 checkpointer | 实现 BaseCheckpointSaver | get/put/list/writes | 存储介质、序列化、索引 | 必须保留接口语义 |
thread_id 策略 | config 约定 | invoke/get_state/history | 会话隔离粒度 | 必须稳定、唯一、长度受控 |
checkpoint_ns | config / 子图 namespace | 子图状态读写 | 父子图隔离 | 不建议业务随意拼接内部 ns |
metadata | config / runtime 自动生成 | put checkpoint | 审计、过滤、history 查询 | 不应放大对象或敏感信息 |
update_state | 公共 API | time travel / fork / 测试 | 修改历史点后的 state | 不是原地回滚 |
| persistent saver | Postgres/SQLite/custom | 生产运行 | 跨进程持久化 | 需要事务、索引、清理策略 |
10.2 自定义 Checkpointer 源码解剖
职责与所处阶段
自定义 checkpointer 需要实现 BaseCheckpointSaver 契约,让 Pregel 可以把 checkpoint 写入任意持久介质。
真实源码签名
class CustomSaver(BaseCheckpointSaver): def get_tuple(self, config: RunnableConfig) -> CheckpointTuple | None: ... def list(self, config: RunnableConfig | None, *, filter=None, before=None, limit=None) -> Iterator[CheckpointTuple]: ... def put(self, config: RunnableConfig, checkpoint: Checkpoint, metadata: CheckpointMetadata, new_versions: ChannelVersions) -> RunnableConfig: ... def put_writes(self, config: RunnableConfig, writes: Sequence[tuple[str, Any]], task_id: str, task_path: str = "") -> None: ...调用方与被调用方
Pregel runtime ↓BaseCheckpointSaver methods ↓CustomSaver ↓DB / Object Storage / KV Store输入、输出与状态变化
| 项目 | 类型 | 说明 |
|---|---|---|
| 输入 | config/checkpoint/metadata/writes | Pregel 运行时状态 |
| 输出 | checkpoint tuple / updated config | 状态读取或写入定位 |
| 状态变化 | 外部存储 | checkpoint rows、writes rows、blob rows |
| 副作用 | 有 | 数据库写入、序列化 |
细粒度伪代码
以下为保留关键控制流的简化伪代码,不是源码逐字复制:
class CustomSaver(BaseCheckpointSaver): def put(self, config, checkpoint, metadata, new_versions): thread_id = config["configurable"]["thread_id"] checkpoint_ns = config["configurable"].get("checkpoint_ns", "") checkpoint_id = checkpoint["id"] parent_id = config["configurable"].get("checkpoint_id")
with transaction(): # 1. 保存 checkpoint metadata 和 parent 链 upsert_checkpoint_row( thread_id=thread_id, checkpoint_ns=checkpoint_ns, checkpoint_id=checkpoint_id, parent_id=parent_id, metadata=serialize(metadata), )
# 2. 保存 channel blobs for channel, version in checkpoint["channel_versions"].items(): value = checkpoint["channel_values"].get(channel, EMPTY) upsert_blob(thread_id, checkpoint_ns, channel, version, serialize(value))
# 3. 返回包含新 checkpoint_id 的 config return patch_configurable(config, {"checkpoint_id": checkpoint_id})
def get_tuple(self, config): thread_id = config["configurable"]["thread_id"] checkpoint_ns = config["configurable"].get("checkpoint_ns", "") checkpoint_id = config["configurable"].get("checkpoint_id")
# 4. 不指定 checkpoint_id 时读取 latest row = select_checkpoint(thread_id, checkpoint_ns, checkpoint_id or latest()) if row is None: return None
checkpoint = reconstruct_checkpoint_from_row_and_blobs(row) pending_writes = select_writes(thread_id, checkpoint_ns, row.checkpoint_id)
return CheckpointTuple( config=make_config(thread_id, checkpoint_ns, row.checkpoint_id), checkpoint=checkpoint, metadata=deserialize(row.metadata), parent_config=make_parent_config(row.parent_id), pending_writes=pending_writes, )逐段解释
put() 应以事务方式保存 checkpoint row 和 channel blobs,否则可能出现 metadata 写入成功但 blob 写入失败的半状态。
put() 返回的 config 应包含新的 checkpoint_id,让后续步骤能定位刚写入的 checkpoint。
get_tuple() 不指定 checkpoint_id 时应读取 latest checkpoint;指定时读取历史 checkpoint。
pending_writes 必须一起返回,否则 _prepare_state_snapshot() 无法正确恢复未完成 task 的中间状态。
正常路径
Pregel checkpoint ↓CustomSaver.put ↓DB rows / blobs ↓CustomSaver.get_tuple ↓CheckpointTuple ↓StateSnapshot关键分支与异常路径
| 条件 | 行为 | 结果 |
|---|---|---|
| checkpoint_id 已存在 | upsert 或幂等忽略 | 支持重试 |
| parent checkpoint 缺失 | 报错或拒绝写入 | 防止历史链断裂 |
| blob 缺失 | 恢复失败 | 需要事务或校验 |
| writes 重复 | 按 task_id/write_idx 幂等 | 防止重试重复写 |
设计原因与工程影响
自定义 saver 的核心难点不在“把 dict 存进去”,而在:
1. checkpoint 与 writes 的一致性。2. checkpoint_id / parent_config 的历史链。3. channel blob 的版本化。4. task writes 的幂等写入。5. list/history 的排序与分页。源码证据
langgraph/checkpoint/base/__init__.py::BaseCheckpointSaverlanggraph/checkpoint/memory/__init__.py::InMemorySaverhttps://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/base/__init__.py
10.3 扩展调用链
Framework Entry: graph.invoke / get_state / get_history ↓Pregel runtime ↓BaseCheckpointSaver interface ↓Concrete checkpointer ↓Storage transaction / memory dict ↓CheckpointTuple ↓Pregel._prepare_state_snapshot ↓StateSnapshot10.4 与相邻框架组件的协作
| 相邻组件 | 输入协议 | 输出协议 | 协作边界 |
|---|---|---|---|
StateGraph | compile(checkpointer=...) | CompiledStateGraph | 构建期注入 checkpointer |
Pregel | channels、nodes、config | checkpoint writes / StateSnapshot | 运行期调度和状态恢复 |
Reducer / Channel | writes | checkpointable channel value | checkpoint 读取 channel 快照 |
interrupt() | interrupt write | StateSnapshot.interrupts | checkpoint 保存暂停点 |
Command(resume=...) | resume input | 继续图执行 | 依赖 thread_id 定位暂停状态 |
update_state() | checkpoint config + values | 新 checkpoint config | 用于 time travel fork |
Store | namespace/key/value | 长期数据 | 不替代 thread checkpoint |
10.5 公共扩展接口与内部实现
业务代码可以依赖: graph.invoke(input, config={"configurable": {"thread_id": ...}}) graph.get_state(config) graph.get_state_history(config) graph.update_state(config, values=...) builder.compile(checkpointer=...) BaseCheckpointSaver 契约
业务代码避免依赖: checkpoint["versions_seen"] 的具体内部布局 InMemorySaver.storage 的嵌套 dict 结构 Pregel._prepare_state_snapshot 私有方法 checkpoint_ns 的内部拼接细节10.6 自定义扩展示例
示例只展示扩展契约,不重新实现完整数据库 saver:
from langgraph.checkpoint.base import BaseCheckpointSaver, CheckpointTuple
class AuditedCheckpointSaver(BaseCheckpointSaver): def __init__(self, inner: BaseCheckpointSaver, audit_logger): super().__init__(serde=inner.serde) self.inner = inner self.audit_logger = audit_logger
def get_tuple(self, config): self.audit_logger.info("checkpoint.get", config=config) return self.inner.get_tuple(config)
def list(self, config, *, filter=None, before=None, limit=None): self.audit_logger.info("checkpoint.list", config=config, limit=limit) yield from self.inner.list(config, filter=filter, before=before, limit=limit)
def put(self, config, checkpoint, metadata, new_versions): self.audit_logger.info( "checkpoint.put", thread_id=config["configurable"].get("thread_id"), checkpoint_id=checkpoint["id"], metadata=metadata, ) return self.inner.put(config, checkpoint, metadata, new_versions)
def put_writes(self, config, writes, task_id, task_path=""): self.audit_logger.info("checkpoint.put_writes", task_id=task_id) return self.inner.put_writes(config, writes, task_id, task_path)说明:
- 扩展点接收
config/checkpoint/metadata/writes。 - 扩展点允许增强日志、加密、压缩、审计、metrics。
- 扩展点必须返回与内部 saver 一致的结果。
- 异常应向上传播,除非有明确 retry 策略。
- 包装 saver 会影响可观测性,但不应改变 checkpoint 语义。
10.7 选择扩展还是重写流程
| 条件 | 选择扩展点 | 选择更底层框架 |
|---|---|---|
| 只是换存储介质 | 是 | 否 |
| 只是增加审计日志 | 是 | 否 |
| 需要改变 checkpoint 保存时机 | 否 | 是,但风险高 |
| 需要节点内长期偏好记忆 | 否 | 用 Store 或外部 DB |
| 需要分布式任务队列和强事务工作流 | 视情况 | 可能需要结合 Temporal / workflow engine |
11. 工程决策与适用场景
11.1 适用场景
| 场景 | 是否推荐 | 原因 |
|---|---|---|
| 多轮旅行规划助手 | 是 | 同一 thread 保存需求、计划、修改历史 |
| HITL 审批流程 | 是 | interrupt 必须依赖 checkpoint resume |
| 长任务研究 Agent | 是 | 可在步骤间恢复、回放、审计 |
| 一次性纯文本改写 | 否 | 没必要引入 checkpoint 成本 |
| 高风险退款/订单流程 | 是,但需持久化 saver | 需要审计、恢复、故障后补偿 |
| 跨用户长期偏好记忆 | 不单独推荐 | 应用 Store 或业务数据库,不是 checkpointer |
11.2 工程决策表
| 决策点 | 推荐选择 | 前提 | 风险 |
|---|---|---|---|
| 本地调试 | InMemorySaver | 单进程、临时运行 | 进程重启丢失 |
| 本地持久开发 | SQLite saver | 需要跨重启调试 | 并发能力有限 |
| 生产服务 | Postgres/托管 saver | 多用户、多进程、审计 | 需要索引、清理、迁移 |
thread_id | 稳定业务 ID 或 UUID | 可定位会话/任务 | 随机乱用导致无法恢复 |
| checkpoint 保留策略 | 设置 retention/prune | 长会话或高频调用 | 存储无限增长 |
| time travel | 用历史 checkpoint config | 工具副作用可控 | 非幂等工具可能重复执行 |
| 故障恢复 | 从 latest checkpoint resume | 节点幂等、工具可补偿 | 外部副作用与 checkpoint 不一致 |
11.3 性能、可靠性与安全边界
性能:checkpoint 读写增加 IO 和序列化开销;大 state 会放大延迟和存储成本。可靠性:checkpoint saver 是恢复能力的关键依赖;生产必须使用持久化存储和幂等写入。安全:checkpoint 中可能包含用户输入、工具结果、业务数据;需要加密、权限和 retention。可观测性:必须记录 thread_id、checkpoint_id、step、node、writes、interrupts、errors。12. 常见误区与源码纠正
12.1 误区:thread_id 只是普通配置,可传可不传
错误原因:
初学者看到 thread_id 在 configurable 中,容易以为它只是日志字段。
源码事实:
BaseCheckpointSaver 文档和源码都把 thread_id 作为 checkpoint 的主定位键。没有 thread_id,checkpointer 无法保存、恢复、time travel 或 resume。工程影响:
生产中如果每次随机生成 thread_id,就无法恢复同一个任务;如果多个用户复用一个 thread_id,会造成状态污染。
12.2 误区:Checkpoint 就是业务 state dict
错误原因:
get_state(config).values 看起来像普通 state,因此容易把底层 checkpoint 等同于 state。
源码事实:
Checkpoint 保存 channel_values、channel_versions、versions_seen、updated_channels 等调度信息。StateSnapshot.values 只是从 channels 读出的用户可见视图。工程影响:
直接操作底层 checkpoint 容易破坏调度版本,导致恢复后 next tasks 计算错误。
12.3 误区:InMemorySaver 是生产级持久化
错误原因:
它能让 get_state 和 history 跑起来,因此容易误以为已经具备 durable execution。
源码事实:
InMemorySaver 把 storage、writes、blobs 保存在当前进程内存中。源码 docstring 明确建议只用于 debugging/testing,生产使用 Postgres 或托管 saver。工程影响:
服务重启、扩容、多进程部署都会丢失或分裂状态。
12.4 误区:Checkpoint 能自动保证外部工具幂等
错误原因:
durable execution 容易被误解为“一切都能自动恢复正确”。
源码事实:
Checkpoint 保存 graph state 和 writes。外部 API 副作用是否幂等,由工具层和业务系统保证。工程影响:
如果 replay 某个 checkpoint 后再次执行支付、退款、发邮件等工具,可能造成重复副作用。工具必须设计 idempotency key 或人工确认。
12.5 误区:Time travel 是回滚原线程
错误原因:
“回到历史状态”听起来像数据库 rollback。
源码事实:
官方 time travel 语义中,replay 是从旧 checkpoint 重新执行后续节点;fork 是通过 update_state 创建新的 checkpoint 分支,原历史保持不变。工程影响:
审计系统中不能把 fork 理解为删除或覆盖原历史;它是新增分支。
12.6 误区:Checkpoint 可以替代长期记忆
错误原因:
checkpoint 能跨调用保留 thread state,因此容易被当作 memory 全部方案。
源码事实:
官方文档区分 checkpointer 和 store:checkpointer 保存 thread graph state;store 保存跨 thread 的应用级长期数据。工程影响:
把用户长期偏好塞进 checkpoint 会导致跨会话复用困难、历史膨胀、权限边界混乱。
13. 最终心智模型与掌握检查
13.1 构建期心智模型
StateGraph builder ↓compile(checkpointer=InMemorySaver()) ↓CompiledStateGraph / Pregel ↓runtime 持有 checkpointer13.2 运行时心智模型
graph.invoke(input, config={thread_id}) ↓读取或创建 checkpoint ↓Pregel super-step 执行节点 ↓保存 writes ↓应用 state update ↓保存 checkpoint ↓返回 output13.3 分支与异常心智模型
正常路径 → latest checkpoint 可由 get_state 读取。历史路径 → get_state_history 列出 checkpoint 序列。可恢复异常 → 从最近 checkpoint 继续或 replay。不可恢复异常 → 存储损坏、缺少 thread_id、非幂等副作用需人工处理。保护上限 → recursion_limit 防止循环无限产生 checkpoint。13.4 一句话总结
LangGraph 通过
checkpointer在 Pregel super-step 边界保存Checkpoint,用thread_id组织同一执行线的 checkpoint 序列,运行时通过get_state()和get_state_history()把底层 checkpoint 转换为StateSnapshot,从而支持短期记忆、HITL、time travel、故障恢复和生产级 durable workflow。
13.5 掌握检查
- 能说清
checkpointer与store的区别。 - 能解释为什么
thread_id是 checkpoint 主键。 - 能画出
graph.invoke → super-step → writes → checkpoint的主链。 - 能说出
Checkpoint至少保存哪些字段。 - 能解释
CheckpointTuple为什么比Checkpoint多 metadata、parent_config、pending_writes。 - 能解释
StateSnapshot.values / next / tasks / interrupts的含义。 - 能解释
get_state()和get_state_history()分别调用 saver 的哪个方法。 - 能说明 time travel replay 与 fork 的区别。
- 能判断
InMemorySaver是否适合生产。 - 能说明外部工具副作用为什么不能只靠 checkpoint 保证安全。
14. 参考资料与下一篇衔接
14.1 官方概念文档
-
LangGraph Persistence
https://docs.langchain.com/oss/python/langgraph/persistence -
LangGraph Time Travel
https://docs.langchain.com/oss/python/langgraph/use-time-travel -
LangGraph Overview
https://docs.langchain.com/oss/python/langgraph/overview -
LangGraph Graph API
https://docs.langchain.com/oss/python/langgraph/graph-api
14.2 官方 API Reference
-
Checkpointing
https://reference.langchain.com/python/langgraph/checkpoints/ -
InMemorySaver
https://reference.langchain.com/python/langgraph.checkpoint/memory/InMemorySaver -
StateGraph.compile
https://reference.langchain.com/python/langgraph/graph/state/StateGraph/compile -
CompiledStateGraph.get_state
https://reference.langchain.com/python/langgraph/graph/state/CompiledStateGraph/get_state -
CompiledStateGraph.get_state_history
https://reference.langchain.com/python/langgraph/graph/state/CompiledStateGraph/get_state_history
14.3 官方源码
-
langgraph/graph/state.py::StateGraph.compile
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/graph/state.py -
langgraph/pregel/main.py::Pregel.get_state
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py -
langgraph/pregel/main.py::Pregel.get_state_history
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py -
langgraph/pregel/main.py::Pregel._prepare_state_snapshot
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/langgraph/langgraph/pregel/main.py -
langgraph/checkpoint/base/__init__.py::Checkpoint
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/base/__init__.py -
langgraph/checkpoint/base/__init__.py::BaseCheckpointSaver
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/base/__init__.py -
langgraph/checkpoint/memory/__init__.py::InMemorySaver
https://github.com/langchain-ai/langgraph/blob/1.2.7/libs/checkpoint/langgraph/checkpoint/memory/__init__.py
14.4 下一篇衔接
下一篇进入:
第 12 篇:Interrupt 与 Human-in-the-loop 源码解剖需要继续回答:
interrupt() 如何暂停图执行?Command(resume=...) 如何恢复执行?interrupt 为什么必须依赖 checkpointer 和 thread_id?HITL 审批如何与高风险工具调用结合?