大模型 API 工程化:兼容层设计、降级与重试策略
问题
生产环境接入大模型 API 时,如何设计统一的兼容层,让 OpenAI、Anthropic、国产大模型(文心、通义、智谱、月之暗面)无缝切换?当模型不可用、限流、超时时,怎么做到优雅降级和自动重试?
这个问题的本质是把模型供应商当作一个外部依赖来管理——和数据库、缓存、消息队列没什么不同。但大模型 API 有它自己的特殊性:
- 不同厂商的 API 格式、参数名、返回结构不一致
- 模型本身是"非确定性"的,同样的输入可能返回不同质量
- Token 消耗直接和成本挂钩,降级不仅仅是"换一个服务",还要考虑成本、延迟、能力
兼容层设计:Adapter 模式 + 统一接口
核心思路是面向接口编程,而不是对着某个 SDK 写死。定义一个 LLMClient 接口,每个模型实现自己的 Adapter。
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import AsyncIterator, Optional
@dataclass
class LLMConfig:
api_key: str
base_url: str
model: str
max_tokens: int = 4096
temperature: float = 0.7
timeout_seconds: int = 30
@dataclass
class LLMResponse:
content: str
model: str
usage: dict # {"prompt_tokens": 100, "completion_tokens": 50}
latency_ms: float
class LLMClient(ABC):
@abstractmethod
async def chat(self, messages: list[dict], **kwargs) -> LLMResponse:
"""非流式对话"""
...
@abstractmethod
async def chat_stream(self, messages: list[dict], **kwargs) -> AsyncIterator[str]:
"""流式对话"""
...
@abstractmethod
async def embed(self, texts: list[str]) -> list[list[float]]:
"""文本向量化"""
...然后每个厂商实现一个 Adapter:
class OpenAIClient(LLMClient):
def __init__(self, config: LLMConfig):
self.config = config
self.client = OpenAI(api_key=config.api_key, base_url=config.base_url)
async def chat(self, messages, **kwargs) -> LLMResponse:
start = time.monotonic()
resp = await self.client.chat.completions.create(
model=self.config.model,
messages=messages,
**kwargs
)
return LLMResponse(
content=resp.choices[0].message.content,
model=self.config.model,
usage=resp.usage.model_dump(),
latency_ms=(time.monotonic() - start) * 1000
)
class ZhipuClient(LLMClient):
"""智谱 GLM 系列 Adapter——API 结构不同,需要适配"""
def __init__(self, config: LLMConfig):
self.config = config
self.client = ZhipuAI(api_key=config.api_key)
async def chat(self, messages, **kwargs) -> LLMResponse:
start = time.monotonic()
# 智谱的 API 参数名和 OpenAI 不完全一致
resp = await self.client.chat.completions.create(
model=self.config.model,
messages=messages,
**kwargs
)
# 返回格式也不一样,需要映射
return LLMResponse(
content=resp.choices[0].message.content,
model=self.config.model,
usage={"prompt_tokens": resp.usage.prompt_tokens,
"completion_tokens": resp.usage.completion_tokens},
latency_ms=(time.monotonic() - start) * 1000
)为什么用 Adapter 而不是直接统一 SDK? 因为每个厂商的 SDK 版本更新节奏不同、参数差异(如 top_p、stop、tools 的命名和格式)、以及错误码和限流返回格式都不一样。Adapter 模式把差异隔离在薄薄一层,主业务代码不感知。
降级策略:三层降级
降级不是"坏了就换一个",而是分场景决策:
@dataclass
class ModelRoute:
"""模型路由配置"""
primary: str # 主模型
fallback_cost: Optional[str] = None # 成本降级目标
fallback_speed: Optional[str] = None # 速率降级目标
fallback_backup: Optional[str] = None # 备线
class LLMRouter:
def __init__(self, clients: dict[str, LLMClient], routes: dict[str, ModelRoute]):
self.clients = clients
self.routes = routes
async def chat_with_fallback(self, route_key: str, messages: list[dict]) -> LLMResponse:
route = self.routes[route_key]
# 第一层:主模型
try:
return await self.clients[route.primary].chat(messages)
except RateLimitError:
# 速率降级——主模型限流,切到更快/更便宜的备选
pass
except ServiceUnavailableError:
# 服务端故障——走备线
pass
except TimeoutError:
# 超时——走备线
pass
# 第二层:速率降级(主模型限流时)
if route.fallback_speed:
try:
return await self.clients[route.fallback_speed].chat(messages)
except Exception:
pass
# 第三层:备线降级(完全换模型)
if route.fallback_backup:
try:
return await self.clients[route.fallback_backup].chat(messages)
except Exception as e:
raise RuntimeError(f"All models unavailable: {e}") from e
raise RuntimeError("No fallback configured")三种降级场景的实际例子:
| 场景 | 主模型 | 降级目标 | 触发条件 |
|---|---|---|---|
| 成本降级 | GPT-4o | GPT-4o-mini | 普通对话,不需要太强推理 |
| 速率降级 | GPT-4o | Claude 3.5 Sonnet | GPT-4o 限流(429) |
| 故障降级 | GPT-4o | 通义千问-Max | 主模型不可用(503) |
关键原则:降级不能让用户感知到"服务坏了",但可以用不同的模型能力。比如:复杂分析用 GPT-4o,简单问答降级到 GPT-4o-mini 或国产模型,保证响应速度。
重试策略:指数退避 + 抖动
重试不能简单"失败了就重试",需要区分可重试和不可重试的错误:
import asyncio
import random
RETRYABLE_STATUS_CODES = {429, 500, 502, 503}
MAX_RETRIES = 3
BASE_DELAY = 1.0 # 秒
MAX_DELAY = 30.0
async def chat_with_retry(client: LLMClient, messages: list[dict]) -> LLMResponse:
last_exception = None
for attempt in range(MAX_RETRIES):
try:
return await client.chat(messages)
except APIError as e:
if e.status_code not in RETRYABLE_STATUS_CODES:
raise # 400、401、404 不重试
last_exception = e
# 指数退避 + 抖动
delay = min(BASE_DELAY * (2 ** attempt), MAX_DELAY)
jitter = random.uniform(0, 1.0)
await asyncio.sleep(delay + jitter)
raise last_exception为什么需要抖动(Jitter)? 假设 100 个请求同时触发重试,没有抖动的话它们会在 1s、2s、4s 这三个时间点同时发起重试,产生"惊群效应"——备线模型也被打爆。加上随机抖动后,请求均匀分布,备线来得及消化。
不可重试的错误(直接上报,不浪费资源):
400 Bad Request:Prompt 格式错误,重试 100 次也一样401 Unauthorized:API Key 过期或无效,需要人工介入404 Not Found:模型名称拼错(比如gpt-4-oo而不是gpt-4o)403 Forbidden:没有权限
高级实践:熔断器 + 请求排队
熔断器(Circuit Breaker)
连续重试失败时,应该打开熔断器,暂时完全停止请求该模型,避免浪费资源:
class CircuitBreaker:
def __init__(self, failure_threshold=5, recovery_timeout=30):
self.failure_count = 0
self.failure_threshold = failure_threshold
self.recovery_timeout = recovery_timeout
self.last_failure_time = 0
self.state = "closed" # closed | open | half-open
def call(self, func, *args, **kwargs):
if self.state == "open":
if time.monotonic() - self.last_failure_time > self.recovery_timeout:
self.state = "half-open" # 允许试探一个请求
else:
raise CircuitBreakerOpenError("Circuit breaker is open")
try:
result = func(*args, **kwargs)
if self.state == "half-open":
self.state = "closed" # 试探成功,关闭熔断器
self.failure_count = 0
return result
except Exception as e:
self.failure_count += 1
self.last_failure_time = time.monotonic()
if self.failure_count >= self.failure_threshold:
self.state = "open"
raise e请求排队与 Token 预算
高并发场景下,不能每个请求都直接打到模型 API。需要队列缓冲 + 优先级调度:
class LLMRequestQueue:
"""基于优先级的请求队列,保护模型 API 不被打爆"""
def __init__(self, max_concurrent=10, token_budget_per_user=100000):
self.queue = asyncio.PriorityQueue()
self.semaphore = asyncio.Semaphore(max_concurrent)
self.token_budget = defaultdict(lambda: token_budget_per_user)
async def enqueue(self, request: LLMRequest, priority: int = 5):
"""priority 越小优先级越高,紧急请求设 priority=1"""
# 检查 Token 预算
if self.token_budget[request.user_id] <= 0:
raise BudgetExceededError("Token budget exceeded")
await self.queue.put((priority, request))
async def process_loop(self, client: LLMClient):
while True:
_, request = await self.queue.get()
async with self.semaphore:
resp = await client.chat(request.messages)
self.token_budget[request.user_id] -= resp.usage["total_tokens"]
# 通过 Webhook 或回调返回结果
await request.callback(resp)超时管理:三层超时
大模型调用的延迟波动很大,需要设置分层超时:
class LLMClient(ABC):
def __init__(self, config: LLMConfig):
# 连接超时:建立 TCP 连接的最大时间
connect_timeout = 5.0
# 读取超时:等待首 Token 返回的最大时间(TTFT)
read_timeout = 30.0
# 流式超时:流式读取时两次 Token 之间的最大间隔
stream_timeout = 60.0为什么需要三层? 连接超时 5s 就够了——DNS 解析和 TCP 握手不应该太久。但首 Token 延迟(TTFT)取决于模型大小和负载,大模型预热慢,30s 是合理的。流式超时 60s 是为了防止模型"卡住"不再输出 Token 时,客户端一直等。
总结
大模型 API 工程化的核心是把模型当作基础设施来管理,而不是当作"黑盒魔法"。
- 兼容层:用 Adapter 模式统一接口,支持多模型切换
- 降级策略:成本降级、速率降级、故障降级三层,匹配不同场景
- 重试策略:指数退避 + 抖动,只重试可重试的错误码
- 熔断 + 排队:保护后端不被冲垮,保护预算不被超支
- 超时管理:分层超时,不同阶段不同预期
如果你正在做 AI 应用,先写一个 LLMClient 接口,然后再对接具体厂商。一个接口 + 多个 Adapter 的成本,远低于以后从 OpenAI 切换成国产模型时的重构代价。