Skip to content

12.5 稳定性与可靠性

上一节我们构建了多 Agent 协作系统——多个 Agent 分工配合,能力更强了。但"能干"和"能扛"是两回事。一个 Agent 系统在实验室里跑得好,不代表在生产环境里扛得住。真实环境有网络波动、API 限流、LLM 幻觉、并发竞争——任何一个环节出问题都可能让整个系统崩溃。

本节就来回答一个关键问题:如何让 Agent 系统从"能跑"升级到"能扛"?我们从不确定性管理、幂等性设计、监控告警、故障恢复四个维度来构建可靠性保障体系。

生活类比:航空公司运营

想象一家航空公司的运营。飞机起飞前要做全面检查(健康检查),飞行中雷达实时监控各项指标(监控告警),遇到气流颠簸有自动稳定系统(不确定性管理),同一航班卖出的票不会重复出票(幂等性),发动机故障了有备用发动机和备降方案(故障恢复与降级)。

Agent 系统的可靠性保障和航空公司惊人地相似。它需要做到四件事:预测和管理不确定性(LLM 输出不靠谱怎么办)、保证操作幂等(重复请求不会重复扣款)、全链路监控(哪里慢了哪里错了立刻知道)、故障自动恢复(挂了能自动拉起来、能降级保核心功能)。下面我们逐个拆解。

一、不确定性的来源与管理

AI Agent 系统的核心挑战在于 LLM 本身的非确定性。同一个问题问两次,可能得到不同的答案。这在聊天场景可以接受,但在涉及金额、日期、订单等关键信息时就不可接受了。

不确定性来源有五类:LLM 输出不确定(同一输入不同输出)、工具执行不确定(API 偶尔超时)、网络波动(连接不稳定)、外部 API 不稳定(第三方服务挂了)、并发竞争(多个请求同时修改数据)。UncertaintyManager 对每类不确定性都定义了对应策略:

python
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 通过唯一键来解决:

python
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 的每次执行都自动走幂等检查:

python
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 负责告警判断:

python
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)

四、故障恢复与降级策略

监控告警让你"知道"出了问题,但知道还不够,还要"能恢复"。故障恢复有三个核心组件:熔断器(连续失败时自动断开,防止雪崩)、重试策略(临时故障时指数退避重试)、优雅降级(核心功能不可用时退而求其次)。

python
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

五、常见误区

  1. 误区一:熔断器只看失败次数,不看时间。有人只在代码里记 failure_count,但不设 recovery_timeout——熔断器一旦打开就永远不恢复。必须设置恢复超时,让熔断器在一段时间后进入半开状态试探。

  2. 误区二:重试不加抖动(jitter)。如果 100 个请求同时失败,不加抖动它们会在同一时刻同时重试,造成"重试风暴"。jitter = True 让每个请求的等待时间有随机偏移,错开重试时间。

  3. 误区三:幂等键只用用户 ID。如果只用 user_id 做幂等键,用户先下了个订单,再查个数据,两个请求的幂等键一样,查数据会返回订单结果。必须把请求参数也纳入幂等键计算。

  4. 误区四:监控只记不报。有些团队采集了一堆指标但不设告警规则,等用户投诉了才发现系统挂了。_setup_default_alerts 中定义的四条核心告警(延迟、错误率、工具失败率、队列深度)是最基本的要求。

  5. 误区五:降级策略只有"全有"和"全无"。有人认为服务要么正常要么挂掉,没有中间态。实际上 GracefulDegradation 支持分级降级——level 1 用简化版功能,level 2 返回最小服务,保证核心流程不中断。

六、本节小结

要点说明
不确定性管理结构化输出验证、一致性检查、多次采样对 LLM 输出进行约束
幂等性设计基于请求内容的唯一键,确保重复请求结果一致
监控告警延迟、错误率、吞吐量、队列深度全链路采集,阈值+速率双重告警
熔断器CLOSED -> OPEN -> HALF_OPEN 状态转换,防止级联失败
优雅降级分级降级策略,从完整功能到最小服务,确保核心可用

稳定性与可靠性是 Agent 系统从"能跑"到"能扛"的分水岭。一个在实验室里 Demo 惊艳的系统,如果扛不住生产环境的网络波动、API 限流、并发竞争,就不能算真正落地。掌握了不确定性管理、幂等性、监控告警、熔断降级这套组合拳,你的 Agent 系统才能在真实世界稳稳运行。下一节我们精选几道高频面试场景题,把前面的知识串起来、用出去。