上一篇介绍了 Interrupt。它处理的是主动暂停:Graph 运行到指定位置,等待人工审批、补充信息或修改数据,再恢复执行。
实际运行过程中可能会出现另外一类情况导致 Agent 停止,如模型请求超时、外部 API 临时不可用、某个并行任务失败,甚至服务进程直接退出。
LangGraph 的持久执行(Durable Execution)主要处理这类问题。它依赖 Checkpoint 保存执行进度,并配合 Retry、Pending Writes 和 Time Travel,让工作流能够重试失败节点、从保存的位置恢复,以及基于历史状态重新执行。
一、故障与恢复边界
以一个订单处理流程为例:
读取订单
↓
查询库存
↓
调用支付服务
↓
生成确认信息
↓
完成
如果支付服务偶尔返回 503,再调用一次可能就能成功。这类故障通常持续时间很短,适合在当前 Node 内重试。
如果流程在某一个节点执行后(如“查询库存”)退出,此时需要保留此前完成的进度,服务恢复以后继续执行后面的步骤。
还有一种情况发生在业务逻辑本身。某次运行已经结束,但结束后才发现中间状态有误,希望回到某个历史节点重新运行,或者修改当时的 State,再观察后续路径。这属于 Time Travel 的使用范围。
这三类问题对应三个不同处理机制:
| 情况 | 主要机制 | 处理范围 |
|---|---|---|
| 临时网络错误、服务抖动 | Retry | 当前 Node |
| 运行中断、进程退出 | Checkpoint | Graph 执行进度 |
| 历史状态重新执行 | Time Travel | 已保存的 Checkpoint |
二、持久执行基础
Durable Execution 的基础仍然是前文介绍过的 Checkpointer。
Graph 编译时传入 Checkpointer:
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:
config = {
"configurable": {
"thread_id": "order-1001"
}
}
thread_id 用来标识这一条执行历史。同一个 Thread 下产生的 Checkpoint 会被关联起来,LangGraph 因此能够找到当前状态和之前的运行记录。
Checkpoint 保存在哪里
对于一个顺序 Graph:
START → A → B → C → END
可以简化为:
输入
↓ 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 传入:
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 也可以直接传入异常类型或自定义判断函数。
下面用一个确定性的例子观察重试过程:
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)
运行结果
attempt=1
attempt=2
attempt=3
{'result': 'service ok'}
前两次执行抛出 ConnectionError,符合 retry_on 的条件,所以当前 Node 再次运行。第三次成功后,Graph 才进入后续步骤。
实际项目中,重试条件最好与错误类型对应。
网络超时、服务暂时不可用等异常通常适合重试。如果是参数格式错误、业务规则不满足、资源确实不存在,重复执行则往往没有意义。否则 Retry 只会把一次明确的失败延长成多次相同的失败。
四、执行恢复过程
Retry 有次数限制。如果外部服务持续不可用,Node 最终仍然可能失败。
有了 Checkpointer,已经完成的步骤可以保留下来。当服务恢复之后,可再从保存的执行位置继续。
下面把过程简化成三个节点:
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 中尚未完成的执行。
运行结果如下:
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:
┌─ fetch_profile ─ 成功
START ──┤
└─ fetch_orders ─ 失败
此时整个 Super-step 还没有完成,但 fetch_profile 已经产生了有效结果。
LangGraph 会保存这种已经完成的节点写入,这些记录称为 Pending Writes。恢复执行时,成功节点的结果可以继续使用,不必因为另一个并行节点失败而全部重算。
因此,恢复后的执行更接近:
fetch_profile 使用已有结果
fetch_orders 重新执行
这类机制对并行检索、多数据源查询、多个 Agent 同时执行任务尤其重要。某一路调用失败时,已经完成的模型请求或远程查询不必全部重复。
Checkpoint 和 Pending Writes 的作用范围也有区别:
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。历史记录默认按时间倒序返回。
例如:
history = list(
graph.get_state_history(config)
)
before_finalize = next(
state
for state in history
if state.next == ("finalize",)
)
这里找到的是 finalize 执行前的 Checkpoint。
随后:
result = graph.invoke(
None,
before_finalize.config,
)
Graph 会从这个历史位置继续执行。
Checkpoint 之前已经完成的 Node 不会重新运行;后面的 Node 会再次执行。因此其中的模型调用、HTTP 请求和 Interrupt 也会重新执行,而新的执行结果也并不保证与原来运行结果完全相同。
Fork
Fork 需要使用 update_state()。
它接收一个历史 Checkpoint 的配置以及新的 State 值:
fork_config = graph.update_state(
before_finalize.config,
values={"risk": "high"},
)
随后继续运行:
result = graph.invoke(
None,
fork_config,
)
update_state() 会创建一个新的 Checkpoint 分支,原来的历史记录仍然保留。后续节点读取修改后的 State,因此可能得到不同的路径或结果。
这种能力很适合调试 Agent。
例如前面的检索和模型调用已经花费几十秒,只想验证“如果风险等级改成 high,后面的决策会怎样”,就可以直接从相应 Checkpoint 建立分支,不必反复运行整条 Graph。
七、外部操作与幂等
Checkpoint 能够保存执行状态,但已经写入外部系统的数据仍然需要应用自身保证一致性。
例如一个支付 Node:
调用支付接口
↓
支付成功
↓
进程中断
↓
Node 尚未完整结束
恢复以后,这个 Node 可能再次执行。如果第二次调用又产生一笔支付,就会造成重复扣款。
因此,发送邮件、创建订单、扣款、写入外部数据库等操作,都应考虑幂等(Idempotency)。
幂等表示同一个业务操作重复执行多次,最终结果与执行一次保持一致。常见做法包括使用幂等键、业务唯一键,或者在执行前查询对应操作是否已经完成。
Node 的粒度也会影响恢复成本。Node 如何划分,会直接影响失败以后需要重复多少工作。
如果一个 Node 内连续完成查询数据库、调用三个 API、调用一次模型,最后一步失败,那么重新进入这个 Node 时,其中一部分操作可能再次执行,正如前面一再强调的。此时,分成多个 Node 通常会更合适。
总结
LangGraph 的容错由几个不同层面的机制共同完成。
RetryPolicy 处理 Node 内部的临时故障;Checkpoint 保存 Graph 的执行状态;Pending Writes 记录并行 Super-step 中已经完成的工作;Time Travel 则利用历史 Checkpoint 完成 Replay 和 Fork。
实际设计时,可以按故障的范围选择处理方式:
临时错误
→ Retry
运行中断
→ Checkpoint 恢复
并行节点部分失败
→ Pending Writes
历史状态重新执行
→ Replay / Fork
外部系统写入
→ 幂等控制
Durable Execution 的价值也体现在这里:一个 Agent 的执行不再依赖某一次进程必须完整运行到底。只要状态、执行边界和外部系统写入设计清楚,长时间任务就有了失败后继续恢复的基础。
社区讨论
参与讨论
有问题或想法?欢迎继续讨论。