Agent 工程体系 · 第 45/98 篇。内容以 2026 年 9 月可验证的公开规范和稳定接口为基线;框架版本敏感能力会明确标注,不把实验行为写成通用保证。

LangGraph 持久化执行:Thread、Checkpoint、Interrupt、Time Travel 和恢复

LangGraph 的持久化执行,解决的不是“把聊天记录保存下来”这么简单的问题,而是把一次长时间运行的 Agent 执行变成可定位、可暂停、可恢复、可分叉的状态机运行。

一个没有持久化的图,通常可以抽象为:

输入执行节点输出\text{输入} \xrightarrow{\text{执行节点}} \text{输出}

一旦进程崩溃、网络中断或需要人工审批,运行时只能从头开始。

一个带 Checkpointer 的图,则更接近:

S0N1S1N2S2N3S_0 \xrightarrow{N_1} S_1 \xrightarrow{N_2} S_2 \xrightarrow{N_3} \cdots

其中每个 SiS_i 都可以被保存、读取、检查和重新执行。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_anode_b 可能属于同一个执行阶段。只有当这一阶段的状态更新被合并并持久化后,图才进入下一阶段。

这一区分很重要,因为故障可能发生在:

  1. node_anode_b 都执行成功之前;
  2. node_a 成功、node_b 失败之后;
  3. 两个节点都成功,但 Checkpoint 尚未安全写入之前;
  4. 外部副作用已经发生,但节点本身没有成功返回。

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-1thread-2 是两条独立的执行历史。

2.2 Thread 与 Store 的边界

Thread 适合保存:

  • 当前会话的消息;
  • 当前任务的中间状态;
  • 哪些节点已经完成;
  • 哪些节点等待恢复;
  • 某一轮执行的 Checkpoint 历史。

Store 适合保存:

  • 用户长期偏好;
  • 用户画像;
  • 跨会话事实;
  • 多个 Thread 共享的业务知识。

二者的区别可以形式化为:

Thread State=f(thread_id,execution history)\text{Thread State} = f(\text{thread\_id}, \text{execution history})

Store Data=f(namespace,key)\text{Store Data} = f(\text{namespace}, \text{key})

Thread 的生命周期通常与一次会话或一个业务任务绑定;Store 中的数据则可以被多个 Thread 读取。把长期用户画像直接塞进每个 Thread 的 State,会造成重复存储、版本同步和隐私边界混乱。(docs.langchain.com)


3. Checkpoint:状态快照加上执行位置

3.1 Checkpoint 保存什么

一个 Checkpoint 至少需要让运行时回答以下问题:

  1. 当前 State 是什么?
  2. 这是哪条 Thread 的哪个版本?
  3. 接下来应该运行哪些节点?
  4. 当前是否存在等待恢复的任务或 Interrupt?
  5. 这个状态由什么来源产生?
  6. 是否存在同一阶段已经完成但尚未汇总的写入?

从使用者角度,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 不是普通缓存

缓存的典型语义是:

相同输入直接返回旧输出\text{相同输入} \rightarrow \text{直接返回旧输出}

Checkpoint 的语义是:

恢复某个执行位置继续运行后继节点\text{恢复某个执行位置} \rightarrow \text{继续运行后继节点}

因此,从历史 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 持久化执行”生成草稿"
}

这里有三个关键前置条件:

  1. 图必须在 compile() 时传入 Checkpointer;
  2. 调用图时必须传入 thread_id
  3. 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() 时:

  1. 节点执行被挂起;
  2. 当前图状态被 Checkpointer 保存;
  3. 中断载荷返回给调用方;
  4. Thread 等待外部输入;
  5. 调用方使用同一个 thread_idCommand(resume=...) 恢复;
  6. interrupt() 表达式得到恢复值;
  7. 节点继续产生状态更新。

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

因此,必须区分:

节点函数成功返回外部副作用一定只发生一次\text{节点函数成功返回} \neq \text{外部副作用一定只发生一次}

一个安全的外部写操作通常需要:

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,
    }

这里的关键不是“随机数变得确定”,而是:

一次运行中的随机结果持久化任务结果\text{一次运行中的随机结果} \rightarrow \text{持久化任务结果}

恢复时读取该次任务的已保存结果,而不是重新生成随机数。当前时间、随机数、外部查询等非确定性操作都应采用类似思路。官方 Functional API 文档要求将这类操作封装到 Task 中,并强调可序列化输入输出和恢复时的确定性。(docs.langchain.com)


12. 代码升级与旧 Thread 恢复

LangGraph 的一个重要边界是:旧 Thread 恢复时,通常会使用当前部署的图代码,而不是自动绑定到创建该 Thread 时的代码版本。

因此,持久化数据实际上也是一种跨版本 API:

当前代码历史 State 结构\text{当前代码} \leftrightarrow \text{历史 State 结构}

假设旧版本 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_astep_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

关键路径是:

  1. thread_id 选择历史;
  2. Checkpointer 保存当前状态和下一步;
  3. Interrupt 使图停止;
  4. Command(resume=...) 将外部输入送回暂停点;
  5. 恢复时节点可能从头进入;
  6. 外部副作用依靠幂等性;
  7. 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 持久化执行=Thread 定位+Checkpoint 状态保存+可恢复控制流+外部副作用治理\text{LangGraph 持久化执行} = \text{Thread 定位} + \text{Checkpoint 状态保存} + \text{可恢复控制流} + \text{外部副作用治理}

其中前三项由 LangGraph 运行时提供,最后一项必须由应用和外部系统共同完成。只有把这四部分放在同一个执行模型中,Interrupt、恢复、Time Travel 和长时间运行 Agent 才能在生产环境中形成一致而可诊断的行为。


系列导航与关联阅读

官方资料

本文依据 Agent、模型、协议与框架官方资料重新梳理;正文、示例与生产清单由 WR BLOG 编写。