Agent工坊

【Agent工坊】Agent 不死机指南:工具调用的四层防护实战

你的 AI Agent 在生产环境崩溃过吗?本文给你一套完整的防御性编程方案,让 Agent 工具调用"打不死"。


AI Agent 最脆弱的环节是什么?

不是模型不够聪明,不是 prompt 写得不好——而是工具调用失败后,整个 Agent 直接崩溃,没有任何恢复机制

我见过无数次这样的场景:Agent 调用一个搜索接口,网络抖了一下,接口超时,Agent 拿到一个 ConnectionError,然后整个任务流程应声倒地。用户只看到终端里那行 Traceback (most recent call last),只能手动重启、重新描述需求、再等 Agent 从头跑一遍。

如果这个 Agent 是定时任务在凌晨 3 点自动运行的呢?如果它是客户支付系统里的一环呢?

这不是模型的问题。是我们没有给 Agent 装安全网

我花了三个月时间,把团队里 6 个 Agent 的崩溃率从"每天至少一次"降到"连续两周无人工介入"。秘诀不是换了更好的模型,而是在工具调用层加了四层防护。

今天完整拆解这套方案。


一、问题全景:你的 Agent 工具调用为什么总是崩?

一个典型的 Agent 工具调用链路长这样:

用户指令 → LLM 决策 → 调用工具 A → 调用工具 B → 解析结果 → 返回用户

这个链条上有五个可能断裂的点。我把过去半年遇到的故障做了统计:

故障类型 触发条件 出现频率 典型后果
网络超时 DNS 解析慢、CDN 节点故障 几乎每天 Agent 拿不到搜索结果,任务中断
API 限流 (429) 并发量突增、额度耗尽 每周 2-3 次 连续请求被拒绝,任务卡死
认证过期 (401) Token 到期未刷新 每月 1 次 所有调用全线失败
JSON 解析失败 API 返回格式变更 每周 1-2 次 Agent 拿到数据但解析报错
下游服务宕机 (502/503) 第三方服务维护 每月 2-3 次 依赖方挂了,Agent 也挂了

这里的核心矛盾是:LLM 本身不具备错误恢复的能力。它看到一个 Error 字符串,要么瞎猜一个结果(幻觉),要么直接告诉用户"我做不到"(放弃)。

解决办法只有一个:把错误处理从模型决策层剥离,下沉到工具调用的基础设施层。让模型只做它擅长的事(理解、推理、决策),基础设施负责保证每一次工具调用都能返回有意义的结果。


二、第一层防护:工具包装器——永远不要裸调 API

我见过最让我血压升高的代码是这样写的:

def web_search(query: str) -> dict:
    resp = requests.get(f"https://api.search.com?q={query}", timeout=10)
    return resp.json()

这段代码作者的想法是"反正搜索 API 平时都好好的"。结果上线第一周,凌晨 3 点 API 维护,整个 Agent 批处理任务跑了 30 条后报错中断,第二天早晨才发现 70% 的数据都没处理到。

正确写法:给工具套上"防弹衣"

import requests
import logging
from typing import Optional, Dict, Any

logger = logging.getLogger("agent.tools")

def safe_web_search(
    query: str,
    fallback: Optional[Dict[str, Any]] = None
) -> Dict[str, Any]:
    """
    安全搜索——网络故障时返回结构化空结果,Agent 能继续工作
    """
    try:
        resp = requests.get(
            "https://api.search.example.com/v1/search",
            params={"q": query},
            timeout=10,
            headers={"User-Agent": "AgentTool/1.0"}
        )
        resp.raise_for_status()
        data = resp.json()
        logger.info(f"搜索成功: {query[:30]}, 返回 {len(data.get('results', []))} 条")
        return data

    except requests.Timeout:
        logger.warning(f"搜索超时 (10s): {query[:50]}")

    except requests.ConnectionError as e:
        logger.warning(f"搜索连接失败: {query[:50]}{e}")

    except requests.HTTPError as e:
        status = e.response.status_code if e.response else '?'
        logger.warning(f"搜索 HTTP {status}: {query[:50]}")

    except json.JSONDecodeError as e:
        logger.error(f"搜索返回非 JSON 数据: {e}")

    except Exception as e:
        logger.error(f"搜索未知错误: {type(e).__name__}: {e}")

    # 关键:返回结构化空结果而不是 None
    return fallback or {
        "results": [],
        "query": query,
        "error": "搜索服务暂时不可用,请稍后重试",
        "timestamp": int(time.time())
    }

