首页
/ LangGraph DeltaChannel 检查点数据恢复实战:用 delta-channel-dump 读取 Postgres 原始检查点数据

LangGraph DeltaChannel 检查点数据恢复实战:用 delta-channel-dump 读取 Postgres 原始检查点数据

2026-09-05 13:01:31作者:田桥桑Industrious

当 LangGraph 线程由 langgraph >= 1.2 的 DeltaChannel 格式写入后,直接回滚到旧版运行时时,消息等渠道数据会因旧版 add_messages 无法识别 EXT_DELTA_SNAPSHOT msgpack 扩展编码而静默丢失。本文以仓库中的 delta-channel-dump 工具 为主体,完整讲解其安装、命令行用法、输出格式与回写方式,并结合 libs/langgraphlibs/checkpointlibs/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 位置参构造(如 Sendpathlib.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_nscheckpoint_id 最大的一行作为目标;查不到则直接退出并报错。

第二步:从 target 的 parent 开始回溯

walk_parent_chain 先取出目标 checkpoint 的 parent_checkpoint_id,再对每个请求渠道调用 walk_channel,从父节点沿父链向根方向回溯,沿途做两件事:

  1. 收集 checkpoint_writes 表中该 checkpoint、该渠道的 writes(task_ididxtypeblob),累积到链上;
  2. 检查该 checkpoint 的 checkpoint.channel_values 中是否出现该渠道——一旦出现,该 checkpoint 就是 seed 点,遍历终止。

回溯到达根仍未找到任何祖先持有该渠道的值时,返回 delta_kind: no_seed,但仍会带上途中收集到的全部 writes。

第三步:识别 seed 的类型

_load_seedchannel_values 中该渠道的存储形态分三种情况,正好对应后文的 delta_kind

  • channel_values[ch] is True(布尔哨兵)→ snapshot 类型:用 channel_versions[ch] 作 version 去 checkpoint_blobs 取 blob,解码后若带 __delta_snapshot__ 包裹则解包,得到完整累积值;blob 行不存在或为 empty 时 seed 记为 null。这与 PostgresSaver 历史重建逻辑 的注释一致:对 _DeltaSnapshotchannel_values 里的内联 true 只是标记,不是值。
  • 内联原始值(int/float/str/bool/Nonelegacy_plainput 会把原始类型直接留在 checkpoint 自身的 channel_values 中而不产生 blob 行(见 base.py 的说明),脚本直接内联返回。
  • 其它复杂值 → 同样按 channel_versionscheckpoint_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 → Nonemsgpack → 走上述 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_idseed_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)。两点注意:

  1. 先人工审查恢复出的 JSON 再调用 update_stateupdate_state 会创建一个新 checkpoint,属于变更操作,务必在确认 values 无误后执行;
  2. 该工具只读、有意不修改数据库(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、人工校验后写回,即可让旧运行时读到完整状态。

延伸路径

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.12 K
2.72 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
528
588
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
906
1.83 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
854
1.34 K
docsdocs
暂无描述
Markdown
891
5.78 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.53 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.34 K
1.45 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
987
506
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
540
384