共计 2925 个字符,预计需要花费 8 分钟才能阅读完成。
问题现象与典型场景
最近在对接某 AI 平台的文本生成 API 时,经常遇到这样的报错:Agent terminated due to error you can prompt the model to try again or start。经过排查发现主要出现在以下场景:

- API 响应时间超过 30 秒被网关切断连接
- 突发流量导致的服务端限流(HTTP 429)
- 访问令牌(token)突然失效或配额耗尽
- 模型负载过高返回的临时错误(5xx 系列)
这类错误的特点是 瞬时性——可能下一秒重试就成功,因此实现自动重试机制是提升系统健壮性的关键。
重试策略对比
遇到错误立即重试?先看看三种常见策略的差异:
- 立即重试
- 优势:实现简单,代码量少
- 劣势:容易加剧服务端压力,形成雪崩效应
-
适用场景:明确知道是偶发错误(如网络抖动)
-
固定间隔重试
- 示例:每隔 2 秒重试一次
- 优势:节奏可控,便于调试
-
劣势:无法应对持续高负载场景
-
指数退避(Exponential Backoff)
- 延迟时间按指数增长:1s, 2s, 4s, 8s…
- 优势:能有效缓解服务端压力
- 劣势:实现复杂度较高
对于 AI 服务调用,推荐组合使用指数退避 + 随机抖动(jitter),既能平滑压力,又避免集群同时重试。
Python 实现方案
基础重试装饰器
from functools import wraps
import random
import time
from typing import Callable, TypeVar, Any
T = TypeVar('T')
def retry_with_backoff(
max_retries: int = 3,
initial_delay: float = 1.0,
max_delay: float = 10.0
) -> Callable[[Callable[..., T]], Callable[..., T]]:
"""
带指数退避的重试装饰器
:param max_retries: 最大重试次数
:param initial_delay: 初始延迟秒数
:param max_delay: 最大延迟秒数
"""
def decorator(func: Callable[..., T]) -> Callable[..., T]:
@wraps(func)
def wrapper(*args, **kwargs) -> T:
delay = initial_delay
for attempt in range(max_retries + 1):
try:
return func(*args, **kwargs)
except Exception as e:
if attempt == max_retries:
raise
# 过滤可重试错误类型
if not _is_retryable_error(e):
raise
# 计算退避时间并增加随机抖动
sleep_time = min(max_delay, delay * 2 ** attempt)
sleep_time *= random.uniform(0.8, 1.2) # ±20% 抖动
print(f'Attempt {attempt + 1} failed, retrying in {sleep_time:.2f}s')
time.sleep(sleep_time)
return wrapper
return decorator
def _is_retryable_error(error: Exception) -> bool:
"""判断错误是否可重试"""
# 示例:只重试特定错误码
if hasattr(error, 'status_code'):
return error.status_code in {429, 500, 502, 503, 504}
return False
进阶功能实现
状态持久化
对于长时间任务,需要记录已完成的进度。假设我们处理文本分块:
import json
import os
from pathlib import Path
class StateManager:
def __init__(self, state_file: str = 'progress.json'):
self.state_file = Path(state_file)
def save_progress(self, last_processed_id: str):
"""保存最后处理成功的 ID"""
with open(self.state_file, 'w') as f:
json.dump({'last_id': last_processed_id}, f)
def load_progress(self) -> str:
"""加载进度,不存在返回空字符串"""
if not self.state_file.exists():
return ""
with open(self.state_file) as f:
return json.load(f).get('last_id', "")
分布式锁
防止多进程同时处理相同任务(使用文件锁示例):
import fcntl
class DistributedLock:
def __init__(self, lock_file: str):
self.lock_file = Path(lock_file)
self.fd = None
def __enter__(self):
self.fd = open(self.lock_file, 'w')
fcntl.flock(self.fd, fcntl.LOCK_EX)
def __exit__(self, exc_type, exc_val, exc_tb):
if self.fd:
fcntl.flock(self.fd, fcntl.LOCK_UN)
self.fd.close()
性能与可靠性设计
关键指标监控
建议采集以下指标用于分析:
- 重试成功率:成功重试次数 / 总重试次数
- 平均延迟时间:所有重试等待时间的平均值
- 错误类型分布:统计各类错误的出现频率
参数调优建议
根据监控数据调整:
- 初始延迟:从 1 秒开始,如果发现大量 429 错误则调大
- 最大重试次数:
- 对实时性要求高的场景:3- 5 次
- 离线批处理任务:可设置 8 -10 次
- 退避上限:通常不超过 30 秒,否则影响用户体验
常见避坑指南
无限重试预防
必须设置最大重试次数,并考虑以下情况:
# 在装饰器中加入超时控制
timeout = 60 * 5 # 5 分钟总超时
start_time = time.time()
while time.time() - start_time < timeout:
try:
return func(*args, **kwargs)
except Exception:
...
else:
raise TimeoutError("Operation timed out after retries")
非幂等操作处理
对于创建订单这类非幂等操作:
- 先查询是否已存在
- 使用唯一 IDempotency-Key
- 服务端实现去重逻辑
分布式竞争条件
- 使用 Redis 锁替代文件锁
- 对关键资源采用 CAS(Compare-And-Swap)操作
- 设计任务分片避免冲突
扩展思考
当重试达到阈值后,系统可以:
- 人工干预:发送告警邮件 / 短信,保留错误上下文
- 降级处理:返回缓存内容或简化版结果
- 任务转移:将请求转移到备用服务节点
具体选择取决于业务场景——实时对话系统可能需要立即降级,而数据分析任务可能更适合等待人工排查。
在你的项目中,更倾向于哪种失败处理策略?欢迎在评论区分享实践经验。
正文完
