上一篇介绍了 Interrupt。它处理的是主动暂停:Graph 运行到指定位置,等待人工审批、补充信息或修改数据,再恢复执行。

实际运行过程中可能会出现另外一类情况导致 Agent 停止,如模型请求超时、外部 API 临时不可用、某个并行任务失败,甚至服务进程直接退出。

LangGraph 的持久执行(Durable Execution)主要处理这类问题。它依赖 Checkpoint 保存执行进度,并配合 Retry、Pending Writes 和 Time Travel,让工作流能够重试失败节点、从保存的位置恢复,以及基于历史状态重新执行。

一、故障与恢复边界

以一个订单处理流程为例:

代码片段Text
读取订单

查询库存

调用支付服务

生成确认信息

完成

如果支付服务偶尔返回 503,再调用一次可能就能成功。这类故障通常持续时间很短,适合在当前 Node 内重试。

如果流程在某一个节点执行后(如“查询库存”)退出,此时需要保留此前完成的进度,服务恢复以后继续执行后面的步骤。

还有一种情况发生在业务逻辑本身。某次运行已经结束,但结束后才发现中间状态有误,希望回到某个历史节点重新运行,或者修改当时的 State,再观察后续路径。这属于 Time Travel 的使用范围。

这三类问题对应三个不同处理机制:

情况主要机制处理范围
临时网络错误、服务抖动Retry当前 Node
运行中断、进程退出CheckpointGraph 执行进度
历史状态重新执行Time Travel已保存的 Checkpoint

二、持久执行基础

Durable Execution 的基础仍然是前文介绍过的 Checkpointer。

Graph 编译时传入 Checkpointer:

代码片段Python
from langgraph.checkpoint.memory import InMemorySaver
 
checkpointer = InMemorySaver()
 
graph = builder.compile(
    checkpointer=checkpointer
)

InMemorySaver 是内存版 Checkpointer,适合示例和测试。它把 Checkpoint 保存在当前进程内存中,进程退出后数据也会丢失。

如果需要在服务重启后继续恢复,应使用持久化实现,例如 SQLite 或 PostgreSQL。LangGraph 也提供了 SQLite、PostgreSQL 等 Checkpointer 独立包。

使用 Checkpointer 后,调用 Graph 时还需要提供 thread_id

代码片段Python
config = {
    "configurable": {
        "thread_id": "order-1001"
    }
}

thread_id 用来标识这一条执行历史。同一个 Thread 下产生的 Checkpoint 会被关联起来,LangGraph 因此能够找到当前状态和之前的运行记录。

Checkpoint 保存在哪里

对于一个顺序 Graph:

代码片段Text
START → A → B → C → END

可以简化为:

代码片段Text
输入
 ↓ Checkpoint
A
 ↓ Checkpoint
B
 ↓ Checkpoint
C
 ↓ Checkpoint
结束

LangGraph 的执行以 Super-step 为基本推进单位。Super-step 是 Graph 的一次执行步。在顺序 Graph 中,一个 Super-step 通常只有一个 Node;存在并行分支时,同一个 Super-step 中可以同时运行多个 Node。

LangGraph 会在 Super-step 边界保存完整的状态快照。这个快照由 StateSnapshot 表示,其中除了 State,还包含下一步待执行节点等运行信息。发生故障后,Runtime 可以依据这些信息判断后续应该执行什么。

这里有一个很重要的边界:恢复发生在 Checkpoint 对应的执行位置,并不能恢复到 Node 函数内部某一行代码。关于这一点,在前一篇中,也强调过。

因此,一个 Node 中包含多少任务,会直接影响失败后需要重复执行多少内容。

三、节点重试机制

临时性故障通常没有必要等待整个 Graph 重新恢复。

LangGraph 可以通过 RetryPolicy 为 Node 设置重试策略。RetryPolicy 是一组重试参数,通常在 add_node() 时通过 retry_policy 传入:

代码片段Python
from langgraph.types import RetryPolicy
 
builder.add_node(
    "call_api",
    call_api,
    retry_policy=RetryPolicy(
        max_attempts=3
    ),
)

其中:

  • max_attempts 表示最多尝试多少次,包含第一次执行。
  • initial_interval 表示第一次重试前等待的时间。
  • backoff_factor 控制后续重试间隔的增长。
  • retry_on 用来指定哪些异常允许重试。

当前默认 max_attempts 为 3,并使用带退避和随机抖动的等待策略。retry_on 也可以直接传入异常类型或自定义判断函数。