三个关键设计点

第一,每种异常类型单独捕获。Timeout 和 401 的处理策略完全不同——Timeout 值得重试,401 重试一百次也没用。分开捕获才能在日志里看到到底是什么原因失败的。

第二,返回语义合法的空结果,而不是 None 或空字符串。Agent 看到 {"results": []} 会理解为"没搜到相关内容",继续执行下一步。但 Agent 看到 None 会困惑,可能瞎编一个搜索结果(这就是幻觉的来源)。

第三,日志带上下文。只打 "搜索失败" 没用,出问题时你根本不知道是哪个查询失败的。至少带上 query 前 50 个字符和异常类型。


三、第二层防护:指数退避重试——瞬时故障的自动修复

包装器解决了"不崩溃"的问题,但很多故障其实是瞬时的——网络抖了一下、API 刚好在扩容、CDN 节点刚好在切换。这些场景重试 1-2 次大概率能恢复。

但重试有讲究。我见过有人这么写:

for i in range(5):
    try:
        return call_api()
    except:
        time.sleep(1)  # ❌ 固定等待 + 无脑重试

这样做有两个问题:第一,固定 1 秒间隔意味着 10 个 Agent 实例同时重试会形成"惊群效应"把服务打得更死;第二,对所有异常都重试(包括 401、403)纯属浪费。

正确实现:指数退避 + 随机抖动

import time
import random
from functools import wraps
from typing import Type, Tuple

def retry_with_backoff(
    max_retries: int = 3,
    base_delay: float = 1.0,
    max_delay: float = 30.0,
    backoff_factor: float = 2.0,
    retryable_exceptions: Tuple[Type[Exception], ...] = (
        requests.Timeout,
        requests.ConnectionError,
    )
):
    """
    指数退避重试装饰器

    退避公式:
        delay = min(base_delay × factor^attempt + jitter, max_delay)

    实际等待时序(base_delay=1.0, factor=2.0):
        第 1 次重试前等待:1.0 ~ 1.5 秒
        第 2 次重试前等待:2.0 ~ 3.0 秒
        第 3 次重试前等待:4.0 ~ 6.0 秒

    为什么需要 jitter(随机抖动)?
        如果 10 个 Agent 同时遇到故障,它们会在几乎相同的时间点
        重试。jitter 让每个 Agent 的等待时间略有不同,避免同时冲击
        目标服务。这是 AWS、Google Cloud 等云服务 SDK 的标准做法。
    """
    def decorator(func):
        @wraps(func)
        def wrapper(*args, **kwargs):
            last_exception = None
            for attempt in range(max_retries + 1):
                try:
                    return func(*args, **kwargs)
                except retryable_exceptions as e:
                    last_exception = e
                    if attempt == max_retries:
                        logger.error(
                            f"{func.__name__} 重试 {max_retries} 次全部失败: {e}"
                        )
                        raise
                    delay = min(
                        base_delay * (backoff_factor ** attempt),
                        max_delay
                    )
                    jitter = random.uniform(0, delay * 0.5)
                    total_delay = delay + jitter
                    logger.warning(
                        f"{func.__name__}{attempt + 1}/{max_retries} 次重试,"
                        f"等待 {total_delay:.1f}s 后重试"
                    )
                    time.sleep(total_delay)
            raise last_exception
        return wrapper
    return decorator


# === 实战:给 LLM API 加上智能重试 ===

