你的 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 次尝试)
四个关键设计点:
- 只重试可恢复的错误:Timeout、ConnectionError、429 值得重试;400、401、403 重试毫无意义
- 必须加 jitter:不加随机抖动的后果是——凌晨批量任务触发限流时,所有 Agent 几乎同时重试,形成二次冲击波
- 设 max_delay 上限:生产环境不要让退避时间无限膨胀,30 秒是合理上限
- 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 也重试
只对瞬时故障重试:Timeout、ConnectionError、429。状态码 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