下面用一个确定性的例子观察重试过程:

代码片段Python
from typing_extensions import TypedDict
 
from langgraph.graph import StateGraph, START, END
from langgraph.types import RetryPolicy
 
 
class State(TypedDict):
    result: str
 
 
attempts = 0
 
 
def call_service(state: State):
    global attempts
 
    attempts += 1
    print(f"attempt={attempts}")
 
    if attempts < 3:
        raise ConnectionError("service unavailable")
 
    return {"result": "service ok"}
 
 
builder = StateGraph(State)
 
builder.add_node(
    "call_service",
    call_service,
    retry_policy=RetryPolicy(
        max_attempts=3,
        retry_on=ConnectionError,
        initial_interval=0.1,
        backoff_factor=1.0,
        jitter=False,
    ),
)
 
builder.add_edge(START, "call_service")
builder.add_edge("call_service", END)
 
graph = builder.compile()
 
result = graph.invoke({"result": ""})
print(result)

运行结果

代码片段Text
attempt=1
attempt=2
attempt=3
{'result': 'service ok'}

前两次执行抛出 ConnectionError,符合 retry_on 的条件,所以当前 Node 再次运行。第三次成功后,Graph 才进入后续步骤。

实际项目中,重试条件最好与错误类型对应。

网络超时、服务暂时不可用等异常通常适合重试。如果是参数格式错误、业务规则不满足、资源确实不存在,重复执行则往往没有意义。否则 Retry 只会把一次明确的失败延长成多次相同的失败。

四、执行恢复过程

Retry 有次数限制。如果外部服务持续不可用,Node 最终仍然可能失败。

有了 Checkpointer,已经完成的步骤可以保留下来。当服务恢复之后,可再从保存的执行位置继续。

下面把过程简化成三个节点:

代码片段Python
from typing_extensions import TypedDict
 
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.graph import StateGraph, START, END
 
 
class OrderState(TypedDict, total=False):
    order_id: str
    prepared: bool
    submitted: bool
    status: str
 
 
service_ready = False
 
 
def prepare(state: OrderState):
    print("prepare")
    return {"prepared": True}
 
 
def submit(state: OrderState):
    print("submit")
 
    if not service_ready:
        raise ConnectionError("remote service unavailable")
 
    return {"submitted": True}
 
 
def finalize(state: OrderState):
    print("finalize")
    return {"status": "done"}
 
 
builder = StateGraph(OrderState)
 
builder.add_node("prepare", prepare)
builder.add_node("submit", submit)
builder.add_node("finalize", finalize)
 
builder.add_edge(START, "prepare")
builder.add_edge("prepare", "submit")
builder.add_edge("submit", "finalize")
builder.add_edge("finalize", END)
 
graph = builder.compile(
    checkpointer=InMemorySaver()
)
 
config = {
    "configurable": {
        "thread_id": "order-1001"
    }
}
 
try:
    graph.invoke(
        {"order_id": "ORDER-1001"},
        config,
    )
except ConnectionError:
    print("run failed")
 
service_ready = True
 
result = graph.invoke(
    None,
    config,
)
 
print(result)

这里第二次调用的 None 有明确含义:不提供新的 Graph 输入,而是继续 config 所指向 Thread 中尚未完成的执行。

运行结果如下:

代码片段Text
prepare
submit
run failed
submit
finalize
{
    'order_id': 'ORDER-1001',
    'prepared': True,
    'submitted': True,
    'status': 'done'
}

prepare 只执行了一次。

第一次运行已经完成这个 Node,并保存了相应进度。submit 失败后,第二次执行继续处理尚未完成的部分,所以又执行了 submit,随后进入 finalize

示例使用 InMemorySaver,因此只能展示同一进程内的恢复。服务真正退出并重新启动时,需要 SQLite、PostgreSQL 等持久化 Checkpointer。

五、并行执行恢复

顺序 Graph 中,每次通常只有一个 Node 在运行。并行 Graph 稍微复杂一些。

例如两个节点处于同一个 Super-step:

代码片段Text
        ┌─ fetch_profile ─ 成功
START ──┤
        └─ fetch_orders  ─ 失败

此时整个 Super-step 还没有完成,但 fetch_profile 已经产生了有效结果。

LangGraph 会保存这种已经完成的节点写入,这些记录称为 Pending Writes。恢复执行时,成功节点的结果可以继续使用,不必因为另一个并行节点失败而全部重算。