@retry_with_backoff(
    max_retries=3,
    base_delay=1.0,
    retryable_exceptions=(requests.Timeout, requests.ConnectionError)
)
def call_llm_api(prompt: str, model: str = "gpt-4o") -> str:
    """
    调用 LLM API,只有网络问题才重试。
    401/403/400 等客户端错误不会重试,直接抛出让上层处理。
    """
    resp = requests.post(
        "https://api.openai.com/v1/chat/completions",
        headers={"Authorization": f"Bearer {API_KEY}"},
        json={
            "model": model,
            "messages": [{"role": "user", "content": prompt}],
            "temperature": 0.7,
        },
        timeout=30
    )

    # 429 (Rate Limit) — 用响应头里的 Retry-After 等待
    if resp.status_code == 429:
        retry_after = int(resp.headers.get("Retry-After", 5))
        logger.warning(f"遇到限流 (429),等待 {retry_after}s")
        time.sleep(retry_after)
        raise requests.ConnectionError("Rate limited")  # 触发装饰器重试

    resp.raise_for_status()
    return resp.json()["choices"][0]["message"]["content"]

输出效果

[WARNING] call_llm_api  1/3 次重试等待 1.3s 后重试
[WARNING] call_llm_api  2/3 次重试等待 2.7s 后重试
[INFO]    call_llm_api 调用成功 3 次尝试

四个关键设计点

  1. 只重试可恢复的错误:Timeout、ConnectionError、429 值得重试;400、401、403 重试毫无意义
  2. 必须加 jitter:不加随机抖动的后果是——凌晨批量任务触发限流时,所有 Agent 几乎同时重试,形成二次冲击波
  3. 设 max_delay 上限:生产环境不要让退避时间无限膨胀,30 秒是合理上限
  4. 429 特殊处理:优先使用 API 返回的 Retry-After 头,比自己的退避算法更准确

四、第三层防护:降级回退链——主力挂了备胎上

重试只能解决瞬时故障。但如果 Claude API 的 Token 过期了呢?如果某个搜索服务整站维护两小时呢?

这时候需要降级(Fallback):提前准备好备用方案,主路不通走备路。

多级回退链路模型

优先级 1:Claude 4 Opus(最强,但也最贵、最容易被限流)
   ↓ 失败
优先级 2:GPT-4o(稳定,备选主力)
   ↓ 失败
优先级 3:DeepSeek V3(便宜,兜底方案)
   ↓ 失败
兜底:返回预设的友好话术

实现代码

from typing import List, Callable, Any, Tuple

class FallbackChain:
    """
    多级回退执行链

    设计要点:
    - 先注册的处理器优先级最高
    - 任一成功即返回,不继续尝试后面的
    - 全部失败则返回预设默认值
    - 对调用方完全透明——Agent 不需要知道内部有降级逻辑
    """

    def __init__(self, name: str = "fallback"):
        self.name = name
        self.handlers: List[Tuple[Callable, str]] = []
        self.default_value: Any = None

    def add_handler(self, handler: Callable, description: str = ""):
        """添加处理器,先添加的优先执行"""
        self.handlers.append((handler, description))
        return self

    def set_default(self, value: Any):
        """全部失败时的兜底值"""
        self.default_value = value
        return self

    def execute(self, *args, **kwargs) -> Any:
        for i, (handler, desc) in enumerate(self.handlers):
            try:
                result = handler(*args, **kwargs)
                if i > 0:
                    logger.info(
                        f"[{self.name}] 主路失败,降级到「{desc}」成功"
                    )
                return result
            except Exception as e:
                logger.warning(
                    f"[{self.name}] 第 {i+1} 路「{desc}」失败: "
                    f"{type(e).__name__}"
                )

        logger.error(
            f"[{self.name}] 全部 {len(self.handlers)} 路处理器失败,使用兜底值"
        )
        return self.default_value


# ═══════════════════════════════════════
# 实战:多模型智能回退
# ═══════════════════════════════════════

def call_claude(prompt: str) -> str:
    """Claude 4 Opus — 主力模型"""
    resp = requests.post(
        "https://api.anthropic.com/v1/messages",
        headers={
            "x-api-key": CLAUDE_KEY,
            "anthropic-version": "2023-06-01"
        },
        json={
            "model": "claude-4-opus-20250514",
            "max_tokens": 2000,
            "messages": [{"role": "user", "content": prompt}]
        },
        timeout=30
    )
    resp.raise_for_status()
    return resp.json()["content"][0]["text"]

