LangGraph DeltaChannel 检查点数据恢复实战:用 delta-channel-dump 读取 Postgres 原始检查点数据
当 LangGraph 线程由 langgraph >= 1.2 的 DeltaChannel 格式写入后,直接回滚到旧版运行时时,消息等渠道数据会因旧版 add_messages 无法识别 EXT_DELTA_SNAPSHOT msgpack 扩展编码而静默丢失。本文以仓库中的 delta-channel-dump 工具 为主体,完整讲解其安装、命令行用法、输出格式与回写方式,并结合 libs/langgraph、libs/checkpoint、libs/checkpoint-postgres 的源码,剖析它如何从 Postgres 的 checkpoints / checkpoint_blobs / checkpoint_writes 三张表中还原 seed 与 writes,帮助你在跨版本回滚时保住线程状态。
背景:DeltaChannel 带来了什么问题
DeltaChannel 的存储模型
自 langgraph >= 1.2 起,DeltaChannel(当前标注为 Beta,API 与磁盘表示可能变化)改变了 reducer 渠道的检查点写入方式:它只在检查点 blob 中存一个哨兵,而不存完整值,重建状态时依赖重放祖先检查点的 writes 经过 reducer 还原。核心实现在 DeltaChannel:
snapshot_frequency(默认1000,必须为正整数):每 N 次更新写一次完整快照 blob;- 系统级
DELTA_MAX_SUPERSTEPS_SINCE_SNAPSHOT(默认5000):即便渠道停止接收写入,超步数达到该上界也会强制快照,从而限制重放深度; - 快照以
_DeltaSnapshot元组形式写入(定义见 serde/types.py),祖先遍历遇到它即终止;非快照步时该渠道不出现在channel_values中。
序列化层面,_DeltaSnapshot 通过 msgpack EXT 扩展编码 7(EXT_DELTA_SNAPSHOT) 写入 checkpoint_blobs,完整 EXT 编码表定义在 jsonplus.py:
| EXT code | 名称 | 用途 |
|---|---|---|
| 0 | EXT_CONSTRUCTOR_SINGLE_ARG |
单参构造(如 uuid.UUID、SecretStr) |
| 1 | EXT_CONSTRUCTOR_POS_ARGS |
位置参构造(如 Send、pathlib.Path) |
| 2 | EXT_CONSTRUCTOR_KW_ARGS |
namedtuple(_asdict) |
| 3 | EXT_METHOD_SINGLE_ARG |
方法调用 |
| 4 / 5 | EXT_PYDANTIC_V1 / V2 |
Pydantic 模型 |
| 6 | EXT_NUMPY_ARRAY |
numpy 数组 |
| 7 | EXT_DELTA_SNAPSHOT |
DeltaChannel 快照 blob |
回滚困境与这个工具的定位
README 指出的故障场景是:线程由 langgraph >= 1.2 / deepagents 0.6.x(Postgres 运行时,包括 LangGraph Server / langgraph-api)写入后,运维方需要回滚到 deepagents 0.5.x / langgraph < 1.2。旧版运行时的 add_messages reducer 不认识 EXT_DELTA_SNAPSHOT(EXT code 7)这个 msgpack ext 编码,会静默返回空列表,受影响的渠道(典型如 messages)在旧运行时里读出来就是空的。
delta-channel-dump 就是为此准备的恢复工具:它绕过 langgraph 运行时,直接读取 Postgres 中的原始 checkpoint 数据,解码 msgpack blob,输出可检查、可回放的 JSON;随后由操作者手动通过 LangGraph Server SDK 的 update_state 或 OSS 的 graph.update_state 把恢复出的值写回。由于 langgraph-api 与 OSS checkpoint-postgres 使用同一套 checkpoints / checkpoint_blobs / checkpoint_writes 表结构,该脚本对两种部署方式都适用。
安装与运行
安装
工具只依赖两个包(无需安装 langgraph 本身):
pip install "psycopg[binary]" ormsgpack
命令行用法
export DATABASE_URI=postgres://user:pass@host:5432/dbname
python3 dump.py \
--thread-id <uuid> \
--channel messages \
[--channel files ...] \
[--checkpoint-id <uuid>] \
[--checkpoint-ns ""] \
--output recovery.json
参数说明(与 dump.py 的参数解析 一一对应):
| 参数 | 是否必需 | 说明 |
|---|---|---|
--thread-id |
必需 | 线程 UUID。脚本会做 uuid.UUID(...) 校验,非法值直接报错 |
--channel |
必需,可重复 | 要恢复的渠道名,如 messages;解析后按出现顺序去重 |
--checkpoint-id |
可选 | 目标 checkpoint;默认取该线程指定命名空间下最新一个(ORDER BY checkpoint_id DESC LIMIT 1) |
--checkpoint-ns |
可选 | 命名空间,默认 ""(主图;子图线程需指定对应 ns) |
--database-uri |
可选 | Postgres URI,默认读取环境变量 DATABASE_URI;两者都没有时报错退出 |
--output |
可选 | 输出 JSON 文件路径,默认 - 表示打印到 stdout |
例如恢复某线程最新消息列表并落盘:
python3 dump.py \
--thread-id 6f9619ff-8b86-4d01-b4f9-c54c1f9c2e33 \
--channel messages \
--output recovery.json
若指定 --checkpoint-id 想恢复到某个历史时间点(时间旅行场景),把该参数指到目标 checkpoint 即可;脚本输出的 parent_checkpoint_id 表明实际遍历起点是其父节点。
工作原理:从三张表拼出 seed + writes
examples/delta-channel-dump/dump.py 是单文件自包含脚本,主流程(main)分四步:解析参数并校验 UUID → 解析目标 checkpoint → 沿父链遍历 → 组装 JSON。
第一步:确定目标 checkpoint
resolve_target_checkpoint_id:若未显式指定 --checkpoint-id,则查询 checkpoints 表中该 thread_id + checkpoint_ns 下 checkpoint_id 最大的一行作为目标;查不到则直接退出并报错。
第二步:从 target 的 parent 开始回溯
walk_parent_chain 先取出目标 checkpoint 的 parent_checkpoint_id,再对每个请求渠道调用 walk_channel,从父节点沿父链向根方向回溯,沿途做两件事:
- 收集
checkpoint_writes表中该 checkpoint、该渠道的 writes(task_id、idx、type、blob),累积到链上; - 检查该 checkpoint 的
checkpoint.channel_values中是否出现该渠道——一旦出现,该 checkpoint 就是 seed 点,遍历终止。
回溯到达根仍未找到任何祖先持有该渠道的值时,返回 delta_kind: no_seed,但仍会带上途中收集到的全部 writes。
第三步:识别 seed 的类型
_load_seed 按 channel_values 中该渠道的存储形态分三种情况,正好对应后文的 delta_kind:
channel_values[ch] is True(布尔哨兵)→snapshot类型:用channel_versions[ch]作 version 去checkpoint_blobs取 blob,解码后若带__delta_snapshot__包裹则解包,得到完整累积值;blob 行不存在或为empty时 seed 记为null。这与 PostgresSaver 历史重建逻辑 的注释一致:对_DeltaSnapshot,channel_values里的内联true只是标记,不是值。- 内联原始值(
int/float/str/bool/None) →legacy_plain:put会把原始类型直接留在 checkpoint 自身的channel_values中而不产生 blob 行(见 base.py 的说明),脚本直接内联返回。 - 其它复杂值 → 同样按
channel_versions查checkpoint_blobs解码为legacy_plain(DeltaChannel 之前的内联/blob 值)。
第四步:writes 的排序保证回放语义
_load_writes_for_checkpoint 对单个 checkpoint 按 ORDER BY task_id DESC, idx DESC 取行,而 walk_channel 在回溯结束后把累积的“新到旧”扁平列表整体 reverse,最终输出的 writes 是从旧到新——正是 reducer 回放时应处理的顺序。dump.py 的 docstring 明确说明这一排序方式与 PostgresSaver._build_delta_channels_writes_history 对齐,实现可对照 base.py#L471-L547:后者同样按 (task_id, idx) 降序收集、最后 collected.reverse() 得到 oldest-first 的 PendingWrite 序列。
msgpack EXT 解码:把不可读字节变成 JSON
真正的“可读化”工作在 ext_hook:它复刻了 官方 serde 的 ext_hook 对 EXT code 7 的处理,但把结果映射为纯 JSON 可序列化对象,而非重建 Python 对象:
| EXT code | dump.py 的输出形式 |
|---|---|
7 EXT_DELTA_SNAPSHOT |
递归解码内层 msgpack,包装为 {"__delta_snapshot__": <inner>},随后由 delta_unwrap 拆包 |
0 EXT_CONSTRUCTOR_SINGLE_ARG |
特判 uuid.UUID:把裸 hex 还原为带连字符的标准 UUID 字符串;其它构造直接取参数 |
1 EXT_CONSTRUCTOR_POS_ARGS |
特判 langgraph.types.Send → {"__send__": {"node", "arg"(, "timeout")}} |
| 2 / 3(kwargs / method) | 直接返回第 3 个元素(参数/dict) |
| 4 / 5(Pydantic v1/v2) | 返回模型 dump 出的 dict |
6 EXT_NUMPY_ARRAY |
{"__numpy_array__": {dtype, shape, order, data_b64}},payload 转 base64 |
decode_blob 则按 checkpoint_blobs.type 分派:empty/null/无 blob → None;msgpack → 走上述 ext_hook 解码;bytes/bytearray → base64 字符串;其余类型(包括加密 blob)抛 RuntimeError,明确声明 AES 加密部署不在 v1 支持范围。json_default 兜底处理 JSON 序列化:bytes/bytearray → base64,UUID → 字符串。
输出结构:读懂 delta_kind 与 writes
完整输出示例(继承自 README 的格式约定):
{
"thread_id": "...",
"checkpoint_ns": "",
"target_checkpoint_id": "...",
"parent_checkpoint_id": "...",
"channels": {
"messages": {
"delta_kind": "snapshot",
"seed_checkpoint_id": "...",
"seed_version": "...",
"seed": [{ "type": "ai", "content": "...", "id": "ai-0" }],
"writes": [
{
"checkpoint_id": "...",
"task_id": "...",
"idx": 0,
"value": [{ "type": "ai", "content": "...", "id": "ai-10" }]
}
]
}
}
}
delta_kind 三值的判读:
| 取值 | 含义 | 处置建议 |
|---|---|---|
snapshot |
命中 DeltaChannel 快照 blob(channel_values[ch] == true 的哨兵形态) |
seed 即完整累积值,writes 只需回放 seed 之后的增量 |
legacy_plain |
命中 DeltaChannel 之前的内联或 blob 值 | seed 已是全量值 |
no_seed |
回溯到根都没找到任何持有该渠道值的祖先 | 只能拿到 writes,需自行判断基线(通常为 typ() 空值) |
writes 恒为从旧到新(reducer 回放顺序),每条携带来源 checkpoint_id / task_id / idx,可审计每条写入出自哪次超步。seed_checkpoint_id 与 seed_version 标出 seed 的确切出处,便于与 checkpoint_blobs 表交叉核对。
还原消息列表:seed + writes 归约
对 deepagents 风格的消息列表,把 seed 与 writes 拼接后按 id 去重、剔除 RemoveMessage 墓碑,即可近似还原 reducer 的最终状态(继承自 README 的参考实现):
import json
data = json.load(open("recovery.json"))
ch = data["channels"]["messages"]
messages = list(ch["seed"] or [])
for w in ch["writes"]:
messages.extend(w["value"] or [])
# Dedup by id, keep last; drop RemoveMessage tombstones
by_id = {}
for m in messages:
if isinstance(m, dict) and m.get("type") == "remove":
by_id.pop(m.get("id"), None)
elif isinstance(m, dict) and m.get("id"):
by_id[m["id"]] = m
else:
by_id[id(m)] = m
reduced = list(by_id.values())
这段逻辑“近似 _messages_delta_reducer 语义”并非随意简化,而是与 DeltaChannel 对 reducer 的契约吻合:DeltaChannel 的文档 要求 reducer 必须确定性且对批次划分不敏感(reducer(reducer(state, xs), ys) == reducer(state, xs + ys)),且 replay_writes 对 Overwrite 写入的处理是“取最后一个 Overwrite 作为重置基线”。因此“seed 打底 + writes 顺序追加 + 按 id 保留最后一条”在 add_messages 语义下等价于官方重放结果。如果你的图使用了自定义 reducer,请按该图实际语义调整归约代码。
通过 update_state 写回
README 给出的回写方式(LangGraph Server / langgraph-api 运行时,走 SDK):
from langgraph_sdk import get_client
client = get_client(url="http://localhost:8123")
await client.threads.update_state(
thread_id,
values={"messages": reduced},
)
OSS 部署则使用 graph.update_state(config, {"messages": reduced})(config 指向对应 thread_id / checkpoint_ns)。两点注意:
- 先人工审查恢复出的 JSON 再调用
update_state。update_state会创建一个新 checkpoint,属于变更操作,务必在确认 values 无误后执行; - 该工具只读、有意不修改数据库(README 原话:“This tool is read-only and intentionally does not mutate the database”),全部写回动作由操作者手动完成。
适用前提与范围边界(v1)
README 的 Scope / non-goals 章节界定了本工具的边界,使用时必须对照确认:
- 仅支持 Postgres —— OSS
PostgresSaver或 langgraph-api 的 Postgres 运行时;不支持 inmemory、gRPC core、Mongo、Redis 等其它 checkpointer 后端; - 不做 AES 解密 —— 使用
LANGGRAPH_AES_KEY或自定义加密的部署无法解码(decode_blob遇到未知 blob 类型直接抛错); - 不内置 reducer —— 只输出原始 seed + writes,归约由使用方按自己图的语义完成(见上节参考实现);
- 不自动执行
update_state—— 回写由操作者手动触发。
移植方面:dump.py 完全自包含,可复制到任何机器运行,运行时仅需 psycopg[binary] 与 ormsgpack 两个依赖,没有任何 langgraph 导入,因此不会与目标环境新旧版本运行时的行为相互干扰——这也是把它设计成独立脚本而非库内命令的原因。
适用前提总结:线程由 langgraph >= 1.2(或 deepagents 0.6.x)的 DeltaChannel 格式写入 Postgres,而你即将或已经回滚到无法解析 EXT_DELTA_SNAPSHOT 的旧运行时;此时用本工具把 messages 等渠道导出为 JSON、人工校验后写回,即可让旧运行时读到完整状态。
延伸路径
- 工具本体:examples/delta-channel-dump/README.md、examples/delta-channel-dump/dump.py
- 渠道实现:libs/langgraph/langgraph/channels/delta.py(Beta,快照频率与重放语义)
- 序列化:libs/checkpoint/langgraph/checkpoint/serde/jsonplus.py(EXT 编码表)、libs/checkpoint/langgraph/checkpoint/serde/types.py(
_DeltaSnapshot) - Postgres 侧重建逻辑:libs/checkpoint-postgres/langgraph/checkpoint/postgres/base.py(
get_delta_channel_history的写入历史重建) - 相关验证:libs/checkpoint-conformance 的 DeltaChannelHistory 一致性规格、libs/langgraph 的迁移测试、libs/checkpoint-sqlite 的迁移测试
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0623
Hy4-previewHy4 preview 是由腾讯混元团队研发的新一代混合专家(MoE)旗舰模型。模型总参数量 770B,每个 token 激活 49B,主干共包含78层,第一层采用标准 FFN,其余 77 层均为 MoE 结构,每层包含 256 个路由专家与 1 个共享专家,每个 token 激活 top-8 路由专家及共享专家。主干之外原生内置 1 层 MTP(总参数量 10B,激活 0.7B)以支持投机解码。Python00
GLM-5.3GLM-5.3 与 GLM-5.2 使用相同的基座模型——所有提升均来自后训练。与 GLM-5.2 相比,它在复杂编程和长程任务上的表现显著提升。Jinja00
GLM-5.3-FlashGLM-5.3-Flash (320B-A18B),是GLM-5系列的首个原生多模态模型。320B总参数,能力超过GLM-5.2Jinja00
Spark-X2.5-4BSpark-X2.5-4B 旨在让强大的 AI 更实用、更高效、更易获得。在广泛日常任务中表现强劲,涵盖对话、写作、翻译、推理、编码、工具调用以及智能体工作流,并在同等规模的开源模型中取得领先成绩。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00