因此,恢复后的执行更接近:

代码片段Text
fetch_profile    使用已有结果
fetch_orders     重新执行

这类机制对并行检索、多数据源查询、多个 Agent 同时执行任务尤其重要。某一路调用失败时,已经完成的模型请求或远程查询不必全部重复。

Checkpoint 和 Pending Writes 的作用范围也有区别:

代码片段Text
Checkpoint
    Super-step 边界上的完整状态

Pending Writes
    当前 Super-step 中已完成节点的写入

Time Travel 选择的仍然是完整 Checkpoint;Pending Writes 主要用于失败后的恢复。

六、Time Travel 机制

Checkpoint 保存了历史状态,也为 Time Travel 提供了基础。

LangGraph 的 Time Travel 主要包含两种操作:

  • Replay:从历史 Checkpoint 重新执行后续节点。
  • Fork:修改历史状态,再从该位置形成新的执行分支。

Replay

get_state_history(config) 可以取得某个 Thread 的历史 StateSnapshot

每个 StateSnapshot 中的 next 表示从这个 Checkpoint 继续执行时,下一步计划运行哪些 Node。历史记录默认按时间倒序返回。

例如:

代码片段Python
history = list(
    graph.get_state_history(config)
)
 
before_finalize = next(
    state
    for state in history
    if state.next == ("finalize",)
)

这里找到的是 finalize 执行前的 Checkpoint。

随后:

代码片段Python
result = graph.invoke(
    None,
    before_finalize.config,
)

Graph 会从这个历史位置继续执行。

Checkpoint 之前已经完成的 Node 不会重新运行;后面的 Node 会再次执行。因此其中的模型调用、HTTP 请求和 Interrupt 也会重新执行,而新的执行结果也并不保证与原来运行结果完全相同。

Fork

Fork 需要使用 update_state()

它接收一个历史 Checkpoint 的配置以及新的 State 值:

代码片段Python
fork_config = graph.update_state(
    before_finalize.config,
    values={"risk": "high"},
)

随后继续运行:

代码片段Python
result = graph.invoke(
    None,
    fork_config,
)

update_state() 会创建一个新的 Checkpoint 分支,原来的历史记录仍然保留。后续节点读取修改后的 State,因此可能得到不同的路径或结果。

这种能力很适合调试 Agent。

例如前面的检索和模型调用已经花费几十秒,只想验证“如果风险等级改成 high,后面的决策会怎样”,就可以直接从相应 Checkpoint 建立分支,不必反复运行整条 Graph。

七、外部操作与幂等

Checkpoint 能够保存执行状态,但已经写入外部系统的数据仍然需要应用自身保证一致性。

例如一个支付 Node:

代码片段Text
调用支付接口

支付成功

进程中断

Node 尚未完整结束

恢复以后,这个 Node 可能再次执行。如果第二次调用又产生一笔支付,就会造成重复扣款。

因此,发送邮件、创建订单、扣款、写入外部数据库等操作,都应考虑幂等(Idempotency)。

幂等表示同一个业务操作重复执行多次,最终结果与执行一次保持一致。常见做法包括使用幂等键、业务唯一键,或者在执行前查询对应操作是否已经完成。

Node 的粒度也会影响恢复成本。Node 如何划分,会直接影响失败以后需要重复多少工作。

如果一个 Node 内连续完成查询数据库、调用三个 API、调用一次模型,最后一步失败,那么重新进入这个 Node 时,其中一部分操作可能再次执行,正如前面一再强调的。此时,分成多个 Node 通常会更合适。

总结

LangGraph 的容错由几个不同层面的机制共同完成。

RetryPolicy 处理 Node 内部的临时故障;Checkpoint 保存 Graph 的执行状态;Pending Writes 记录并行 Super-step 中已经完成的工作;Time Travel 则利用历史 Checkpoint 完成 Replay 和 Fork。

实际设计时,可以按故障的范围选择处理方式:

代码片段Text
临时错误
    → Retry

运行中断
    → Checkpoint 恢复

并行节点部分失败
    → Pending Writes

历史状态重新执行
    → Replay / Fork

外部系统写入
    → 幂等控制

Durable Execution 的价值也体现在这里:一个 Agent 的执行不再依赖某一次进程必须完整运行到底。只要状态、执行边界和外部系统写入设计清楚,长时间任务就有了失败后继续恢复的基础。