def call_gpt4o(prompt: str) -> str:
    """GPT-4o — 稳定备选"""
    resp = requests.post(
        "https://api.openai.com/v1/chat/completions",
        headers={"Authorization": f"Bearer {OPENAI_KEY}"},
        json={
            "model": "gpt-4o",
            "messages": [{"role": "user", "content": prompt}]
        },
        timeout=30
    )
    resp.raise_for_status()
    return resp.json()["choices"][0]["message"]["content"]

def call_deepseek(prompt: str) -> str:
    """DeepSeek V3 — 便宜兜底,API 兼容 OpenAI 格式"""
    resp = requests.post(
        "https://api.deepseek.com/v1/chat/completions",
        headers={"Authorization": f"Bearer {DEEPSEEK_KEY}"},
        json={
            "model": "deepseek-chat",
            "messages": [{"role": "user", "content": prompt}]
        },
        timeout=30
    )
    resp.raise_for_status()
    return resp.json()["choices"][0]["message"]["content"]

# 组装回退链
smart_llm = (
    FallbackChain("LLM调用")
    .add_handler(call_claude, "Claude 4 Opus(主力)")
    .add_handler(call_gpt4o, "GPT-4o(降级1)")
    .add_handler(call_deepseek, "DeepSeek V3(降级2)")
    .set_default("抱歉,当前所有 AI 服务不可用,请稍后重试。")
)

# Agent 调用时完全感知不到内部的降级逻辑
result = smart_llm.execute("用 200 字介绍 AI Agent 的概念")

实际输出示例

[WARNING] [LLM调用]  1 Claude 4 Opus主力)」失败: RateLimitError
[INFO]    [LLM调用] 主路失败降级到GPT-4o降级1)」成功

关键设计原则

  • 成本优先排序:最便宜的放最后。目标是尽量不触发降级,但降级方案不能太贵
  • 兜底值必须有意义:空字符串、None 会让 Agent 困惑。预设一段友好的错误提示
  • 降级深度 2-3 级最优:超过 3 级说明你的第一选择本身有问题
  • 对上游完全透明:Agent 调用 smart_llm.execute() 和使用单个函数毫无区别

五、第四层防护:熔断器——避免故障雪崩

前面积三层防护能处理绝大部分故障。但还有一种场景更危险:某个服务持续故障 10 分钟

如果每次都重试 3 次 × 30 秒超时 = 90 秒,10 分钟下来就是 6-7 次无效等待。如果同时有 20 个 Agent 实例在跑,这就是 120+ 次无效 API 调用。

熔断器的工作原理

跟电路保险丝一个思路:

  • 正常状态(Closed):请求正常通过
  • 连续失败 N 次:熔断器打开(Open),拒绝所有请求,直接返回失败
  • 冷却 T 秒后:进入半开状态(Half-Open),允许少量探测请求
  • 探测成功:熔断器关闭,恢复正常
  • 探测失败:重新打开,继续冷却

实现代码

import time
from enum import Enum

class CircuitState(Enum):
    CLOSED = "closed"        # 正常
    OPEN = "open"            # 熔断中
    HALF_OPEN = "half_open"  # 探测恢复中

