Agent 工程体系 · 第 45/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。
LangGraph 持久化执行:Thread、Checkpoint、Interrupt、Time Travel 和恢复
LangGraph 的持久化执行,解决的不是“把聊天记录保存下来”这么简单的问题,而是把一次长时间运行的 Agent 执行变成可定位、可暂停、可恢复、可分叉的状态机运行。
一个没有持久化的图,通常可以抽象为:
一旦进程崩溃、网络中断或需要人工审批,运行时只能从头开始。
一个带 Checkpointer 的图,则更接近:
其中每个 都可以被保存、读取、检查和重新执行。LangGraph 将这种能力用于线程级短期记忆、故障恢复、人机协作和时间旅行;跨线程的用户偏好、事实和共享知识,则由 Store 负责。(docs.langchain.com)
1. 先建立执行模型:State、Node、Super-step
理解 Thread 和 Checkpoint,必须先明确 LangGraph 执行的对象是什么。
一个图通常由以下部分组成:
- State:节点之间传递的状态;
- Node:读取 State 并产生状态更新的函数;
- Edge:决定后续执行哪些节点;
- Reducer:合并多个节点对同一个 State 字段的更新;
- Checkpointer:在执行过程中保存状态和运行位置;
- Thread:一条可持续推进的图执行历史。
示例状态:
from typing import TypedDict
class State(TypedDict, total=False):
topic: str
draft: str
approved: bool
节点不是直接修改共享对象,而是返回更新:
def generate_topic(state: State):
return {"topic": "如何解释持久化执行"}
def write_draft(state: State):
return {
"draft": f"围绕主题“{state['topic']}”生成一份草稿"
}
如果图是:
START -> generate_topic -> write_draft -> END
那么一次执行可以表示为:
| 执行阶段 | 当前状态 | 下一步 |
|---|---|---|
| 初始 | {} |
generate_topic |
generate_topic 完成 |
{"topic": "如何解释持久化执行"} |
write_draft |
write_draft 完成 |
{"topic": "...", "draft": "..."} |
END |
LangGraph 的检查点通常对应图执行过程中的状态快照,而不是单纯的最终输出。官方文档将其描述为:编译图时配置 Checkpointer 后,运行时会将图状态组织在 Thread 中,并在执行步骤中保存 Checkpoint。(docs.langchain.com)
1.1 什么是 Super-step
LangGraph 的执行受到 Pregel 等计算模型的启发。工程上可以把一次 Super-step 理解为一轮“当前已调度节点执行并提交状态更新”的过程。
如果两个节点从同一个入口并行执行:
-> node_a ->
START END
-> node_b ->
那么 node_a 和 node_b 可能属于同一个执行阶段。只有当这一阶段的状态更新被合并并持久化后,图才进入下一阶段。
这一区分很重要,因为故障可能发生在:
node_a和node_b都执行成功之前;node_a成功、node_b失败之后;- 两个节点都成功,但 Checkpoint 尚未安全写入之前;
- 外部副作用已经发生,但节点本身没有成功返回。
LangGraph 的 Checkpointer 不只保存完整 Checkpoint,也支持保存某个 Super-step 中已经成功产生的 pending writes。节点失败后恢复时,已经成功写入的节点不必再次运行;失败节点仍然需要重新执行。(docs.langchain.com)
2. Thread:持久化执行的逻辑游标
2.1 Thread 不是线程池线程
LangGraph 中的 Thread 不是 Python threading.Thread,也不是数据库连接线程。
Thread 是一条具有稳定标识的图执行历史。调用图时,需要通过配置传入:
config = {
"configurable": {
"thread_id": "article-demo-001"
}
}
这个 thread_id 告诉 Checkpointer:
本次输入应该接在哪一条已有的图状态历史上?
官方文档把 thread_id 称为持久化游标:复用相同的 ID,会加载同一条 Thread 的状态;使用新的 ID,则会创建一条空状态的新 Thread。(docs.langchain.com)
因此,下面两个调用的语义不同:
graph.invoke(
{"topic": "A"},
{"configurable": {"thread_id": "thread-1"}},
)
graph.invoke(
{"topic": "B"},
{"configurable": {"thread_id": "thread-2"}},
)
它们不会共享 Thread 级 State。thread-1 和 thread-2 是两条独立的执行历史。
2.2 Thread 与 Store 的边界
Thread 适合保存:
- 当前会话的消息;
- 当前任务的中间状态;
- 哪些节点已经完成;
- 哪些节点等待恢复;
- 某一轮执行的 Checkpoint 历史。
Store 适合保存:
- 用户长期偏好;
- 用户画像;
- 跨会话事实;
- 多个 Thread 共享的业务知识。
二者的区别可以形式化为:
Thread 的生命周期通常与一次会话或一个业务任务绑定;Store 中的数据则可以被多个 Thread 读取。把长期用户画像直接塞进每个 Thread 的 State,会造成重复存储、版本同步和隐私边界混乱。(docs.langchain.com)
3. Checkpoint:状态快照加上执行位置
3.1 Checkpoint 保存什么
一个 Checkpoint 至少需要让运行时回答以下问题:
- 当前 State 是什么?
- 这是哪条 Thread 的哪个版本?
- 接下来应该运行哪些节点?
- 当前是否存在等待恢复的任务或 Interrupt?
- 这个状态由什么来源产生?
- 是否存在同一阶段已经完成但尚未汇总的写入?
从使用者角度,graph.get_state(config) 返回的状态快照通常包含:
values:当前 State;next:下一批待运行节点;tasks:调度任务及其错误、Interrupt 等信息;config:包含 Thread 和 Checkpoint 标识;metadata:步骤、来源等元数据。
可以使用如下方式检查当前状态:
state = graph.get_state(config)
print("values =", state.values)
print("next =", state.next)
print("config =", state.config)
print("tasks =", state.tasks)
而:
history = list(graph.get_state_history(config))
可以得到该 Thread 的历史 Checkpoint。官方示例指出,历史通常按逆时间顺序返回,最新的 Checkpoint 在前。(docs.langchain.com)
3.2 Checkpoint 不是普通缓存
缓存的典型语义是:
Checkpoint 的语义是:
因此,从历史 Checkpoint 恢复,并不意味着所有后续结果都从缓存中读取。节点之后的计算可能重新发生,包括:
- LLM 调用;
- 外部 API 请求;
- Interrupt;
- 数据库操作;
- 随机或时间相关逻辑。
时间旅行文档明确区分了 Replay 与缓存:Replay 会重新执行 Checkpoint 之后的节点,只有 Checkpoint 之前已保存的结果不会再次执行。(docs.langchain.com)
3.3 Checkpoint 不是完整业务事件日志
Checkpoint 历史具有审计和调试价值,但不能自动等价于业务事件日志。
例如,以下两个事件并不完全相同:
节点返回 {"status": "paid"}
支付服务已经扣款 100 元
前者是图状态更新,后者是外部世界中的业务事实。如果节点在支付请求已经成功后进程崩溃,Checkpoint 可能还没有记录 status="paid"。恢复时再次调用支付接口,就可能重复扣款。
因此:
- Checkpoint 记录的是图的执行状态;
- 业务事件日志记录的是已经提交的领域事实;
- 外部副作用必须通过幂等键、状态查询或 Outbox 等机制保护。
LangGraph 文档要求将外部 API 和副作用封装在 Task 中,并使其具备幂等性,因为任务可能在启动但未成功完成后被重新执行。(docs.langchain.com)
4. 配置 Checkpointer:从内存演示到持久存储
下面给出一个可以独立运行的最小示例。
from typing import TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
class State(TypedDict, total=False):
topic: str
draft: str
def generate_topic(state: State):
return {"topic": "LangGraph 持久化执行"}
def write_draft(state: State):
return {
"draft": f"围绕“{state['topic']}”生成草稿"
}
builder = StateGraph(State)
builder.add_node("generate_topic", generate_topic)
builder.add_node("write_draft", write_draft)
builder.add_edge(START, "generate_topic")
builder.add_edge("generate_topic", "write_draft")
builder.add_edge("write_draft", END)
checkpointer = InMemorySaver()
graph = builder.compile(checkpointer=checkpointer)
config = {
"configurable": {
"thread_id": "demo-thread-1"
}
}
result = graph.invoke({}, config)
print(result)
print(graph.get_state(config).values)
预期输出类似:
{
"topic": "LangGraph 持久化执行",
"draft": "围绕“LangGraph 持久化执行”生成草稿"
}
这里有三个关键前置条件:
- 图必须在
compile()时传入 Checkpointer; - 调用图时必须传入
thread_id; - State 和需要持久化的节点结果必须可序列化。
如果没有 Checkpointer,就无法依赖 Thread 级持久化能力。InMemorySaver 只把数据放在进程内存中,进程重启后数据会丢失,因此适合测试和局部调试,不适合生产恢复。官方文档列出了 SQLite、PostgreSQL 等持久化 Checkpointer 集成。(docs.langchain.com)
生产配置的结构通常类似:
from langgraph.checkpoint.postgres import PostgresSaver
checkpointer = PostgresSaver.from_conn_string(
"postgresql://user:password@localhost:5432/langgraph"
)
checkpointer.setup()
graph = builder.compile(checkpointer=checkpointer)
实际项目应根据所使用的 LangGraph 版本和独立 Checkpointer 包确认导入路径、连接池和异步用法。Checkpointer 的核心接口包括保存 Checkpoint、保存 pending writes、读取单个 Checkpoint 和列出历史 Checkpoint。(docs.langchain.com)
5. Interrupt:在图内部暂停并等待外部输入
5.1 Interrupt 的语义
interrupt() 是动态中断点。它可以放在节点内部,并根据运行时状态决定是否暂停:
from langgraph.types import interrupt
def approval_node(state: State):
approved = interrupt({
"question": "是否允许继续?",
"draft": state["draft"],
})
return {"approved": approved}
执行到 interrupt() 时:
- 节点执行被挂起;
- 当前图状态被 Checkpointer 保存;
- 中断载荷返回给调用方;
- Thread 等待外部输入;
- 调用方使用同一个
thread_id和Command(resume=...)恢复; interrupt()表达式得到恢复值;- 节点继续产生状态更新。
Interrupt 的载荷必须是 JSON 可序列化的值,例如字符串、数字、布尔值、字典或数组。LangGraph 文档明确要求 Interrupt 依赖 Checkpointer 和 Thread ID。(docs.langchain.com)
5.2 完整的审批示例
from typing import Literal, TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
from langgraph.types import Command, interrupt
class ApprovalState(TypedDict, total=False):
action: str
approved: bool
status: str
def approval(state: ApprovalState):
decision = interrupt({
"type": "approval",
"question": "是否执行该操作?",
"action": state["action"],
})
return {"approved": bool(decision)}
def route(state: ApprovalState) -> Command[Literal["proceed", "cancel"]]:
if state["approved"]:
return Command(goto="proceed")
return Command(goto="cancel")
def proceed(state: ApprovalState):
return {"status": "completed"}
def cancel(state: ApprovalState):
return {"status": "cancelled"}
builder = StateGraph(ApprovalState)
builder.add_node("approval", approval)
builder.add_node("route", route)
builder.add_node("proceed", proceed)
builder.add_node("cancel", cancel)
builder.add_edge(START, "approval")
builder.add_edge("approval", "route")
builder.add_edge("proceed", END)
builder.add_edge("cancel", END)
graph = builder.compile(checkpointer=InMemorySaver())
config = {
"configurable": {
"thread_id": "approval-thread-1"
}
}
paused = graph.invoke(
{
"action": "发送一封外部邮件",
"status": "waiting",
},
config,
)
print(paused.get("__interrupt__"))
首次调用的结果会包含类似:
(
Interrupt(
value={
"type": "approval",
"question": "是否执行该操作?",
"action": "发送一封外部邮件",
}
),
)
此时不能使用新的 Thread ID 恢复:
# 错误语义:这会启动一条新 Thread,而不是回答原来的 Interrupt
graph.invoke(
Command(resume=True),
{"configurable": {"thread_id": "another-thread"}},
)
正确方式是复用原配置:
resumed = graph.invoke(
Command(resume=True),
config,
)
print(resumed)
恢复值 True 会成为 approval() 中 interrupt(...) 的返回值,随后状态变为:
approval
-> interrupt
-> resume=True
-> approved=True
-> route
-> proceed
-> END
Command(resume=...) 是作为外部输入恢复 Interrupt 的用法;而 Command(goto=...) 通常由节点返回,用于在节点内部改变后续路由。不要把两者混为一谈。(docs.langchain.com)
6. 恢复时节点会从哪里开始?
这是 Interrupt 最容易被误解的地方。
恢复并不是把 Python 函数的栈帧冻结后继续执行。节点会从包含 Interrupt 的节点开头重新进入,然后在运行到对应的 interrupt() 时取出恢复值。
例如:
def review_node(state: State):
print("before")
value = interrupt("请确认")
print("after", value)
return {"approved": value}
执行过程是:
第一次运行:
print("before")
interrupt("请确认")
节点暂停
恢复运行:
print("before") # 再次执行
interrupt(...) # 得到 resume 值
print("after", True)
返回状态更新
所以 Interrupt 之前的代码必须满足以下条件:
- 不应重复产生不可逆副作用;
- 如果有副作用,应当幂等;
- 随机数、当前时间、外部查询结果等非确定性数据,应被封装为可持久化的 Task 或在 State 中显式保存;
- 多个 Interrupt 在同一个节点中的顺序不能随意改变。
官方 Interrupt 规则明确指出:不要在 interrupt() 外层捕获其控制流异常,不要重排同一节点中的 Interrupt 调用,Interrupt 之前的副作用必须具备幂等性。(docs.langchain.com)
6.1 错误示例:Interrupt 前发送邮件
def bad_node(state: State):
send_email("user@example.com", "开始处理")
approved = interrupt("是否继续?")
return {"approved": approved}
恢复后,send_email() 可能再次执行。
更安全的结构是:
def notify_task():
# 通过任务封装,并使用业务幂等键
send_email_once(
idempotency_key="task-123-start-notification",
to="user@example.com",
subject="开始处理",
)
def safer_node(state: State):
notify_task()
approved = interrupt("是否继续?")
return {"approved": approved}
仅仅把代码放进 Task 并不自动使外部系统幂等。Task 负责让运行时能够记录和恢复任务结果;业务系统仍然需要使用幂等键、唯一约束或状态检查防止重复提交。LangGraph 文档同时强调了 Task 封装与幂等设计这两个要求。(docs.langchain.com)
7. 多个 Interrupt:为什么需要 ID 到恢复值的映射
当多个并行节点同时中断时,恢复值不能只依赖列表位置。
示例:
from typing import Annotated, TypedDict
import operator
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
from langgraph.types import Command, interrupt
class ParallelState(TypedDict):
results: Annotated[list[str], operator.add]
def node_a(state: ParallelState):
answer = interrupt("问题 A")
return {"results": [f"A={answer}"]}
def node_b(state: ParallelState):
answer = interrupt("问题 B")
return {"results": [f"B={answer}"]}
builder = StateGraph(ParallelState)
builder.add_node("a", node_a)
builder.add_node("b", node_b)
builder.add_edge(START, "a")
builder.add_edge(START, "b")
builder.add_edge("a", END)
builder.add_edge("b", END)
graph = builder.compile(checkpointer=InMemorySaver())
config = {"configurable": {"thread_id": "parallel-1"}}
paused = graph.invoke({"results": []}, config)
interrupts = paused["__interrupt__"]
for item in interrupts:
print(item.id, item.value)
中断集合可能类似:
id=abc, value=问题 A
id=xyz, value=问题 B
恢复时应建立映射:
resume_map = {
item.id: f"回答:{item.value}"
for item in interrupts
}
result = graph.invoke(
Command(resume=resume_map),
config,
)
print(result["results"])
结果可能是:
[
"A=回答:问题 A",
"B=回答:问题 B",
]
映射的意义是将每个外部回答绑定到具体的 Interrupt,而不是假设运行时会永远以固定列表顺序返回它们。并行分支同时中断时,官方文档要求使用 Interrupt ID 到恢复值的映射。(docs.langchain.com)
8. Time Travel:Replay 与 Fork
LangGraph 的 Time Travel 是围绕历史 Checkpoint 进行的两类操作:
- Replay:从过去的 Checkpoint 重新执行;
- Fork:从过去的 Checkpoint 修改状态,创建新的后续轨迹。
二者都不会修改原来的历史节点,而是在已有历史上继续创建新的执行分支。(docs.langchain.com)
8.1 Replay:重跑历史之后的节点
先构建一个简单图:
from typing import TypedDict
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
class JokeState(TypedDict, total=False):
topic: str
joke: str
def generate_topic(state: JokeState):
return {"topic": "洗衣机里的袜子"}
def write_joke(state: JokeState):
return {
"joke": f"为什么{state['topic']}总会消失?因为它们私奔了。"
}
builder = StateGraph(JokeState)
builder.add_node("generate_topic", generate_topic)
builder.add_node("write_joke", write_joke)
builder.add_edge(START, "generate_topic")
builder.add_edge("generate_topic", "write_joke")
builder.add_edge("write_joke", END)
graph = builder.compile(checkpointer=InMemorySaver())
config = {
"configurable": {
"thread_id": "joke-1"
}
}
graph.invoke({}, config)
history = list(graph.get_state_history(config))
for snapshot in history:
print(
"next=",
snapshot.next,
"checkpoint_id=",
snapshot.config["configurable"].get("checkpoint_id"),
)
可以选择 next == ("write_joke",) 的 Checkpoint:
before_joke = next(
snapshot
for snapshot in history
if snapshot.next == ("write_joke",)
)
然后重放:
replayed = graph.invoke(None, before_joke.config)
print(replayed)
这里的逻辑是:
已有历史:
generate_topic -> Checkpoint C1
write_joke -> Checkpoint C2
从 C1 Replay:
generate_topic 不再执行
write_joke 重新执行
如果 write_joke 是 LLM 节点,那么 Replay 可能得到不同结果;如果它调用了外部 API,也可能再次调用。Replay 适合调试和验证,但不应被误认为“读取原结果”。(docs.langchain.com)
8.2 Fork:修改过去状态并探索新路径
Fork 的核心是 update_state():
fork_config = graph.update_state(
before_joke.config,
{
"topic": "代码审查里的隐藏 Bug",
},
)
forked = graph.invoke(None, fork_config)
print(forked)
这里不会把原来的 before_joke 修改成新内容,而是:
原分支:
C1(topic="洗衣机里的袜子")
-> C2(joke="...")
新分支:
C1(topic="洗衣机里的袜子")
-> C1' (topic="代码审查里的隐藏 Bug")
-> C2' (重新生成 joke)
update_state() 产生的是新的 Checkpoint。状态更新仍会遵循 State 字段配置的 Reducer;对于带 Reducer 的字段,更新可能是追加或合并,而不是简单覆盖。(docs.langchain.com)
8.3 as_node:决定 Fork 从哪里继续
如果运行时无法判断这次状态更新由哪个节点产生,或者你希望显式指定恢复位置,可以使用 as_node:
fork_config = graph.update_state(
before_joke.config,
{"topic": "新的主题"},
as_node="generate_topic",
)
这表示:
把这次状态更新视为
generate_topic产生的结果,因此接下来运行它的后继节点。
它在以下场景尤其重要:
- 并行分支同时写入状态;
- 新 Thread 没有执行历史;
- 需要跳过某个节点;
- 希望将人工编辑后的状态接到指定节点之后。
如果 as_node 设置错误,可能导致:
- 节点被重复执行;
- 节点被错误跳过;
- 路由进入意外分支;
- 出现状态更新冲突。
因此,as_node 不是普通的“修改状态参数”,而是同时修改了图的执行语义。(docs.langchain.com)
9. Interrupt 与 Time Travel 的关系
Interrupt 是运行中的暂停;Time Travel 是对历史 Checkpoint 的定位和重新启动。
可以把二者放在同一条时间轴上:
C0: 初始状态
|
v
C1: 节点 A 完成
|
v
C2: 节点 B 调用 interrupt()
| \
| \ resume=True
| v
| C3: 节点 B 完成
|
\ fork with approved=False
v
C2': 修改后的分支
从 C2 恢复属于正常 Resume:
graph.invoke(Command(resume=True), config)
从 C1 重新执行属于 Replay:
graph.invoke(None, c1_config)
从 C2 修改状态并执行属于 Fork:
fork_config = graph.update_state(
c2_config,
{"approved": False},
)
graph.invoke(None, fork_config)
三种操作的区别:
| 操作 | 是否修改原历史 | 是否重新执行后续节点 | 是否注入新状态 |
|---|---|---|---|
| Resume | 否 | 是 | 以 Interrupt 恢复值注入 |
| Replay | 否 | 是 | 通常不修改 State |
| Fork | 否 | 是 | 通过 update_state() 修改 |
Replay 和 Fork 都可能重新触发 LLM、API 或 Interrupt,因此它们不是无副作用的只读查询。(docs.langchain.com)
10. 故障恢复:恢复的是图状态,不是外部世界
10.1 节点失败的基本路径
假设图为:
node_a -> node_b -> node_c
执行过程:
1. node_a 成功
2. 保存 Checkpoint C1
3. node_b 开始
4. node_b 抛出异常
5. Thread 进入 error 状态
6. 运维或调用方重新执行 Thread
7. 从 C1 继续 node_b
如果 node_b 在外部系统中已经完成操作,但在返回 LangGraph 前崩溃,那么恢复时可能再次执行 node_b。
因此,必须区分:
一个安全的外部写操作通常需要:
def charge_payment(order_id: str, amount: int):
key = f"charge:{order_id}"
existing = payment_service.query_by_idempotency_key(key)
if existing is not None:
return existing
return payment_service.charge(
order_id=order_id,
amount=amount,
idempotency_key=key,
)
节点中只保存可恢复的业务结果:
def payment_node(state):
payment = charge_payment(
order_id=state["order_id"],
amount=state["amount"],
)
return {
"payment_id": payment["id"],
"payment_status": payment["status"],
}
LangGraph 提供的是持久化执行框架,不会替外部支付、邮件、库存或数据库系统自动提供 exactly-once 语义。对于外部副作用,通常只能通过幂等设计实现 effectively-once。
10.2 Pending writes 的作用
并行执行时,假设:
node_a 成功
node_b 失败
如果 node_a 的结果已经作为 pending write 被保存,那么恢复时可以复用 node_a 的结果,只重跑 node_b。这避免了无意义的重复计算,也减少了并行节点中的重复副作用。(docs.langchain.com)
但这并不消除所有风险。假如 node_a 已经调用外部系统成功,而它的 pending write 保存过程也失败,恢复时仍然可能再次调用外部系统。因此:
- Pending writes 减少图内部的重复执行;
- 幂等键解决外部系统的重复提交;
- 业务事件日志确认外部事实;
- Checkpoint 连接图的执行位置。
这四个机制解决的是不同问题,不能互相替代。
11. 确定性:恢复为什么要求稳定的控制流
设节点中存在随机控制流:
import random
def unstable_node(state):
branch = random.choice(["a", "b"])
if branch == "a":
value = interrupt("分支 A")
else:
value = interrupt("分支 B")
return {"value": value}
第一次运行可能得到:
随机结果:A
暂停点:interrupt("分支 A")
恢复时,如果节点重新随机得到 B,运行时面对的 Interrupt 序列就发生了变化。这会破坏“恢复同一个执行”的前提。
更可靠的方式是把随机结果先记录为可恢复结果:
from langgraph.func import task
@task
def choose_branch():
return random.choice(["a", "b"])
def stable_node(state):
branch = choose_branch().result()
if branch == "a":
value = interrupt("分支 A")
else:
value = interrupt("分支 B")
return {
"branch": branch,
"value": value,
}
这里的关键不是“随机数变得确定”,而是:
恢复时读取该次任务的已保存结果,而不是重新生成随机数。当前时间、随机数、外部查询等非确定性操作都应采用类似思路。官方 Functional API 文档要求将这类操作封装到 Task 中,并强调可序列化输入输出和恢复时的确定性。(docs.langchain.com)
12. 代码升级与旧 Thread 恢复
LangGraph 的一个重要边界是:旧 Thread 恢复时,通常会使用当前部署的图代码,而不是自动绑定到创建该 Thread 时的代码版本。
因此,持久化数据实际上也是一种跨版本 API:
假设旧版本 State 是:
class State(TypedDict):
draft: str
新版本改成:
class State(TypedDict):
content: str
那么停在旧字段 draft 上的 Thread,恢复时就可能无法继续。
升级时应考虑:
- 保留旧字段一段时间;
- 节点兼容旧字段和新字段;
- 对 State 做显式迁移;
- 不要直接删除仍可能被活跃 Thread 使用的节点;
- 在预发布环境使用
get_state()和 Time Travel 检查旧 Thread; - 记录当前图版本与 State schema 版本。
官方向后兼容文档明确指出,LangGraph 会让现有 Thread 立即使用最新图代码,因此持久化数据与图代码之间的契约必须由应用负责维护。(docs.langchain.com)
一个兼容读取方式:
def write_draft(state):
content = state.get("content") or state.get("draft")
if not content:
raise ValueError("missing content/draft")
return {
"content": content.strip()
}
这不是永久方案,但可以为历史 Thread 提供迁移窗口。
13. Subgraph 的 Checkpoint 粒度
默认情况下,子图可以继承父图的 Checkpointer。此时,父图可能把整个子图视为一个较大的执行单元:
父图:
START -> subgraph_node -> END
子图:
START -> step_a -> step_b -> END
如果子图没有自己的 Checkpointer,父图层面通常只能从 subgraph_node 之前或之后进行时间旅行,不能细化到 step_a 和 step_b 之间。
如果需要在子图内部进行恢复或时间旅行,可以为子图配置自己的 Checkpointer:
subgraph = (
StateGraph(SubState)
.add_node("step_a", step_a)
.add_node("step_b", step_b)
.add_edge(START, "step_a")
.add_edge("step_a", "step_b")
.compile(checkpointer=True)
)
这会增加 Checkpoint 数量和持久化开销,但可以获得更细粒度的 Interrupt、Replay 和 Fork 能力。官方 Time Travel 文档将“继承父 Checkpointer”和“子图拥有独立 Checkpointer”视为两种不同的恢复粒度。(docs.langchain.com)
14. 生产诊断:先判断 Thread 处于什么状态
当一次 Agent “卡住”时,不能直接假设是网络问题。至少需要区分:
- Thread 尚未开始;
- Thread 正在执行;
- Thread 已正常结束;
- Thread 正在等待 Interrupt;
- Thread 因节点异常进入错误状态;
- Thread 已有历史,但调用方错误地使用了新的
thread_id; - Thread 已从旧代码恢复,而 State schema 不兼容。
对已知 Thread,可以先检查:
state = graph.get_state(config)
print("values:", state.values)
print("next:", state.next)
print("tasks:", state.tasks)
print("metadata:", state.metadata)
如果希望检查完整轨迹:
history = list(graph.get_state_history(config))
for snapshot in history:
print({
"checkpoint_id": snapshot.config["configurable"].get("checkpoint_id"),
"next": snapshot.next,
"metadata": snapshot.metadata,
})
查找中断点:
interrupted = next(
snapshot
for snapshot in history
if snapshot.tasks
and any(task.interrupts for task in snapshot.tasks)
)
常见错误表现及其原因如下:
| 表现 | 常见原因 |
|---|---|
| 恢复后像新任务一样从头开始 | 使用了新的 thread_id |
| 读取不到历史 | 使用了内存 Checkpointer 且进程已重启 |
| 恢复后副作用重复 | Interrupt 前代码或失败节点不幂等 |
| Replay 结果变化 | 后续节点含 LLM、随机数或外部 API |
| 状态字段越来越大 | 长 Thread 持续保存完整状态 |
| Fork 后节点走错 | as_node 或 Reducer 语义不符合预期 |
| 子图无法从内部节点恢复 | 子图没有独立 Checkpointer |
长会话还会带来 Checkpoint 增长问题。官方文档建议配置保留策略、定期清理旧 Checkpoint;部分增量存储能力可以减少 append-heavy State 的存储量,但应根据当前版本确认其稳定性和 API 状态。(docs.langchain.com)
15. 一个完整的持久化执行时序
把上述机制合并起来,一次需要人工审批、可能失败并支持分叉的 Agent,大致经历以下过程:
sequenceDiagram
participant C as Caller
participant G as LangGraph
participant P as Checkpointer
participant H as Human
participant X as External Service
C->>G: invoke(input, thread_id)
G->>P: load latest checkpoint
G->>G: execute deterministic nodes
G->>P: save checkpoint
G->>G: execute approval node
G->>P: save interrupted state
G-->>C: __interrupt__ payload
C->>H: display approval request
H-->>C: approve/reject
C->>G: invoke(Command(resume=value), same thread_id)
G->>P: load interrupted checkpoint
G->>G: re-enter interrupted node
G->>G: obtain resume value
G->>X: execute idempotent side effect
X-->>G: result
G->>P: save completed checkpoint
G-->>C: final state
C->>G: invoke(None, old checkpoint config)
G->>P: load historical checkpoint
G->>G: replay later nodes
G->>P: save replay branch
关键路径是:
thread_id选择历史;- Checkpointer 保存当前状态和下一步;
- Interrupt 使图停止;
Command(resume=...)将外部输入送回暂停点;- 恢复时节点可能从头进入;
- 外部副作用依靠幂等性;
- Replay 和 Fork 从历史 Checkpoint 产生新的轨迹。
16. 最终边界:LangGraph 保证什么,不保证什么
在正确配置 Checkpointer、Thread ID 和可序列化 State 的前提下,LangGraph 提供了这些运行时能力:
- 将图状态组织到 Thread 中;
- 保存执行过程中的 Checkpoint;
- 从历史状态恢复;
- 在节点内部动态 Interrupt;
- 通过
Command(resume=...)继续执行; - 通过 Replay 重跑历史之后的节点;
- 通过 Fork 修改过去状态并探索替代路径;
- 在部分节点失败时利用已保存的 pending writes 恢复。
但以下语义不能自动获得:
- 外部 API 的 exactly-once;
- 邮件、支付、库存等副作用的自动幂等;
- 旧图代码与新 State schema 的自动迁移;
- Checkpoint 历史天然等价于领域事件日志;
- Replay 对 LLM 和外部 API 结果的稳定复现;
- 进程内存 Checkpointer 跨重启持久化。
可以用一句更精确的话概括:
其中前三项由 LangGraph 运行时提供,最后一项必须由应用和外部系统共同完成。只有把这四部分放在同一个执行模型中,Interrupt、恢复、Time Travel 和长时间运行 Agent 才能在生产环境中形成一致而可诊断的行为。
系列导航与关联阅读
- 系列入口:Agent 工程完整路线:从运行循环、记忆与协议到安全、评测和生产交付
- 上一篇:LangGraph 完整基础:State、Node、Edge、Command、Checkpoint 和中断
- 下一篇:LangChain Agent:Model、Tool、Middleware、State 与适用边界
- 延伸:Agent 持久化执行:事件日志、Checkpoint、租约、恢复和确定性
官方资料
本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论