12.5 稳定性与可靠性
上一节我们构建了多 Agent 协作系统——多个 Agent 分工配合,能力更强了。但"能干"和"能扛"是两回事。一个 Agent 系统在实验室里跑得好,不代表在生产环境里扛得住。真实环境有网络波动、API 限流、LLM 幻觉、并发竞争——任何一个环节出问题都可能让整个系统崩溃。
本节就来回答一个关键问题:如何让 Agent 系统从"能跑"升级到"能扛"?我们从不确定性管理、幂等性设计、监控告警、故障恢复四个维度来构建可靠性保障体系。
生活类比:航空公司运营
想象一家航空公司的运营。飞机起飞前要做全面检查(健康检查),飞行中雷达实时监控各项指标(监控告警),遇到气流颠簸有自动稳定系统(不确定性管理),同一航班卖出的票不会重复出票(幂等性),发动机故障了有备用发动机和备降方案(故障恢复与降级)。
Agent 系统的可靠性保障和航空公司惊人地相似。它需要做到四件事:预测和管理不确定性(LLM 输出不靠谱怎么办)、保证操作幂等(重复请求不会重复扣款)、全链路监控(哪里慢了哪里错了立刻知道)、故障自动恢复(挂了能自动拉起来、能降级保核心功能)。下面我们逐个拆解。
一、不确定性的来源与管理
AI Agent 系统的核心挑战在于 LLM 本身的非确定性。同一个问题问两次,可能得到不同的答案。这在聊天场景可以接受,但在涉及金额、日期、订单等关键信息时就不可接受了。
不确定性来源有五类:LLM 输出不确定(同一输入不同输出)、工具执行不确定(API 偶尔超时)、网络波动(连接不稳定)、外部 API 不稳定(第三方服务挂了)、并发竞争(多个请求同时修改数据)。UncertaintyManager 对每类不确定性都定义了对应策略:
from typing import Dict, Any, Optional, List # 类型注解
from dataclasses import dataclass, field # 数据类
from enum import Enum # 枚举类型
import hashlib # 哈希
import json # JSON处理
import time # 时间模块
class UncertaintyType(Enum): # 不确定性类型枚举
LLM_OUTPUT = "llm_output" # LLM输出不确定
TOOL_EXECUTION = "tool_execution" # 工具执行结果不确定
NETWORK = "network" # 网络波动
EXTERNAL_API = "external_api" # 外部API不稳定
CONCURRENCY = "concurrency" # 并发竞争
class UncertaintyManager: # 不确定性管理器
"""不确定性管理器""" # 对每类不确定性定义策略
def __init__(self): # 初始化
self.strategies = { # 类型->处理策略映射
UncertaintyType.LLM_OUTPUT: self._handle_llm_uncertainty, # LLM不确定性
UncertaintyType.TOOL_EXECUTION: self._handle_tool_uncertainty, # 工具不确定性
UncertaintyType.NETWORK: self._handle_network_uncertainty, # 网络不确定性
UncertaintyType.EXTERNAL_API: self._handle_external_uncertainty, # 外部API不确定性
}
def _handle_llm_uncertainty(self, response: str, expected_schema: Dict) -> Dict:
"""处理LLM输出不确定性""" # 验证+置信度评估
is_valid, error = self._validate_output(response, expected_schema) # 验证格式
if not is_valid: # 格式不对
return {"status": "retry", "error": error, "fix_strategy": "ask_llm_to_fix"} # 要求重试
return {"status": "success", "data": response, "confidence": self._calculate_confidence(response)} # 成功
def _validate_output(self, response: str, schema: Dict) -> tuple[bool, str]:
"""验证输出是否符合预期格式""" # 结构化校验
if schema.get("type") == "json": # 要求JSON格式
try:
data = json.loads(response) # 尝试解析JSON
for field in schema.get("required", []): # 检查必需字段
if field not in data: # 缺少字段
return False, f"缺少必需字段: {field}"
return True, "" # 验证通过
except json.JSONDecodeError as e: # JSON解析失败
return False, f"JSON解析失败: {str(e)}"
return True, "" # 非JSON格式不校验
def _calculate_confidence(self, response: str) -> float:
"""计算输出置信度""" # 简化版:基于长度和结构
score = 0.5 # 基础分
if len(response) > 50: # 内容够长
score += 0.2 # 加分
if '{' in response and '}' in response: # 有JSON结构
score += 0.2 # 加分
return min(score, 1.0) # 不超过1.0
def _handle_tool_uncertainty(self, result: Any, tool_name: str) -> Dict:
"""处理工具执行不确定性""" # 检查返回值
if result is None: # 返回空
return {"status": "retry", "error": "工具返回空结果"} # 要求重试
if isinstance(result, dict) and result.get("error"): # 返回错误
return {"status": "retry", "error": result["error"]} # 要求重试
return {"status": "success", "data": result} # 成功
def _handle_network_uncertainty(self, error: Exception) -> Dict:
"""处理网络不确定性""" # 指数退避重试
return {"status": "retry", "error": str(error), "strategy": "exponential_backoff"}
def _handle_external_uncertainty(self, error: Exception) -> Dict:
"""处理外部API不确定性""" # 降级到缓存或默认值
return {"status": "fallback", "error": str(error), "fallback_action": "use_cached_or_default"}关键设计思路是差异化策略:LLM 输出不对就要求重试并让 LLM 修复;工具返回空就重试;网络波动用指数退避;外部 API 挂了就降级到缓存。不同问题用不同药方,而不是一刀切地重试。
二、幂等性设计
幂等性是 Agent 系统可靠性的基石。幂等的意思是:同一个请求执行一次和执行多次,结果完全一致。就像你按电梯按钮——按一次和连按十次,电梯来的都是同一趟。
为什么幂等性重要?因为分布式系统中请求可能因为网络超时而重发——用户点了"提交订单",网络卡了一下,前端自动重试,如果系统不是幂等的,就会创建两个订单。IdempotencyManager 通过唯一键来解决:
class IdempotencyManager: # 幂等性管理器
"""幂等性管理器""" # 确保重复请求不重复执行
def __init__(self, storage_backend=None): # 初始化
self.storage = storage_backend or {} # 存储后端,生产用Redis
self.key_prefix = "idempotent:" # 键前缀
self.ttl = 86400 # 24小时过期
def generate_key(self, *args, **kwargs) -> str:
"""生成幂等键""" # 基于请求内容生成唯一键
content = json.dumps({"args": args, "kwargs": kwargs}, sort_keys=True) # 序列化参数
return self.key_prefix + hashlib.sha256(content.encode()).hexdigest()[:16] # 取前16位哈希
def check_or_create(self, idempotent_key: str) -> tuple[bool, Optional[Any]]:
"""检查是否已处理过""" # 返回(是否重复, 缓存结果)
if idempotent_key in self.storage: # 键已存在
record = self.storage[idempotent_key] # 取记录
if record["status"] == "completed": # 已完成
return True, record["result"] # 返回缓存结果
elif record["status"] == "processing": # 正在处理中
return True, None # 等待结果
self.storage[idempotent_key] = { # 创建处理记录
"status": "processing", # 标记为处理中
"created_at": time.time(), # 创建时间
"result": None # 结果暂空
}
return False, None # 不是重复请求
def mark_completed(self, idempotent_key: str, result: Any):
"""标记处理完成""" # 存结果,后续重复请求直接返回
if idempotent_key in self.storage: # 键存在
self.storage[idempotent_key] = { # 更新记录
"status": "completed", # 标记完成
"completed_at": time.time(), # 完成时间
"result": result # 存结果
}
def mark_failed(self, idempotent_key: str, error: str):
"""标记处理失败""" # 允许后续重试
if idempotent_key in self.storage: # 键存在
self.storage[idempotent_key] = { # 更新记录
"status": "failed", # 标记失败
"failed_at": time.time(), # 失败时间
"error": error # 错误信息
}IdempotentAgent 把幂等管理器包装成一个装饰器,Agent 的每次执行都自动走幂等检查:
class IdempotentAgent: # 幂等的Agent包装器
"""幂等的Agent包装器""" # 给任意Agent加幂等能力
def __init__(self, agent, idempotency: IdempotencyManager): # 初始化
self.agent = agent # 被包装的Agent
self.idempotency = idempotency # 幂等管理器
async def execute(self, action: str, **params) -> Dict:
"""幂等执行Agent动作""" # 先查再执行
key = self.idempotency.generate_key(action, **params) # 生成幂等键
is_dup, cached = self.idempotency.check_or_create(key) # 检查是否重复
if is_dup: # 是重复请求
if cached is not None: # 有缓存结果
return {"status": "cached", "data": cached} # 直接返回
else: # 正在处理中
return {"status": "processing", "message": "请求正在处理中"}
try:
result = await self.agent.execute(action, **params) # 执行Agent
self.idempotency.mark_completed(key, result) # 标记完成
return {"status": "success", "data": result}
except Exception as e: # 执行失败
self.idempotency.mark_failed(key, str(e)) # 标记失败
raise # 重新抛出三、监控告警体系
幂等解决了"重复执行"的问题,但怎么知道系统"哪里出问题了"?答案是监控告警。就像飞机驾驶舱的仪表盘——油量低了亮黄灯、发动机温度高了亮红灯,飞行员不用去摸发动机,看仪表盘就知道。
Agent 系统的监控指标分四类:延迟(LLM 响应慢不慢)、错误率(调用失败比例)、吞吐量(每秒处理多少请求)、队列深度(积压了多少任务)。AgentMetrics 负责采集,AgentMonitor 负责告警判断:
import logging # 日志模块
from collections import defaultdict # 默认字典
from datetime import datetime, timedelta # 时间处理
class AgentMetrics: # Agent监控指标
"""Agent监控指标""" # 采集和查询时序数据
def __init__(self): # 初始化
self.metrics = defaultdict(list) # 指标名->值列表
self.alert_rules = [] # 告警规则列表
def record(self, metric_name: str, value: float, tags: Dict = None):
"""记录指标""" # 每次调用都记录
self.metrics[metric_name].append({ # 追加到指标列表
"value": value, # 指标值
"timestamp": time.time(), # 时间戳
"tags": tags or {} # 标签
})
def get_window(self, metric_name: str, window_seconds: int = 300) -> List[float]:
"""获取时间窗口内的指标值""" # 只看最近N秒的数据
now = time.time() # 当前时间
cutoff = now - window_seconds # 截止时间
return [ # 过滤并提取值
m["value"] for m in self.metrics[metric_name]
if m["timestamp"] >= cutoff
]
def check_alert(self, metric_name: str, rule: Dict) -> Optional[str]:
"""检查告警规则""" # 看指标是否触发告警
window = self.get_window(metric_name, rule.get("window", 300)) # 取窗口数据
if not window: # 没有数据
return None # 不告警
if rule["type"] == "threshold": # 阈值告警
threshold = rule["threshold"] # 取阈值
if rule.get("direction", "above") == "above": # 超过阈值
if max(window) > threshold: # 最大值超标
return f"{metric_name} 超过阈值 {threshold},当前值 {max(window)}"
else: # 低于阈值
if min(window) < threshold: # 最小值超标
return f"{metric_name} 低于阈值 {threshold},当前值 {min(window)}"
elif rule["type"] == "rate": # 速率告警
rate = sum(window) / len(window) / rule.get("window", 300) # 计算速率
if rate > rule.get("max_rate", float('inf')): # 超过最大速率
return f"{metric_name} 速率过高:{rate:.2f}/s"
return None # 没触发告警
class AgentMonitor: # Agent监控器
"""Agent监控器""" # 统一采集+告警+处理
def __init__(self): # 初始化
self.metrics = AgentMetrics() # 指标采集器
self.logger = logging.getLogger("AgentMonitor") # 日志器
self.alert_handlers = [] # 告警处理器列表
self._setup_default_alerts() # 设置默认告警规则
def _setup_default_alerts(self):
"""设置默认告警规则""" # 四条核心告警
self.default_alerts = [
{"metric": "llm_latency", # LLM响应延迟
"rule": {"type": "threshold", "threshold": 30, "direction": "above"},
"severity": "warning", # 警告级别
"message": "LLM响应延迟超过30秒"},
{"metric": "error_rate", # 错误率
"rule": {"type": "rate", "max_rate": 0.1},
"severity": "critical", # 严重级别
"message": "错误率超过10%"},
{"metric": "tool_failure_rate", # 工具失败率
"rule": {"type": "threshold", "threshold": 0.05},
"severity": "critical",
"message": "工具调用失败率超过5%"},
{"metric": "queue_depth", # 队列深度
"rule": {"type": "threshold", "threshold": 100},
"severity": "warning",
"message": "任务队列积压超过100"}
]
def record_llm_call(self, latency: float, model: str, success: bool):
"""记录LLM调用""" # 每次LLM调用后记录
self.metrics.record("llm_latency", latency, {"model": model}) # 记延迟
self.metrics.record("llm_calls", 1, {"model": model, "success": success}) # 记次数
if not success: # 调用失败
self.metrics.record("llm_errors", 1, {"model": model}) # 记错误
def record_tool_call(self, tool_name: str, latency: float, success: bool):
"""记录工具调用""" # 每次工具调用后记录
self.metrics.record("tool_latency", latency, {"tool": tool_name}) # 记延迟
self.metrics.record("tool_calls", 1, {"tool": tool_name, "success": success}) # 记次数
if not success: # 调用失败
self.metrics.record("tool_failures", 1, {"tool": tool_name}) # 记失败
def check_alerts(self) -> List[Dict]:
"""检查并触发告警""" # 定期调用
alerts = [] # 收集告警
for alert_config in self.default_alerts: # 遍历告警规则
alert_msg = self.metrics.check_alert( # 检查是否触发
alert_config["metric"],
alert_config["rule"]
)
if alert_msg: # 触发了
alert = { # 构造告警对象
"metric": alert_config["metric"],
"severity": alert_config["severity"],
"message": alert_config["message"],
"detail": alert_msg,
"timestamp": time.time()
}
alerts.append(alert) # 加入列表
for handler in self.alert_handlers: # 遍历处理器
handler(alert) # 触发处理(发钉钉/邮件等)
return alerts # 返回所有告警
def add_alert_handler(self, handler: callable):
"""添加告警处理器""" # 如发送钉钉/飞书/邮件通知
self.alert_handlers.append(handler)四、故障恢复与降级策略
监控告警让你"知道"出了问题,但知道还不够,还要"能恢复"。故障恢复有三个核心组件:熔断器(连续失败时自动断开,防止雪崩)、重试策略(临时故障时指数退避重试)、优雅降级(核心功能不可用时退而求其次)。
import asyncio # 异步IO
from functools import wraps # 装饰器工具
class CircuitBreaker: # 熔断器
"""熔断器""" # 连续失败时断开,防止级联雪崩
def __init__(self, name: str, failure_threshold: int = 5,
recovery_timeout: float = 60.0, half_open_max: int = 3): # 初始化
self.name = name # 熔断器名称
self.failure_threshold = failure_threshold # 失败阈值,连续5次失败就熔断
self.recovery_timeout = recovery_timeout # 恢复超时,60秒后尝试半开
self.half_open_max = half_open_max # 半开状态最大试探次数
self.failure_count = 0 # 当前失败计数
self.half_open_count = 0 # 半开状态成功计数
self.last_failure_time = 0 # 最后一次失败时间
self.state = "CLOSED" # 状态:CLOSED/OPEN/HALF_OPEN
def call(self, func): # 装饰器入口
"""熔断器装饰器""" # 包装异步函数
@wraps(func)
async def wrapper(*args, **kwargs): # 包装函数
if self.state == "OPEN": # 熔断器打开
if time.time() - self.last_failure_time > self.recovery_timeout: # 超过恢复时间
self.state = "HALF_OPEN" # 切到半开,试探
self.half_open_count = 0 # 重置计数
else: # 还没到恢复时间
raise CircuitBreakerOpenError(f"熔断器 {self.name} 已打开") # 直接拒绝
try:
result = await func(*args, **kwargs) # 执行函数
self._on_success() # 成功处理
return result # 返回结果
except Exception as e: # 执行失败
self._on_failure() # 失败处理
raise # 重新抛出异常
return wrapper
def _on_success(self): # 成功时的处理
if self.state == "HALF_OPEN": # 半开状态成功
self.half_open_count += 1 # 成功计数+1
if self.half_open_count >= self.half_open_max: # 连续成功达标
self.state = "CLOSED" # 恢复正常
self.failure_count = 0 # 清零失败计数
else: # 正常状态成功
self.failure_count = 0 # 清零失败计数
def _on_failure(self): # 失败时的处理
self.failure_count += 1 # 失败计数+1
self.last_failure_time = time.time() # 记录失败时间
if self.failure_count >= self.failure_threshold: # 达到失败阈值
self.state = "OPEN" # 打开熔断器
class CircuitBreakerOpenError(Exception): # 熔断器打开异常
pass
class RetryPolicy: # 重试策略
"""重试策略""" # 指数退避+随机抖动
@staticmethod
async def exponential_backoff(func, max_retries: int = 3, base_delay: float = 1.0,
max_delay: float = 60.0, jitter: bool = True):
"""指数退避重试""" # 每次等待时间翻倍
last_exception = None # 记录最后一次异常
for attempt in range(max_retries + 1): # 最多重试max_retries次
try:
return await func() # 执行函数
except Exception as e: # 执行失败
last_exception = e # 记录异常
if attempt == max_retries: # 已经是最后一次
break # 不再重试
delay = min(base_delay * (2 ** attempt), max_delay) # 计算延迟,指数增长
if jitter: # 加随机抖动
import random # 防止所有请求同时重试
delay *= (0.5 + random.random()) # 0.5~1.5倍随机
await asyncio.sleep(delay) # 等待
raise last_exception # 重试耗尽,抛出异常
class GracefulDegradation: # 优雅降级
"""优雅降级""" # 核心功能不可用时退而求其次
def __init__(self): # 初始化
self.fallbacks = {} # 降级函数表
self.degradation_level = 0 # 降级级别:0正常/1降级/2最小服务
def register_fallback(self, service: str, fallback_func: callable):
"""注册降级函数""" # 为每个服务注册兜底方案
self.fallbacks[service] = {
"level_1": fallback_func, # 轻量降级:用简化版功能
"level_2": lambda: {"status": "degraded", "message": "服务暂时不可用"} # 最小服务
}
async def execute_with_fallback(self, service: str, primary_func: callable) -> Any:
"""带降级的执行""" # 先试主功能,失败走降级
try:
return await primary_func() # 先执行主功能
except Exception: # 主功能失败
if service in self.fallbacks: # 有降级方案
level = self.degradation_level # 当前降级级别
fallback = self.fallbacks[service].get( # 取对应级别的降级函数
f"level_{level+1}",
self.fallbacks[service]["level_2"] # 默认用最小服务
)
return await fallback() # 执行降级
raise # 没有降级方案就抛出
class HealthChecker: # 健康检查器
"""健康检查器""" # 定期检查各组件是否存活
def __init__(self): # 初始化
self.checks = {} # 检查项注册表
self.health_status = {} # 健康状态缓存
def register_check(self, name: str, check_func: callable, interval: float = 30):
"""注册健康检查""" # 每个组件注册一个检查函数
self.checks[name] = { # 存入检查项
"func": check_func, # 检查函数
"interval": interval, # 检查间隔(秒)
"last_check": 0 # 上次检查时间
}
async def check_all(self) -> Dict[str, str]:
"""执行所有健康检查""" # 定期调用
now = time.time() # 当前时间
results = {} # 结果字典
for name, check_info in self.checks.items(): # 遍历所有检查项
if now - check_info["last_check"] < check_info["interval"]: # 没到检查时间
results[name] = self.health_status.get(name, "unknown") # 返回缓存状态
continue
try:
healthy = await check_info["func"]() # 执行检查
status = "healthy" if healthy else "unhealthy" # 判断状态
except Exception: # 检查本身出错
status = "unhealthy"
self.health_status[name] = status # 缓存状态
check_info["last_check"] = now # 更新检查时间
results[name] = status # 存入结果
return results # 返回所有状态
@property
def is_healthy(self) -> bool: # 整体是否健康
"""整体是否健康""" # 所有组件都健康才健康
return all(s == "healthy" for s in self.health_status.values()) if self.health_status else True五、常见误区
误区一:熔断器只看失败次数,不看时间。有人只在代码里记
failure_count,但不设recovery_timeout——熔断器一旦打开就永远不恢复。必须设置恢复超时,让熔断器在一段时间后进入半开状态试探。误区二:重试不加抖动(jitter)。如果 100 个请求同时失败,不加抖动它们会在同一时刻同时重试,造成"重试风暴"。
jitter = True让每个请求的等待时间有随机偏移,错开重试时间。误区三:幂等键只用用户 ID。如果只用
user_id做幂等键,用户先下了个订单,再查个数据,两个请求的幂等键一样,查数据会返回订单结果。必须把请求参数也纳入幂等键计算。误区四:监控只记不报。有些团队采集了一堆指标但不设告警规则,等用户投诉了才发现系统挂了。
_setup_default_alerts中定义的四条核心告警(延迟、错误率、工具失败率、队列深度)是最基本的要求。误区五:降级策略只有"全有"和"全无"。有人认为服务要么正常要么挂掉,没有中间态。实际上
GracefulDegradation支持分级降级——level 1 用简化版功能,level 2 返回最小服务,保证核心流程不中断。
六、本节小结
| 要点 | 说明 |
|---|---|
| 不确定性管理 | 结构化输出验证、一致性检查、多次采样对 LLM 输出进行约束 |
| 幂等性设计 | 基于请求内容的唯一键,确保重复请求结果一致 |
| 监控告警 | 延迟、错误率、吞吐量、队列深度全链路采集,阈值+速率双重告警 |
| 熔断器 | CLOSED -> OPEN -> HALF_OPEN 状态转换,防止级联失败 |
| 优雅降级 | 分级降级策略,从完整功能到最小服务,确保核心可用 |
稳定性与可靠性是 Agent 系统从"能跑"到"能扛"的分水岭。一个在实验室里 Demo 惊艳的系统,如果扛不住生产环境的网络波动、API 限流、并发竞争,就不能算真正落地。掌握了不确定性管理、幂等性、监控告警、熔断降级这套组合拳,你的 Agent 系统才能在真实世界稳稳运行。下一节我们精选几道高频面试场景题,把前面的知识串起来、用出去。