class CircuitBreaker:
    """
    熔断器实现

    配置建议:
    - 搜索/爬虫工具:failure_threshold=3, recovery=60s(快速熔断)
    - LLM API:failure_threshold=5, recovery=300s(成本敏感,谨慎熔断)
    """

    def __init__(
        self,
        name: str,
        failure_threshold: int = 5,
        recovery_timeout: float = 60.0,
        half_open_max: int = 1,
    ):
        self.name = name
        self.failure_threshold = failure_threshold
        self.recovery_timeout = recovery_timeout
        self.half_open_max = half_open_max

        self.state = CircuitState.CLOSED
        self.failure_count = 0
        self.last_failure_time = 0.0
        self.half_open_count = 0

    def call(self, func: Callable, *args, **kwargs) -> Any:
        if self.state == CircuitState.OPEN:
            elapsed = time.time() - self.last_failure_time
            if elapsed > self.recovery_timeout:
                logger.info(
                    f"[熔断器:{self.name}] 冷却 {elapsed:.0f}s 完成,进入半开探测"
                )
                self.state = CircuitState.HALF_OPEN
                self.half_open_count = 0
            else:
                raise CircuitBreakerOpenError(
                    f"[熔断器:{self.name}] 熔断中,"
                    f"剩余冷却 {self.recovery_timeout - elapsed:.0f}s"
                )

        if self.state == CircuitState.HALF_OPEN:
            if self.half_open_count >= self.half_open_max:
                raise CircuitBreakerOpenError(
                    f"[熔断器:{self.name}] 半开探测名额已用完"
                )
            self.half_open_count += 1

        try:
            result = func(*args, **kwargs)
            # 成功!恢复熔断器
            if self.state != CircuitState.CLOSED:
                logger.info(f"[熔断器:{self.name}] ✅ 恢复成功,关闭熔断")
                self.state = CircuitState.CLOSED
                self.failure_count = 0
            return result
        except Exception as e:
            self.failure_count += 1
            self.last_failure_time = time.time()
            if self.failure_count >= self.failure_threshold:
                logger.error(
                    f"[熔断器:{self.name}] 🔴 连续失败 {self.failure_count} 次,"
                    f"触发熔断!冷却 {self.recovery_timeout}s"
                )
                self.state = CircuitState.OPEN
            raise


class CircuitBreakerOpenError(Exception):
    """熔断器打开时抛出的异常,供上层做降级处理"""
    pass

使用效果

[ERROR]   [熔断器:web_search] 🔴 连续失败 5 触发熔断冷却 120s
[WARNING] [熔断器:web_search] 熔断中剩余冷却 87s  降级到本地缓存
[WARNING] [熔断器:web_search] 熔断中剩余冷却 43s  降级到本地缓存
[INFO]    [熔断器:web_search] 冷却 123s 完成进入半开探测
[INFO]    [熔断器:web_search]  恢复成功关闭熔断

六、组装四层防护:Agent 工具调用框架

把上面四层串起来,就得到一个生产级的工具执行器:

import random
import logging
from dataclasses import dataclass, field
from typing import Any, Dict, List, Optional, Callable

logger = logging.getLogger("agent.executor")

@dataclass
class ToolConfig:
    """单个工具的四层防护配置"""
    name: str
    handler: Callable
    max_retries: int = 3
    retry_delay: float = 1.0
    fallback_handlers: List[Callable] = field(default_factory=list)
    default_result: Any = None
    circuit_breaker: bool = True
    cb_threshold: int = 5
    cb_recovery: float = 120.0


