Skip to content

容错与控制(Fault Tolerance)

Agent 工作流不是"一次成功或整体失败"的脚本。模型会超时、检索器会宕机、数据库会拒绝写入。本页讲清楚:哪些错误应该重试,哪些应该进入 fallback,以及如何在运行时优雅地停止一个正在执行的长任务。

先修知识

必须深刻理解,不能跳过:容错的本质不是"消灭错误",而是"给每个可能的故障匹配正确的响应策略"。一个超时的模型调用应该重试;一个数据库连接拒绝应该降级或告警,而不是无限重试把连接池打满;一个用户主动取消的请求应该优雅停机,而不是抛异常把半成品状态写进去。把这三类故障混为一谈,是生产环境最常见的容错设计错误。

前端类比:先建立直觉

前端概念LangGraph 容错机制说明
AbortController + setTimeoutadd_node(..., timeout=)给单个操作设定期限
React Error Boundaryadd_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。

🔗 Fault Tolerance 官方文档


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_policyRetryPolicy重试策略,所有节点继承
error_handlerCallable[[NodeError], Command]错误处理函数
timeoutint / 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_infodict当前执行元信息(如 node_attempt 尝试次数)
runtime.heartbeat()method发送心跳,重置 idle_timeout 计时器
runtime.drain_requestedbool是否收到停机请求
runtime.drain_reasonstr / 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_schemaruntime: Runtime 是运行时执行上下文(心跳、停机状态),context_schema 是业务上下文(用户 ID 等)。两者用途不同。

要点回顾

概念一句话
错误分类瞬时故障重试,逻辑错误降级,用户取消优雅停机
timeout=给节点设定期限,仅 async 节点生效
error_handler=捕获节点错误,返回 Command 决定下一步
set_node_defaults()图级默认容错配置,减少重复
retry_policy=瞬时故障自动重试,配合 retry_on 精确控制
RunControl优雅停机,完成当前 superstep 后干净退出
runtime: Runtime运行时注入,心跳 + 停机感知 + 执行信息
DeltaChannel长线程增量 checkpoint 降本(beta)

先修与下一步

参考

学习文档整合站点