Appearance
容错与控制(Fault Tolerance)
Agent 工作流不是"一次成功或整体失败"的脚本。模型会超时、检索器会宕机、数据库会拒绝写入。本页讲清楚:哪些错误应该重试,哪些应该进入 fallback,以及如何在运行时优雅地停止一个正在执行的长任务。
先修知识
- 持久化 - Checkpoint 是恢复的基础
- Durable Execution - 幂等性与恢复语义
必须深刻理解,不能跳过:容错的本质不是"消灭错误",而是"给每个可能的故障匹配正确的响应策略"。一个超时的模型调用应该重试;一个数据库连接拒绝应该降级或告警,而不是无限重试把连接池打满;一个用户主动取消的请求应该优雅停机,而不是抛异常把半成品状态写进去。把这三类故障混为一谈,是生产环境最常见的容错设计错误。
前端类比:先建立直觉
| 前端概念 | LangGraph 容错机制 | 说明 |
|---|---|---|
AbortController + setTimeout | add_node(..., timeout=) | 给单个操作设定期限 |
| React Error Boundary | add_node(..., error_handler=) | 捕获子树错误,渲染 fallback UI |
axios.retry + 指数退避 | retry_policy= | 网络抖动自动重试 |
navigator.sendBeacon 关页清理 | control.request_drain() | 优雅停机,完成当前工作后退出 |
useEffect cleanup 函数 | Runtime 注入 | 运行时感知自身执行状态 |
LangGraph 原生语义:LangGraph v1.2(2026 年 5 月)引入了 Python 专属的节点级执行控制能力。这些能力在 add_node() 和 set_node_defaults() 上配置,运行时通过 Runtime 注入暴露执行上下文。它们与已有的 Checkpoint/Durable Execution 层配合,构成从"单节点故障"到"整图优雅停机"的完整容错链路。
版本要求
本页所有 API 均需要 langgraph >= 1.2.0(Python 专属)。超时和错误处理是 Python 独有能力,JavaScript 版 LangGraph 不包含这些 API。
1. 错误分类:重试还是 Fallback
在写任何容错代码之前,先把错误分到正确的桶里:
| 错误类型 | 典型表现 | 正确策略 | 错误策略 |
|---|---|---|---|
| 瞬时故障 | 模型 API 429/503、网络抖动、连接重置 | 重试(指数退避) | 直接 fallback(浪费已建立的连接) |
| 超时故障 | 模型 60 秒不返回、检索器卡死 | 超时 + 重试或 fallback | 无限等待(阻塞整个图) |
| 逻辑错误 | 输入格式错误、状态字段缺失 | error_handler 降级 | 重试(结果不会变) |
| 基础设施故障 | 数据库连接拒绝、磁盘满 | 降级 + 告警 | 无限重试(打满连接池) |
| 用户取消 | 用户主动停止请求 | 优雅停机 | 抛异常写半成品状态 |
EnviroNexus 映射:在环保知识库场景中--模型超时(瞬时,重试);Retriever 向量检索失败(降级到 keyword 检索);数据库写入失败(降级 + 告警,不阻塞已有答案返回)。这三类故障的处理策略完全不同,不能混用。
2. 节点超时(timeout)
最小示例
python
import os
from langgraph.graph import StateGraph, START, END, MessagesState
from langgraph.types import TimeoutPolicy
# 模型配置用环境变量
model_name = os.environ["LLM_MODEL"]
def slow_node(state: MessagesState):
"""可能很慢的节点"""
# 模拟一个可能超时的操作
import time
time.sleep(30)
return {"messages": [{"role": "assistant", "content": "完成"}]}
builder = StateGraph(MessagesState)
# timeout=10 表示该节点最多执行 10 秒
builder.add_node("slow", slow_node, timeout=10)
builder.add_edge(START, "slow")
builder.add_edge("slow", END)
graph = builder.compile()timeout 参数的三种形式
python
from datetime import timedelta
from langgraph.types import TimeoutPolicy
# 形式 1:数字(秒)
builder.add_node("node_a", fn_a, timeout=30)
# 形式 2:timedelta
builder.add_node("node_b", fn_b, timeout=timedelta(minutes=2))
# 形式 3:TimeoutPolicy(精细控制)
builder.add_node("node_c", fn_c, timeout=TimeoutPolicy(
run_timeout=30, # 单次执行最长 30 秒
idle_timeout=60, # 两次活动之间最长间隔 60 秒
refresh_on="heartbeat" # 收到心跳时重置 idle 计时器
))| 参数 | 含义 | 适用场景 |
|---|---|---|
run_timeout | 从开始到结束的总时间上限 | 防止单次执行卡死 |
idle_timeout | 两次 I/O 活动之间的间隔上限 | 检测异步节点的"静默卡死" |
refresh_on | 什么事件重置 idle 计时器 | 配合 runtime.heartbeat() 使用 |
超时后的行为
当节点超时,LangGraph 抛出 NodeTimeoutError(仅限 async 节点)。如果配置了 error_handler,会进入错误处理流程;否则异常向上传播。
python
from langgraph.errors import NodeTimeoutError
async def risky_async_node(state: MessagesState):
# async 节点超时会触发 NodeTimeoutError
...
builder.add_node("risky", risky_async_node, timeout=5)
# 如果没有 error_handler,NodeTimeoutError 会被抛出
try:
result = await graph.ainvoke(input, config)
except NodeTimeoutError as e:
print(f"节点 {e} 超时")仅限 async 节点
timeout= 和 NodeTimeoutError 仅对 async 节点 生效。同步节点无法被中断(Python 的 GIL 限制)。如果你的节点需要超时保护,请使用 async def。
3. 节点错误处理(error_handler)
核心概念
error_handler 是一个函数,当节点抛出异常时被调用。它可以返回一个 Command 来决定图的下一步行为--更新状态、跳转到 fallback 节点,或者终止执行。
python
from langgraph.graph import StateGraph, START, END, MessagesState
from langgraph.types import Command
from langgraph.errors import NodeError
def call_model(state: MessagesState):
"""主节点:调用模型生成回答"""
# 可能抛出异常的模型调用
...
def handle_model_error(error: NodeError) -> Command:
"""错误处理函数:收到 NodeError,返回 Command"""
print(f"节点 [{error.node}] 出错: {error.error}")
# 策略 1:更新状态并跳转到 fallback 节点
return Command(
update={"messages": [{"role": "assistant", "content": "服务暂时不可用,请稍后重试"}]},
goto="fallback_response"
)
def fallback_response(state: MessagesState):
"""降级响应节点"""
return {"messages": [{"role": "assistant", "content": "降级处理已完成"}]}
builder = StateGraph(MessagesState)
builder.add_node("model", call_model, error_handler=handle_model_error)
builder.add_node("fallback_response", fallback_response)
builder.add_edge(START, "model")
# 注意:error_handler 返回的 goto 会覆盖正常边
builder.add_edge("model", END)
builder.add_edge("fallback_response", END)
graph = builder.compile()error_handler 的签名
python
from langgraph.errors import NodeError
from langgraph.types import Command
def my_error_handler(error: NodeError) -> Command:
# error.node -> 出错的节点名(str)
# error.error -> 原始异常对象(BaseException)
...
return Command(update={...}, goto="...")NodeError 是一个 frozen dataclass,包含:
.node: str - 出错的节点名称.error: BaseException - 原始异常
返回值必须是 Command,通过 update= 更新状态,通过 goto= 指定下一步去哪个节点。
error_handler 返回值策略
| 返回值 | 效果 | 适用场景 |
|---|---|---|
Command(update={...}, goto="fallback") | 更新状态并跳转 | 降级到备用方案 |
Command(update={...}, goto=END) | 更新状态后终止 | 记录错误并返回部分结果 |
Command(update={...}) | 只更新状态,继续正常边 | 注入错误信息到状态中 |
4. set_node_defaults(图级默认值)
当多个节点需要相同的容错配置时,用 set_node_defaults() 设置图级默认值,避免逐节点重复:
python
from langgraph.graph import StateGraph, START, END, MessagesState
from langgraph.pregel import RetryPolicy
# 图级默认值:所有节点都继承这些配置
builder = StateGraph(MessagesState)
builder.set_node_defaults(
retry_policy=RetryPolicy(max_attempts=3),
timeout=30,
error_handler=handle_model_error,
)
# 后续 add_node 不需要再重复这些参数
builder.add_node("research", research_node)
builder.add_node("analyze", analyze_node)
builder.add_node("report", report_node)set_node_defaults() 接受四个参数:
| 参数 | 类型 | 说明 |
|---|---|---|
retry_policy | RetryPolicy | 重试策略,所有节点继承 |
error_handler | Callable[[NodeError], Command] | 错误处理函数 |
timeout | int / timedelta / TimeoutPolicy | 超时配置 |
cache_policy | 缓存策略 | 节点结果缓存(v1.2) |
优先级:
add_node()上的显式参数会覆盖set_node_defaults()的默认值。如果两者都没配,则没有容错保护。
5. retry_policy(重试策略)
基本用法
python
from langgraph.pregel import RetryPolicy
# 默认重试策略
retry = RetryPolicy(max_attempts=3)
# 自定义退避
from langchain_core.runnables import RunnableConfig
import random
retry = RetryPolicy(
max_attempts=5,
initial_interval=1.0, # 首次重试等待 1 秒
backoff_factor=2.0, # 每次翻倍:1, 2, 4, 8
max_interval=60.0, # 单次最长等 60 秒
jitter=True, # 加随机抖动避免惊群
)
builder.add_node("flaky_api", call_flaky_api, retry_policy=retry)哪些错误应该重试
python
from langgraph.pregel import RetryPolicy
# 只重试瞬时异常
def should_retry(error: BaseException) -> bool:
"""判断异常是否值得重试"""
# 瞬时故障 -> 重试
if isinstance(error, (ConnectionError, TimeoutError)):
return True
# 逻辑错误 -> 不重试(重试结果不会变)
if isinstance(error, (ValueError, KeyError, TypeError)):
return False
# 默认重试
return True
retry = RetryPolicy(
max_attempts=3,
retry_on=should_retry, # 自定义重试判断
)关键原则:重试只对瞬时故障有意义。如果一个节点因为输入格式错误而失败,重试 100 次结果也一样。
retry_on参数让你精确控制哪些异常值得重试。
6. 优雅停机(RunControl)
场景
当用户取消请求、服务需要缩容、或检测到异常需要停止当前执行时,你需要一种方式让图"优雅地停下来"--不是抛异常丢失状态,而是完成当前 superstep 后干净退出。
使用 RunControl
python
from langgraph.runtime import RunControl
from langgraph.errors import GraphDrained
control = RunControl()
# 启动一个长运行图
config = {
"configurable": {"thread_id": "long-task-1"},
"run_control": control, # 注入 RunControl
}
# 在另一个协程/线程中请求停机
import asyncio
async def cancel_after_5s():
await asyncio.sleep(5)
# 请求排空:完成当前 superstep 后停止
control.request_drain(reason="用户取消请求")
asyncio.create_task(cancel_after_5s())
# 图会在完成当前 superstep 后抛出 GraphDrained
try:
result = await graph.ainvoke(input, config)
except GraphDrained as e:
print(f"图已优雅停机: {e}")
# checkpoint 已保存,可以用 invoke(None, config) 恢复request_drain vs 强制中断
| 方式 | 行为 | checkpoint 状态 |
|---|---|---|
control.request_drain(reason) | 完成当前 superstep 后停止 | 保存到最后完成的 superstep |
| 强制 kill 进程 | 立即终止 | 依赖上一次 checkpoint(可能不完整) |
前端类比:request_drain() 类似于 Kubernetes 的 kubectl drain--不是直接杀 Pod,而是先标记为不可调度,等现有请求处理完再驱逐。
7. Runtime 注入
什么是 Runtime 注入
LangGraph v1.2 允许节点函数声明 runtime: Runtime 参数,运行时自动注入。Runtime 对象提供执行上下文信息,让节点感知自身运行状态。
python
from langgraph.runtime import Runtime
from langgraph.graph import StateGraph, START, END, MessagesState
async def long_running_node(state: MessagesState, runtime: Runtime):
"""需要长时间运行的节点,通过 runtime 汇报心跳"""
# 1. 检查是否被请求停机
if runtime.drain_requested:
return {"messages": [{"role": "assistant", "content": "任务被中断"}]}
# 2. 定期发送心跳,重置 idle_timeout
for i in range(10):
runtime.heartbeat() # 重置 idle 计时器
# ... 执行一批工作 ...
# 3. 读取执行信息
info = runtime.execution_info
print(f"当前尝试次数: {info.get('node_attempt', 1)}")
return {"messages": [{"role": "assistant", "content": "长任务完成"}]}
builder = StateGraph(MessagesState)
builder.add_node("worker", long_running_node, timeout=300)
builder.add_edge(START, "worker")
builder.add_edge("worker", END)
graph = builder.compile()Runtime 的属性
| 属性/方法 | 类型 | 说明 |
|---|---|---|
runtime.execution_info | dict | 当前执行元信息(如 node_attempt 尝试次数) |
runtime.heartbeat() | method | 发送心跳,重置 idle_timeout 计时器 |
runtime.drain_requested | bool | 是否收到停机请求 |
runtime.drain_reason | str / None | 停机原因(来自 request_drain(reason)) |
heartbeat 与 idle_timeout 配合
python
from langgraph.types import TimeoutPolicy
from langgraph.runtime import Runtime
async def streaming_process(state, runtime: Runtime):
"""流式处理节点:持续工作但定期汇报"""
items = state["items"]
results = []
for item in items:
# 每处理一个 item 就发一次心跳
runtime.heartbeat()
# 如果被请求停机,保存当前进度后退出
if runtime.drain_requested:
return {"partial_results": results, "interrupted": True}
results.append(await process_item(item))
return {"results": results, "interrupted": False}
# idle_timeout=30: 如果 30 秒没有心跳就认为卡死了
builder.add_node(
"streaming_process",
streaming_process,
timeout=TimeoutPolicy(idle_timeout=30, refresh_on="heartbeat")
)8. 完整示例:带容错的 EnviroNexus 检索流程
python
import os
from langgraph.graph import StateGraph, START, END, MessagesState
from langgraph.types import Command, TimeoutPolicy
from langgraph.errors import NodeError, NodeTimeoutError
from langgraph.pregel import RetryPolicy
from langgraph.runtime import Runtime
from langchain.chat_models import init_chat_model
from typing import Any
model_name = os.environ["LLM_MODEL"]
model = init_chat_model(model_name)
class QueryState(MessagesState):
query: str
retrieved_docs: list[dict]
answer: str
degraded: bool
async def vector_retrieve(state: QueryState, runtime: Runtime) -> dict:
"""向量检索节点:可能超时或失败"""
runtime.heartbeat()
# 模拟向量检索
try:
docs = await hybrid_search(state["query"])
return {"retrieved_docs": docs}
except Exception as e:
# 如果是连接错误,向上抛出(由 error_handler 处理)
raise e
def handle_retrieval_error(error: NodeError) -> Command:
"""检索失败的降级策略:回退到 keyword 检索"""
print(f"向量检索失败 [{error.node}]: {error.error}")
return Command(
update={
"degraded": True,
"retrieved_docs": keyword_fallback(query=""), # keyword 降级检索
},
goto="generate_answer",
)
async def generate_answer(state: QueryState) -> dict:
"""生成回答节点"""
context = "\n".join(d.get("content", "") for d in state["retrieved_docs"])
prompt = f"基于以下资料回答问题:\n{context}\n\n问题:{state['query']}"
response = await model.ainvoke([
{"role": "user", "content": prompt}
])
return {"answer": response.content}
# 构建图
builder = StateGraph(QueryState)
# 向量检索节点:超时 10 秒,重试 2 次,失败降级到 keyword
builder.add_node(
"retrieve",
vector_retrieve,
timeout=TimeoutPolicy(run_timeout=10, idle_timeout=5, refresh_on="heartbeat"),
retry_policy=RetryPolicy(max_attempts=2),
error_handler=handle_retrieval_error,
)
# 生成节点:超时 60 秒,重试 3 次
builder.add_node(
"generate_answer",
generate_answer,
timeout=60,
retry_policy=RetryPolicy(max_attempts=3),
)
builder.add_edge(START, "retrieve")
builder.add_edge("retrieve", "generate_answer")
builder.add_edge("generate_answer", END)
# 图级默认值(可选)
builder.set_node_defaults(
retry_policy=RetryPolicy(max_attempts=2),
)
graph = builder.compile()容错恢复时序
9. DeltaChannel(beta):长线程降本
问题
在长对话线程中(如几百轮的消息历史),每次 superstep 都要序列化完整状态写 checkpoint。随着消息列表增长,checkpoint 体积和写入延迟线性增加。
DeltaChannel 的思路
DeltaChannel(v1.2 beta)不每次写完整快照,而是只存增量 delta。通过 snapshot_frequency=K 控制每 K 次增量后做一次完整快照(类似 Git 的 packed objects)。
python
# DeltaChannel 目前是 beta API,配置方式可能调整
# 核心参数:snapshot_frequency=K
# 每 K 次增量写入后做一次完整快照适用场景:消息列表很长(>100 轮)的对话型 Agent。对于短工作流,DeltaChannel 的收益不明显,甚至可能因为增量计算增加开销。
注意:DeltaChannel 仍是 beta 特性,API 可能在后续版本调整。生产使用前请核实最新官方文档。
10. 容错配置速查
| 能力 | API | 适用故障类型 | 版本 |
|---|---|---|---|
| 节点超时 | add_node(..., timeout=) | 超时故障 | >= 1.2 |
| 精细超时 | TimeoutPolicy(run_timeout=, idle_timeout=, refresh_on=) | 超时 + 静默卡死 | >= 1.2 |
| 错误处理 | add_node(..., error_handler=fn) | 逻辑错误 + 基础设施故障 | >= 1.2 |
| 重试 | add_node(..., retry_policy=) / set_node_defaults(retry_policy=) | 瞬时故障 | >= 1.0 |
| 图级默认值 | set_node_defaults(...) | 统一配置 | >= 1.2 |
| 优雅停机 | RunControl + request_drain() | 用户取消 / 缩容 | >= 1.2 |
| 运行时感知 | runtime: Runtime 注入 | 心跳 + 停机感知 | >= 1.2 |
| 长线程降本 | DeltaChannel(beta) | checkpoint 体积膨胀 | >= 1.2 beta |
11. 不适用场景与能力边界
- 同步节点无法超时:
timeout=和NodeTimeoutError仅对 async 节点生效。如果你的节点是同步的,要么改成 async,要么在节点内部自行实现超时逻辑。 - error_handler 不能"重试当前节点":
error_handler返回Command只能跳转到其他节点或终止。如果需要重试,用retry_policy=。 - DeltaChannel 是 beta:不要在核心生产链路依赖它,除非你愿意跟进 API 变更。
- Runtime 注入不等同于 context_schema:
runtime: Runtime是运行时执行上下文(心跳、停机状态),context_schema是业务上下文(用户 ID 等)。两者用途不同。
要点回顾
| 概念 | 一句话 |
|---|---|
| 错误分类 | 瞬时故障重试,逻辑错误降级,用户取消优雅停机 |
timeout= | 给节点设定期限,仅 async 节点生效 |
error_handler= | 捕获节点错误,返回 Command 决定下一步 |
set_node_defaults() | 图级默认容错配置,减少重复 |
retry_policy= | 瞬时故障自动重试,配合 retry_on 精确控制 |
RunControl | 优雅停机,完成当前 superstep 后干净退出 |
runtime: Runtime | 运行时注入,心跳 + 停机感知 + 执行信息 |
DeltaChannel | 长线程增量 checkpoint 降本(beta) |
先修与下一步
- 先修:持久化 | Durable Execution
- 下一步:事件流 | LangGraph Studio