class ResilientToolExecutor:
    """
    Agent 工具执行器——自动四层防护

    使用方式:
        executor = ResilientToolExecutor()
        executor.register(ToolConfig(name="web_search", handler=..., ...))
        result = executor.execute("web_search", "AI Agent 最新进展")
    """

    def __init__(self):
        self.tools: Dict[str, ToolConfig] = {}
        self.breakers: Dict[str, CircuitBreaker] = {}

    def register(self, cfg: ToolConfig):
        """注册工具配置"""
        self.tools[cfg.name] = cfg
        if cfg.circuit_breaker:
            self.breakers[cfg.name] = CircuitBreaker(
                name=cfg.name,
                failure_threshold=cfg.cb_threshold,
                recovery_timeout=cfg.cb_recovery,
            )
        logger.info(f"✅ 注册工具: {cfg.name}(重试{cfg.max_retries}次 + "
                     f"降级{len(cfg.fallback_handlers)}路 + "
                     f"{'熔断' if cfg.circuit_breaker else '无熔断'})")

    def execute(self, tool_name: str, *args, **kwargs) -> Any:
        """执行工具——Agent 唯一需要调用的入口"""
        cfg = self.tools.get(tool_name)
        if not cfg:
            return {"error": f"未知工具: {tool_name}"}

        breaker = self.breakers.get(tool_name)
        handlers = [(cfg.handler, f"{tool_name}(主)")]
        for i, fb in enumerate(cfg.fallback_handlers):
            handlers.append((fb, f"{tool_name}(降级{i+1})"))

        for handler, desc in handlers:
            try:
                if breaker:
                    return breaker.call(
                        self._retry_wrapper(
                            handler, cfg.max_retries, cfg.retry_delay
                        ),
                        *args, **kwargs
                    )
                else:
                    return self._retry_wrapper(
                        handler, cfg.max_retries, cfg.retry_delay
                    )(*args, **kwargs)
            except CircuitBreakerOpenError:
                continue  # 熔断中,跳到下一个处理器
            except Exception as e:
                logger.warning(f"[{tool_name}] {desc} 失败: {type(e).__name__}")
                continue

        logger.error(f"[{tool_name}] 所有处理器均失败,返回默认值")
        return cfg.default_result

    @staticmethod
    def _retry_wrapper(func, max_retries, delay):
        def wrapper(*args, **kwargs):
            last_err = None
            for attempt in range(max_retries + 1):
                try:
                    return func(*args, **kwargs)
                except (requests.Timeout, requests.ConnectionError) as e:
                    last_err = e
                    if attempt < max_retries:
                        wait = delay * (2 ** attempt) + random.uniform(0, 1)
                        time.sleep(wait)
            raise last_err
        return wrapper

Agent 使用时只需一行代码,完全感知不到内部的四层防护

# 初始化(启动时执行一次)
executor = ResilientToolExecutor()

executor.register(ToolConfig(
    name="web_search",
    handler=google_search,
    max_retries=2,
    fallback_handlers=[bing_search, local_cache_search],
    default_result={"results": [], "status": "degraded"},
    cb_threshold=3, cb_recovery=60,
))

executor.register(ToolConfig(
    name="llm_call",
    handler=call_claude,
    max_retries=2,
    fallback_handlers=[call_gpt4o, call_deepseek],
    default_result="AI 服务暂时不可用。",
    cb_threshold=5, cb_recovery=300,
))

# Agent 调用——简洁、安全
search_results = executor.execute("web_search", "AI Agent 框架对比 2026")
llm_answer = executor.execute("llm_call", "总结以下内容:\n" + text)

七、八个踩坑实录

这些都是我在生产环境中亲自踩过的坑。

坑 1:捕获异常后返回 None

# ❌ 致命错误
try:
    return api.search(q)
except:
    return None  # Agent 拿到 None 会干什么?不知道。

Agent 看到 None 的反应是不可预测的——可能编造数据(幻觉),可能直接报错,可能在 5 步之后才炸。永远返回结构化数据

坑 2:对所有异常都重试

# ❌ 浪费:401 重试 5 次还是 401
@retry_with_backoff(max_retries=5)
def bad_api_call():
    resp = requests.get(url)
    resp.raise_for_status()  # 401 也重试

只对瞬时故障重试:TimeoutConnectionError429。状态码 400-499 不应该重试。

坑 3:不使用 jitter

这是分布式系统里的经典陷阱。假设你有 50 个定时任务 Agent,它们在整点同时发起请求。如果都没加 jitter,第一次失败后 50 个 Agent 会在第 1 秒同时重试——等于对目标服务发起了一波 DDoS。

坑 4:熔断器阈值设太高

failure_threshold=20 的后果是——服务已经挂了 2 分钟,熔断器还在"积累失败次数"。经验值:搜索/爬虫类 3-5 次,LLM 类 5-8 次

坑 5:降级链太深

四级以上降级是设计失误。如果你备了 5 个搜索 API,说明你的主搜索 API 选错了。主路 + 2 路降级 = 最优

坑 6:不记日志

出了线上故障,打开日志发现只有"搜索失败"四个字——完全没法排查是哪个 query、什么错误码、发生在几点几分。每条日志至少包含:工具名、操作描述、异常类型、时间戳

坑 7:兜底值太敷衍

# ❌ 兜底值完全无用
default_result = "error"

