AI Agent 架构设计与多 Agent 协作系统搭建:模型输出异常时的降级边界

1. 生产环境的隐患:模型突然返回空 JSON 或直接超时

并发请求叠加上游模型抖动时,Agent 工作流要先处理超时、格式错误和回退路径。

当线上 Agent 工作流在处理批量报表解析时,如果最上游的 LLM 调用接口大量报出 HTTP 504 响应,检查日志可以发现 HTTP 建连正常,但数据包在 API 网关处被卡住,Read Timeout 持续攀升。部分未超时的请求甚至可能返回包含 JSON 截断的字符串,导致后方解析器抛出 json.decoder.JSONDecodeError 异常,挂起整个任务队列。

内存中的 Goroutine 和 Python 异步 Task 积压瞬间增加,服务器 Load Average 会迅速上升。

网络抖动、上游故障和格式不合约定,都应在上线前演练。简单地“报错就重试三次”在持续超时时常会放大压力:请求继续堆积,网关和队列的恢复会更慢。

系统需要清楚的降级闸门。模型输出可能不稳定,但每一种失败都应有可测试的处理路径。


2. 三层降级防线:从网络超时到语义兜底

为了在 LLM 出错时保障核心业务不断流,系统必须在架构设计层面建立三层物理防线。不能把希望寄托在 API 供应商的 SLA 承诺上,工程实现需要考虑最坏情况下的保底手段。

第一层是网络与时间闸门。模型请求需要连接与读取超时,并配合熔断器。具体时长、错误率和连续失败次数要从端到端时延预算、上游配额和压测结果反推;达到阈值后暂时停止远端调用,转入可解释的回退路径。

第二层是格式自我修复与 Prompt 降级。当模型返回的 JSON 缺失字段或语法报错时,优先使用轻量级本地正则或修复器(Auto-repair)救场;若救场失败,迅速切换为更简短、降低 Output Complexity 的备用 Prompt(降级 Prompt),或者调用参数量更小、响应更快的备选模型(例如从 GPT-4o 降级至 GPT-4o-mini 或本地部署的 DeepSeek-R1 8B 蒸馏版)。

第三层是规则引擎兜底。当备用模型依然无法工作、或者熔断器处于打开状态时,降级引擎直接拦截 LLM 调用,由本地准备好的规则引擎、静态模板或历史缓存(Semantic Cache)返回保底结果,确保下游业务拿到安全的默认数据。

在实际落地过程中,降级链路的演进不能只靠人工判断。必须为每个降级节点注入清晰的状态上下文,这样当运维和开发人员查看 Trace 日志时,能明确看清这次请求究竟是在哪一道防线拦截救回来的。


3. 生产级 Python 降级引擎实现

下面的 Python 实现演示了如何在不依赖 LangChain 或 LlamaIndex 等高层抽象库的前提下,手写一套具备熔断、超时控制、格式修复与降级兜底的 Agent 处理器。代码包含完整的线程安全熔断器、正则表达式清洗以及备用模型切分逻辑。

import json
import time
import re
import logging
from typing import Dict, Any, Callable, Optional

logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")

class CircuitBreakerOpenException(Exception):
    """熔断器打开异常"""
    pass

class CircuitBreaker:
    """生产级熔断器,支持 CLOSED、OPEN、HALF-OPEN 状态转换"""
    def __init__(self, failure_threshold: int = 3, recovery_time: float = 30.0):
        self.failure_threshold = failure_threshold
        self.recovery_time = recovery_time
        self.failure_count = 0
        self.last_failure_time = 0.0
        self.state = "CLOSED"  # CLOSED, OPEN, HALF-OPEN

    def allow_request(self) -> bool:
        now = time.time()
        if self.state == "OPEN":
            if now - self.last_failure_time > self.recovery_time:
                self.state = "HALF-OPEN"
                logging.warning("熔断器进入 HALF-OPEN 状态,尝试放行探针请求")
                return True
            return False
        return True

    def record_success(self):
        self.failure_count = 0
        self.state = "CLOSED"

    def record_failure(self):
        self.failure_count += 1
        self.last_failure_time = time.time()
        if self.failure_count >= self.failure_threshold:
            self.state = "OPEN"
            logging.error(f"连续失败 {self.failure_count} 次,熔断器切换至 OPEN 状态!")