Agent 看到 "error" 会怎么理解?没任何上下文。兜底值应该是完整的、包含说明的、Agent 能直接展示给用户的友好内容。

坑 8:熔断后不通知人类

熔断器默默打开了 30 分钟,开发者毫不知情。熔断触发应该发出告警(发微信消息 / 写入监控系统 / 触发 Webhook),让人知道服务出问题了。


七点五、异步场景的特殊处理

如果你的 Agent 用的是 async/await(比如 FastAPI 服务里的 Agent),上面的同步代码需要适配。异步场景下最容易犯的错误是——用 time.sleep() 阻塞了整个事件循环。

异步版重试装饰器

import asyncio
import random

def async_retry_with_backoff(
    max_retries: int = 3,
    base_delay: float = 1.0,
    max_delay: float = 30.0,
):
    """异步版指数退避重试装饰器"""
    def decorator(func):
        @wraps(func)
        async def wrapper(*args, **kwargs):
            last_err = None
            for attempt in range(max_retries + 1):
                try:
                    return await func(*args, **kwargs)
                except (asyncio.TimeoutError, ConnectionError) as e:
                    last_err = e
                    if attempt == max_retries:
                        raise
                    delay = min(base_delay * (2 ** attempt), max_delay)
                    jitter = random.uniform(0, delay * 0.5)
                    logger.warning(
                        f"{func.__name__} 异步重试 {attempt+1}/{max_retries}, "
                        f"等待 {delay+jitter:.1f}s"
                    )
                    await asyncio.sleep(delay + jitter)
            raise last_err
        return wrapper
    return decorator

关键区别:用 await asyncio.sleep() 代替 time.sleep()。如果用同步 sleep,整个事件循环会被阻塞,其他并发请求全部卡住——这在 Web 服务里是致命的。

异步熔断器

熔断器的状态变更在异步高并发场景下要注意竞态条件。多个协程可能同时检查 self.state == CircuitState.HALF_OPEN 并都认为自己可以发探测请求。生产环境建议用 asyncio.Lock 保护状态变更:

class AsyncCircuitBreaker:
    def __init__(self, *args, **kwargs):
        # ... 同同步版
        self._lock = asyncio.Lock()

    async def call(self, func, *args, **kwargs):
        async with self._lock:
            if self.state == CircuitState.OPEN:
                # ... 检查冷却时间,状态变更
                pass
        # 释放锁后再执行实际调用(避免长时间持锁)
        try:
            result = await func(*args, **kwargs)
            async with self._lock:
                self.state = CircuitState.CLOSED
                self.failure_count = 0
            return result
        except Exception as e:
            async with self._lock:
                self.failure_count += 1
                if self.failure_count >= self.failure_threshold:
                    self.state = CircuitState.OPEN
            raise

要点:锁只保护状态读写,不保护业务调用。如果锁住整个 func() 调用,高并发下会变成串行执行。


八、实战案例:凌晨三点的批处理任务救了我

说一段真实经历。我们有一个夜间批处理 Agent,每天凌晨 3 点自动运行,扫描 200+ 条数据源,调用 LLM 做摘要,最后推送到飞书群。

初始架构(裸奔版)

# 凌晨 3 点 cron 触发
for source in data_sources:
    raw = requests.get(source.url, timeout=10).json()  # 第 1 个炸弹
    summary = call_claude(raw["content"])                # 第 2 个炸弹
    post_to_feishu(summary)

上线第一周就出了三次事故:两次是某数据源 CDN 凌晨维护超时导致任务中断,一次是 Claude API 在凌晨 3:15 遇到限流,后面 150+ 条数据全没处理。

改造后(四层防护版)

executor = ResilientToolExecutor()

# 数据源抓取——3 次重试,失败记日志但不中断
executor.register(ToolConfig(
    name="fetch_source",
    handler=fetch_with_timeout,
    max_retries=3,
    fallback_handlers=[fetch_via_cdn_fallback],
    default_result={"content": "", "status": "fetch_failed"},
    cb_threshold=5,
))

# LLM 摘要——多模型降级
executor.register(ToolConfig(
    name="summarize",
    handler=call_claude,
    max_retries=2,
    fallback_handlers=[call_gpt4o, call_deepseek],
    default_result="[摘要生成失败,请查看原文]",
    cb_threshold=5,
    cb_recovery=600,  # 批处理场景,冷却 10 分钟
))

# 主循环——不会中断了
failed_count = 0
for i, source in enumerate(data_sources, 1):
    raw = executor.execute("fetch_source", source.url)
    if raw.get("status") == "fetch_failed":
        failed_count += 1
        continue  # 单条失败不中断整个批次

    summary = executor.execute("summarize", raw.get("content", ""))
    post_to_feishu(summary)

logger.info(f"批处理完成: {len(data_sources)} 条, 成功 {len(data_sources)-failed_count}, 失败 {failed_count}")

改造前后对比

指标 改造前 改造后
任务完成率 ~70%(经常中断) 99.8%(个别源失败不影响整体)
凌晨人工介入 每周 2-3 次 连续 3 周零介入
单条失败影响 整批中断 仅该条跳过
故障恢复时间 人工发现→修复→重跑(1-2 小时) 自动重试/降级(秒级)

这个案例的核心启示就是:Agent 任务应该像数据库事务一样,单条失败不回滚全局。你的 Agent 任务里如果有 200 步操作,第 199 步失败不应该让前 198 步的成果白费。

扩展:如何处理"部分成功"的结果

很多 Agent 任务不是"全或无"的。比如上面这个批处理,99 条成功 + 1 条失败,应该怎么汇报?

def batch_execute_with_report(items, executor, tool_name):
    """带结果汇总的批量执行"""
    results = {"success": [], "failed": [], "degraded": []}

    for item in items:
        try:
            result = executor.execute(tool_name, item)
            if result and not result.get("error"):
                results["success"].append({"item": item, "result": result})
            else:
                # 有结果但是降级返回的(比如用了 fallback)
                results["degraded"].append({"item": item, "result": result})
        except Exception as e:
            results["failed"].append({"item": item, "error": str(e)})

    return results


# 使用
batch_result = batch_execute_with_report(data_sources, executor, "summarize")

# 汇报给用户——清晰区分成功/降级/失败
report = f"""
📊 批处理完成
✅ 成功: {len(batch_result['success'])}
⚠️ 降级: {len(batch_result['degraded'])}
❌ 失败: {len(batch_result['failed'])}
"""

这样 Agent 汇报的不是"我失败了",而是"成功了 198 条,有 2 条用了备用方案,请检查"——用户就知道该关注什么,而不是一脸懵地看到报错堆栈。


九、总结

防护层 解决什么问题 一句话要诀
try/except 包装器 裸调用崩溃 永远返回结构化数据,永远不返回 None
指数退避重试 瞬时网络抖动 只重试可恢复错误,必须加 jitter
降级回退链 服务不可用 主路 + 2 路降级,兜底值必须有意义
熔断器 持续故障雪崩 快速失败比无限等待更好

这四层套上去之后,我们团队 6 个 Agent 的故障率变化:

  • 实施前:平均每天至少 1 次人工介入(重启任务 / 重跑)
  • 实施后:连续 14 天零人工介入

核心逻辑只有一句:让 Agent 专注思考,让框架负责保活


附录:快速检查清单

上生产前逐项确认:

  • [ ] 每个工具函数都包了 try/except,返回结构化结果
  • [ ] 重试只针对 Timeout / ConnectionError / 429,不过多
  • [ ] 重试加了随机 jitter,设了 max_delay 上限
  • [ ] 关键工具有至少 1 个降级方案
  • [ ] 兜底值是对用户友好的完整文本,不是 "error"
  • [ ] 高频调用工具启用了熔断器
  • [ ] 每次异常都打日志(工具名 + query + 异常类型)
  • [ ] 熔断触发时发送了告警通知

参考来源
- AWS Builders' Library: Timeouts, retries, and backoff with jitter
- Martin Fowler: Circuit Breaker Pattern
- Google SRE Book: Handling Overload
- Python: tenacity — retrying library