class ResilientAgentEngine:
    """高可用 Agent 执行引擎,提供确定性降级防线"""
    def __init__(self, primary_llm: Callable, fallback_llm: Callable):
        self.primary_llm = primary_llm
        self.fallback_llm = fallback_llm
        self.breaker = CircuitBreaker(failure_threshold=3, recovery_time=20.0)

    def _clean_json_string(self, text: str) -> str:
        """剥离 Markdown 代码块标记,提取有效 JSON"""
        match = re.search(r'```json\s*(\{.*?\})\s*```', text, re.DOTALL)
        if match:
            return match.group(1)
        match_raw = re.search(r'\{.*\}', text, re.DOTALL)
        if match_raw:
            return match_raw.group(0)
        return text

    def _try_repair_json(self, raw_text: str, required_keys: list) -> Optional[Dict[str, Any]]:
        """自动修复器:尝试清洗并补齐缺省字段"""
        cleaned = self._clean_json_string(raw_text)
        try:
            data = json.loads(cleaned)
            for key in required_keys:
                if key not in data:
                    data[key] = "UNKNOWN_DEGRADED"
            return data
        except json.JSONDecodeError:
            return None

    def execute_task(self, prompt: str, required_keys: list) -> Dict[str, Any]:
        """执行 Agent 任务,包含完整降级链路"""
        # 1. 检查熔断状态
        if not self.breaker.allow_request():
            logging.warning("熔断器拦截请求,直接触发规则兜底")
            return self._static_rule_fallback(prompt, reason="Circuit breaker open")

        # 2. 尝试主模型调用
        try:
            logging.info("正在调用主 LLM API...")
            raw_response = self.primary_llm(prompt, timeout=5.0)
            parsed = self._try_repair_json(raw_response, required_keys)
            if parsed:
                self.breaker.record_success()
                parsed["_execution_meta"] = {"source": "primary_llm", "degraded": False}
                return parsed
            else:
                logging.warning("主 LLM 返回格式解析失败,尝试修复无效")
                self.breaker.record_failure()
        except Exception as e:
            logging.error(f"主 LLM 调用异常: {str(e)}")
            self.breaker.record_failure()

        # 3. 降级到备用轻量模型
        try:
            logging.info("切换至备用 LLM 进行降级尝试...")
            fallback_prompt = f"以极简格式输出 JSON,必须包含字段 {required_keys}: {prompt}"
            raw_response = self.fallback_llm(fallback_prompt, timeout=3.0)
            parsed = self._try_repair_json(raw_response, required_keys)
            if parsed:
                parsed["_execution_meta"] = {"source": "fallback_llm", "degraded": True}
                return parsed
        except Exception as e:
            logging.error(f"备用 LLM 调用依然失败: {str(e)}")

        # 4. 兜底规则引擎
        return self._static_rule_fallback(prompt, reason="All LLMs unavailable")

    def _static_rule_fallback(self, prompt: str, reason: str) -> Dict[str, Any]:
        """本地静态规则/缓存兜底"""
        logging.warning(f"触发终极兜底逻辑,原因: {reason}")
        return {
            "status": "degraded_fallback",
            "summary": "系统当前繁忙或服务不可用,已自动切换至安全输出模式",
            "data": [],
            "_execution_meta": {
                "source": "static_rule_engine",
                "degraded": True,
                "reason": reason
            }
        }


# --- 模拟测试 ---
if __name__ == "__main__":
    def mock_primary_llm(p: str, timeout: float):
        raise TimeoutError("LLM API connection read timeout after 5.0s")

    def mock_fallback_llm(p: str, timeout: float):
        return '```json\n{"summary": "备用模型解析结果", "status": "ok"}\n```'

    engine = ResilientAgentEngine(mock_primary_llm, mock_fallback_llm)
    
    print("--- 第一次执行(主模型失败,降级至备用模型)---")
    res1 = engine.execute_task("解析 2026 年 Q2 财报数据", required_keys=["summary", "status"])
    print("结果:", json.dumps(res1, ensure_ascii=False, indent=2))

    print("\n--- 连续触发失败以验证熔断器 ---")
    engine.execute_task("测试 1", ["status"])
    engine.execute_task("测试 2", ["status"])

    print("\n--- 熔断打开后的请求(直接触发规则兜底,无需发起 API)---")
    res_breaker = engine.execute_task("测试 3", ["status"])
    print("结果:", json.dumps(res_breaker, ensure_ascii=False, indent=2))

4. 灰度验证与防护边界:什么场景绝不能自动降级

降级机制虽然能保障系统可用性,但如果缺乏边界控制也会带来潜在风险。

在进行系统治理时,常见误区是盲目对所有接口一刀切降级

在架构设计时,必须将业务场景按风险等级严格划分。对于只读查询、内容摘要、推荐词生成等容错率高、离线度高的场景,降级策略可以较为激进。哪怕返回的是半小时前的静态缓存,或者由几条规则拼凑出来的备用结果,用户体验也优于直接返回系统故障提示。

但是,一旦涉及资金、权限控制、数据修改或合规审批等高风险写操作,自动静默降级存在重大风险

例如,在财务自动化转账 Agent 中,如果模型因为输出格式报错导致 JSON 解析失败,降级引擎如果默认填入一个 amount: 0approve: true 传给下游系统,就会造成严重逻辑漏洞。在这些高敏感场景下,降级策略必须是显式阻断——立即记录现场堆栈,保存原始 Prompt 和 Raw Response,向告警管道发送高级别告警,并要求人工接入审查。

此外,系统的可观测性也是降级方案能否持续运转的关键。所有降级响应在 API 返回头中必须携带 X-Agent-Degraded: true 的 Tag,并且将降级事件录入 OpenTelemetry Trace 中。在监控面板上必须挂载降级率指标:一旦某条业务线 5 分钟内的降级率升至 3% 以上,就必须触发告警;升至 10% 则必须启动自动灰度回滚。


5. 总结

设计高可用 Agent 架构的关键,在于当模型输出异常或遭遇网络故障时,系统能够保持稳定且符合预期的运行逻辑。

依赖远端 API 始终稳定是不切实际的。通过设置严格的 Timeout 缩减无效等待,借助熔断器割断异常流量,利用本地 Schema 校验和备用模型组建多层防护,再结合静态规则保底,才能提升 Agent 系统在生产环境下的抗风险能力。

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