06-大模型API调用失败-超时重试限流和熔断设计
大模型 API 调用失败怎么办?超时、重试、限流和熔断设计
系列:Python + FastAPI 大模型应用基础(第 6 篇)
目标:在客户端级异常处理之上,构建包含整体超时、有限重试、令牌桶限流和熔断器状态机的系统级韧性层。
1. 本篇与第 3 篇有什么区别
第 3 篇解决的是单个模型客户端如何处理:
- 连接超时和读取超时;
- 可重试与不可重试错误;
- 指数退避和
Retry-After; - 流式响应中断。
当模型客户端被放进 FastAPI 并面对并发用户后,仅有客户端重试还不够:
- 100 个请求同时失败,每个重试 3 次,可能放大成 300 次上游请求;
- 某个模型持续故障时,系统仍然等待到超时才返回;
- 恶意或异常流量可能耗尽连接池和模型额度;
- 单次超时设置为 20 秒,重试 3 次后整体可能超过 60 秒;
- 上游已经不可用时,继续请求只会加重故障。
因此,本篇关注的是系统级 Resilience(韧性)设计。
2. 四种机制分别解决什么问题
| 机制 | 解决的问题 | 不能解决的问题 |
|---|---|---|
| Timeout(超时) | 限制一次或整个调用链等待多久 | 不能让故障自动恢复 |
| Retry(重试) | 处理少量、短暂、可能恢复的故障 | 不能处理持续故障,可能放大流量 |
| Rate Limiter(限流器) | 控制请求进入速度,保护资源和额度 | 不能判断上游是否健康 |
| Circuit Breaker(熔断器) | 上游持续失败时快速拒绝,给依赖恢复时间 | 不能代替输入校验和用户限流 |
合理的调用链:
用户请求
↓
身份认证与参数校验
↓
限流:当前请求是否允许进入?
↓
熔断:上游当前是否允许尝试?
↓
整体调用截止时间
↓
有限重试 + 单次 HTTP 超时
↓
模型服务
3. Timeout:必须同时考虑单次超时和整体截止时间
假设配置如下:
单次读取超时:20 秒
最大请求次数:3 次
两次重试等待:1 秒、2 秒
最坏等待时间可能超过:
20 + 1 + 20 + 2 + 20 = 63 秒
如果业务要求用户最多等待 45 秒,仅设置单次 HTTP 超时无法满足要求,还需要 Total Deadline(整体截止时间)。
本文使用两层限制:
httpx.Timeout:限制每次 HTTP 连接、写入和读取;asyncio.timeout(45):限制重试和等待在内的整个模型调用。
asyncio.timeout() 需要 Python 3.11 及以上版本,因此本文使用 Python 3.11+。
4. Retry:为什么不能把所有错误都重试
推荐分类:
| 状态或异常 | 是否重试 | 原因 |
|---|---|---|
| 网络连接失败 | 有限重试 | 可能是短暂网络抖动 |
| 读取超时 | 有条件重试 | 可能恢复,但也可能重复计费 |
| 408 | 有限重试 | 请求超时可能暂时恢复 |
| 429 | 尊重 Retry-After 后重试 | 服务端明确要求降速 |
| 500、502、503、504 | 有限重试 | 上游可能短暂故障 |
| 400、422 | 不重试 | 参数不变,结果不会改变 |
| 401、403 | 不重试 | 密钥或权限需要修复 |
| 404 | 不重试 | URL 或模型标识需要修复 |
重试必须包含:
- 最大次数;
- 指数退避;
- 随机抖动;
- 整体截止时间;
- 只捕获明确的短暂错误。
5. Rate Limiter:令牌桶如何工作
Token Bucket(令牌桶)可以理解为一个持续补充令牌的容器:
桶容量:10 个令牌
补充速度:每秒 2 个令牌
每个请求:消耗 1 个令牌
当桶中有令牌,请求可以立即通过;没有令牌时,请求被拒绝或等待。
令牌桶允许有限突发流量:如果一段时间没有请求,桶中可以积累令牌,但最多不超过容量。
6. Circuit Breaker:熔断器状态机
熔断器通常有三种状态:
CLOSED(关闭)
正常放行请求并统计连续失败
│
│ 达到失败阈值
▼
OPEN(打开)
快速拒绝请求,不再访问上游
│
│ 恢复等待时间结束
▼
HALF_OPEN(半开)
只允许少量探测请求
│ │
│ 成功 │ 失败
▼ ▼
CLOSED OPEN
“熔断器打开”听起来像允许流量通过,但在工程语义中恰好相反:电路断开,请求被快速拒绝。
7. 创建项目
项目结构:
resilient_model_service/
├── app/
│ ├── __init__.py
│ ├── resilience.py
│ ├── model_service.py
│ └── main.py
└── requirements.txt
requirements.txt:
fastapi>=0.115,<1
uvicorn[standard]>=0.30,<1
httpx>=0.27,<1
pydantic>=2.7,<3
创建环境并安装依赖:
python -m venv .venv
.\.venv\Scripts\python.exe -m pip install -r requirements.txt
app/__init__.py 是空文件。
8. 实现重试、令牌桶和熔断器
新建 app/resilience.py:
import asyncio
import random
import time
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
from enum import Enum
from typing import TypeVar
T = TypeVar("T")
class TransientUpstreamError(RuntimeError):
"""可能通过等待和重试恢复的上游错误。"""
def __init__(
self,
message: str,
retry_after: float | None = None,
) -> None:
super().__init__(message)
self.retry_after = retry_after
class PermanentUpstreamError(RuntimeError):
"""参数、鉴权或权限等不应自动重试的错误。"""
class RateLimitExceeded(RuntimeError):
"""当前服务的本地限流器拒绝了请求。"""
def __init__(self, retry_after: float) -> None:
super().__init__("请求过于频繁,请稍后重试")
self.retry_after = retry_after
class CircuitOpen(RuntimeError):
"""熔断器打开,当前不允许访问上游。"""
def __init__(self, retry_after: float) -> None:
super().__init__("模型服务暂时熔断,请稍后重试")
self.retry_after = retry_after
@dataclass(frozen=True)
class RetryPolicy:
"""有限重试参数。"""
# 包含第一次调用,例如 3 表示最多调用 3 次
max_attempts: int = 3
base_delay: float = 0.5
max_delay: float = 5.0
def __post_init__(self) -> None:
if self.max_attempts < 1:
raise ValueError("max_attempts 不能小于 1")
if self.base_delay < 0 or self.max_delay < 0:
raise ValueError("重试等待时间不能为负数")
async def retry_async(
operation: Callable[[], Awaitable[T]],
policy: RetryPolicy,
) -> T:
"""只重试 TransientUpstreamError。"""
for attempt in range(1, policy.max_attempts + 1):
try:
return await operation()
except TransientUpstreamError as exc:
if attempt >= policy.max_attempts:
# 最后一次失败后保留原异常,不再等待
raise
if exc.retry_after is not None:
# 优先参考上游 Retry-After,同时限制最长等待
delay = min(max(0.0, exc.retry_after), policy.max_delay)
else:
exponential_delay = (
policy.base_delay * (2 ** (attempt - 1))
)
capped_delay = min(exponential_delay, policy.max_delay)
# 加入随机抖动,减少大量请求同时重试
jitter = random.uniform(
0.0,
min(0.5, capped_delay * 0.25),
)
delay = capped_delay + jitter
# 异步等待不会阻塞整个事件循环
await asyncio.sleep(delay)
# 正常情况下循环一定会 return 或 raise
raise RuntimeError("重试循环未返回结果")
class AsyncTokenBucket:
"""进程内异步令牌桶限流器。"""
def __init__(self, capacity: int, refill_rate: float) -> None:
if capacity <= 0:
raise ValueError("capacity 必须大于 0")
if refill_rate <= 0:
raise ValueError("refill_rate 必须大于 0")
self.capacity = float(capacity)
self.refill_rate = refill_rate
self._tokens = float(capacity)
self._last_refill = time.monotonic()
self._lock = asyncio.Lock()
async def try_acquire(self, cost: float = 1.0) -> float | None:
"""尝试消费令牌;成功返回 None,失败返回建议等待秒数。"""
if cost <= 0:
raise ValueError("cost 必须大于 0")
if cost > self.capacity:
raise ValueError("单次 cost 不能超过桶容量")
async with self._lock:
now = time.monotonic()
elapsed = now - self._last_refill
# 根据经过时间补充令牌,但不能超过桶容量
self._tokens = min(
self.capacity,
self._tokens + elapsed * self.refill_rate,
)
self._last_refill = now
if self._tokens >= cost:
self._tokens -= cost
return None
missing_tokens = cost - self._tokens
retry_after = missing_tokens / self.refill_rate
return retry_after
class CircuitState(str, Enum):
"""熔断器的三种状态。"""
CLOSED = "closed"
OPEN = "open"
HALF_OPEN = "half_open"
class CircuitBreaker:
"""简化的进程内异步熔断器。"""
def __init__(
self,
failure_threshold: int = 3,
recovery_timeout: float = 30.0,
) -> None:
if failure_threshold < 1:
raise ValueError("failure_threshold 不能小于 1")
if recovery_timeout <= 0:
raise ValueError("recovery_timeout 必须大于 0")
self.failure_threshold = failure_threshold
self.recovery_timeout = recovery_timeout
self._state = CircuitState.CLOSED
self._failure_count = 0
self._opened_at = 0.0
self._half_open_call_in_progress = False
self._lock = asyncio.Lock()
def _open_circuit(self, now: float) -> None:
"""调用者必须已经持有锁。"""
self._state = CircuitState.OPEN
self._opened_at = now
self._half_open_call_in_progress = False
async def before_call(self) -> None:
"""在访问上游前判断本次请求是否允许通过。"""
async with self._lock:
now = time.monotonic()
if self._state == CircuitState.OPEN:
elapsed = now - self._opened_at
remaining = self.recovery_timeout - elapsed
if remaining > 0:
raise CircuitOpen(retry_after=remaining)
# 恢复时间结束,进入半开状态进行一次探测
self._state = CircuitState.HALF_OPEN
if self._state == CircuitState.HALF_OPEN:
if self._half_open_call_in_progress:
# 半开状态只允许一个探测请求访问上游
raise CircuitOpen(retry_after=1.0)
self._half_open_call_in_progress = True
async def record_success(self) -> None:
"""调用成功后关闭熔断器并清零失败次数。"""
async with self._lock:
self._state = CircuitState.CLOSED
self._failure_count = 0
self._opened_at = 0.0
self._half_open_call_in_progress = False
async def record_failure(self) -> None:
"""记录一次最终失败,必要时打开熔断器。"""
async with self._lock:
now = time.monotonic()
if self._state == CircuitState.HALF_OPEN:
# 探测请求失败,立即重新打开熔断器
self._failure_count = self.failure_threshold
self._open_circuit(now)
return
self._failure_count += 1
if self._failure_count >= self.failure_threshold:
self._open_circuit(now)
async def record_abandoned(self) -> None:
"""半开探测被取消时释放占用,并重新进入恢复等待。"""
async with self._lock:
if self._state == CircuitState.HALF_OPEN:
self._open_circuit(time.monotonic())
async def snapshot(self) -> dict[str, object]:
"""返回当前状态快照,用于内部监控。"""
async with self._lock:
return {
"state": self._state.value,
"failure_count": self._failure_count,
"failure_threshold": self.failure_threshold,
}
9. 实现具有韧性的模型服务
新建 app/model_service.py:
import asyncio
from dataclasses import dataclass
from typing import Any
import httpx
from app.resilience import (
AsyncTokenBucket,
CircuitBreaker,
PermanentUpstreamError,
RateLimitExceeded,
RetryPolicy,
TransientUpstreamError,
retry_async,
)
@dataclass(frozen=True)
class ModelSettings:
"""模型服务配置。"""
api_key: str
base_url: str
model: str
class ResilientModelService:
"""统一编排限流、熔断、整体超时和重试。"""
def __init__(
self,
http_client: httpx.AsyncClient,
settings: ModelSettings,
rate_limiter: AsyncTokenBucket,
circuit_breaker: CircuitBreaker,
retry_policy: RetryPolicy,
total_timeout: float = 45.0,
) -> None:
self.http_client = http_client
self.settings = settings
self.rate_limiter = rate_limiter
self.circuit_breaker = circuit_breaker
self.retry_policy = retry_policy
self.total_timeout = total_timeout
def _parse_retry_after(self, value: str | None) -> float | None:
"""解析秒数形式的 Retry-After。"""
if not value:
return None
try:
return max(0.0, float(value))
except ValueError:
# HTTP 日期格式可复用第 3 篇的完整解析逻辑
return None
async def _call_once(self, user_message: str) -> str:
"""只执行一次上游模型调用,不在此方法内部重试。"""
request_url = (
f"{self.settings.base_url.rstrip('/')}/chat/completions"
)
request_headers = {
"Authorization": f"Bearer {self.settings.api_key}",
"Content-Type": "application/json",
}
request_body = {
"model": self.settings.model,
"messages": [
{
"role": "system",
"content": "你是一名严谨的 Python 助手。",
},
{"role": "user", "content": user_message},
],
"temperature": 0.2,
}
try:
response = await self.http_client.post(
request_url,
headers=request_headers,
json=request_body,
)
except httpx.TimeoutException as exc:
raise TransientUpstreamError("单次模型请求超时") from exc
except (httpx.ConnectError, httpx.NetworkError) as exc:
raise TransientUpstreamError(
"模型服务发生网络异常"
) from exc
except httpx.HTTPError as exc:
raise TransientUpstreamError(
"模型 HTTP 请求异常"
) from exc
status_code = response.status_code
if status_code in {408, 429} or status_code >= 500:
retry_after = self._parse_retry_after(
response.headers.get("Retry-After")
)
raise TransientUpstreamError(
f"模型服务暂时异常,状态码:{status_code}",
retry_after=retry_after,
)
if status_code >= 400:
# 参数、鉴权、权限和地址错误不能通过原样重试修复
raise PermanentUpstreamError(
f"模型请求不可重试,状态码:{status_code}"
)
try:
data: dict[str, Any] = response.json()
content = data["choices"][0]["message"]["content"]
except (ValueError, KeyError, IndexError, TypeError) as exc:
# 未确认可以自行恢复前,契约错误不应盲目重试
raise PermanentUpstreamError(
"模型响应结构与当前适配器契约不一致"
) from exc
if not isinstance(content, str) or not content.strip():
raise PermanentUpstreamError("模型返回了空内容")
return content.strip()
async def generate(self, user_message: str) -> str:
"""按照限流、熔断、截止时间、重试的顺序调用模型。"""
cleaned_message = user_message.strip()
if not cleaned_message:
raise ValueError("用户问题不能为空")
# 第一道保护:控制进入模型调用链的请求速度
retry_after = await self.rate_limiter.try_acquire(cost=1.0)
if retry_after is not None:
raise RateLimitExceeded(retry_after=retry_after)
# 第二道保护:上游持续故障时快速失败
await self.circuit_breaker.before_call()
try:
# 整体截止时间包含每次请求、重试等待和网络耗时
async with asyncio.timeout(self.total_timeout):
result = await retry_async(
operation=lambda: self._call_once(cleaned_message),
policy=self.retry_policy,
)
except TimeoutError as exc:
await self.circuit_breaker.record_failure()
raise TransientUpstreamError(
"模型调用超过整体截止时间"
) from exc
except TransientUpstreamError:
# 多次重试后的最终失败才计入一次熔断失败
await self.circuit_breaker.record_failure()
raise
except PermanentUpstreamError:
# 永久错误说明上游有响应,不计入服务可用性失败
await self.circuit_breaker.record_success()
raise
except asyncio.CancelledError:
# 客户端取消半开探测时,不能永远占用探测名额
await self.circuit_breaker.record_abandoned()
raise
except Exception:
# 未知异常默认记为失败并继续抛出,不能静默吞掉
await self.circuit_breaker.record_failure()
raise
await self.circuit_breaker.record_success()
return result
为什么只把“最终失败”计入熔断器
本文把一次用户请求及其内部重试看作一次逻辑调用:
用户请求 1
├── 第一次上游调用失败
├── 第二次上游调用失败
└── 第三次上游调用成功
如果每个内部重试都计入熔断器,一次用户请求可能直接累计三次失败并打开熔断器。本文只在所有重试最终失败后记录一次。
这是一个设计选择,不是唯一答案。生产系统需要根据请求量、错误率窗口和熔断库语义决定统计方式。
10. 接入 FastAPI
新建 app/main.py:
import os
from collections.abc import AsyncIterator
from contextlib import asynccontextmanager
import httpx
from fastapi import FastAPI, Request
from fastapi.responses import JSONResponse
from pydantic import BaseModel, Field, field_validator
from app.model_service import ModelSettings, ResilientModelService
from app.resilience import (
AsyncTokenBucket,
CircuitBreaker,
CircuitOpen,
PermanentUpstreamError,
RateLimitExceeded,
RetryPolicy,
TransientUpstreamError,
)
class ChatRequest(BaseModel):
"""聊天接口请求。"""
user_message: str = Field(min_length=1, max_length=4000)
@field_validator("user_message")
@classmethod
def message_must_not_be_blank(cls, value: str) -> str:
cleaned_value = value.strip()
if not cleaned_value:
raise ValueError("user_message 不能为空白字符串")
return cleaned_value
class ChatResponse(BaseModel):
"""聊天接口成功响应。"""
content: str
model: str
def required_env(name: str) -> str:
"""读取必需环境变量。"""
value = os.getenv(name, "").strip()
if not value:
raise RuntimeError(f"缺少环境变量:{name}")
return value
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
"""创建连接池和韧性组件。"""
settings = ModelSettings(
api_key=required_env("LLM_API_KEY"),
base_url=required_env("LLM_BASE_URL"),
model=required_env("LLM_MODEL"),
)
# 这是单次 HTTP 调用的超时,不是整个重试链路的超时
http_timeout = httpx.Timeout(
connect=5.0,
read=20.0,
write=10.0,
pool=5.0,
)
http_client = httpx.AsyncClient(timeout=http_timeout)
rate_limiter = AsyncTokenBucket(
capacity=10, # 最多允许积累 10 个令牌
refill_rate=2.0, # 每秒补充 2 个令牌
)
circuit_breaker = CircuitBreaker(
failure_threshold=3, # 连续 3 次最终失败后熔断
recovery_timeout=30.0,
)
retry_policy = RetryPolicy(
max_attempts=3,
base_delay=0.5,
max_delay=5.0,
)
app.state.model_service = ResilientModelService(
http_client=http_client,
settings=settings,
rate_limiter=rate_limiter,
circuit_breaker=circuit_breaker,
retry_policy=retry_policy,
total_timeout=45.0,
)
app.state.circuit_breaker = circuit_breaker
app.state.model_name = settings.model
yield
await http_client.aclose()
app = FastAPI(
title="具有韧性的大模型服务",
version="1.0.0",
lifespan=lifespan,
)
def retry_after_header(seconds: float) -> dict[str, str]:
"""生成合法且至少为 1 秒的 Retry-After 响应头。"""
rounded_seconds = max(1, int(seconds + 0.999))
return {"Retry-After": str(rounded_seconds)}
@app.exception_handler(RateLimitExceeded)
async def handle_rate_limit(
request: Request,
exc: RateLimitExceeded,
) -> JSONResponse:
return JSONResponse(
status_code=429,
content={"detail": str(exc)},
headers=retry_after_header(exc.retry_after),
)
@app.exception_handler(CircuitOpen)
async def handle_circuit_open(
request: Request,
exc: CircuitOpen,
) -> JSONResponse:
return JSONResponse(
status_code=503,
content={"detail": str(exc)},
headers=retry_after_header(exc.retry_after),
)
@app.exception_handler(TransientUpstreamError)
async def handle_transient_error(
request: Request,
exc: TransientUpstreamError,
) -> JSONResponse:
return JSONResponse(
status_code=503,
content={"detail": str(exc)},
)
@app.exception_handler(PermanentUpstreamError)
async def handle_permanent_error(
request: Request,
exc: PermanentUpstreamError,
) -> JSONResponse:
# 这是后端调用上游失败,不应该把上游 401 直接变成用户的 401
return JSONResponse(
status_code=502,
content={"detail": str(exc)},
)
@app.post("/api/v1/chat", response_model=ChatResponse)
async def chat(
chat_request: ChatRequest,
request: Request,
) -> ChatResponse:
"""通过韧性层调用模型。"""
model_service: ResilientModelService = (
request.app.state.model_service
)
content = await model_service.generate(chat_request.user_message)
return ChatResponse(
content=content,
model=request.app.state.model_name,
)
@app.get("/internal/resilience")
async def resilience_status(request: Request) -> dict[str, object]:
"""仅供内部监控使用,生产环境不能公开暴露。"""
circuit_breaker: CircuitBreaker = (
request.app.state.circuit_breaker
)
return await circuit_breaker.snapshot()
设置环境变量并启动:
$env:LLM_API_KEY = "替换为真实密钥"
$env:LLM_BASE_URL = "https://替换为模型服务地址/v1"
$env:LLM_MODEL = "替换为真实模型标识"
.\.venv\Scripts\python.exe -m uvicorn app.main:app --reload
11. 不连接真实模型也能测试重试逻辑
先用一个前两次失败、第三次成功的函数测试 retry_async():
import asyncio
from app.resilience import (
RetryPolicy,
TransientUpstreamError,
retry_async,
)
async def main() -> None:
attempt_count = 0
async def unstable_operation() -> str:
"""模拟一个前两次失败、第三次成功的上游服务。"""
nonlocal attempt_count
attempt_count += 1
print(f"正在执行第 {attempt_count} 次调用")
if attempt_count < 3:
raise TransientUpstreamError("模拟临时故障")
return "第三次调用成功"
result = await retry_async(
operation=unstable_operation,
policy=RetryPolicy(
max_attempts=3,
base_delay=0.1,
max_delay=0.5,
),
)
assert attempt_count == 3
assert result == "第三次调用成功"
print(result)
if __name__ == "__main__":
asyncio.run(main())
测试韧性机制时,优先使用 Fake(模拟对象)和可控制的失败,不要通过高频请求真实模型来“测试限流”,否则可能产生费用或违反服务商使用限制。
12. 为什么限流应该在重试之前
如果每个用户请求进入系统后都先创建任务、占用连接,再判断是否超限,限流就失去了保护入口的意义。
本文先执行本地令牌桶,再执行熔断和重试:
请求到达
↓
本地限流
↓
熔断判断
↓
最多 3 次上游调用
需要注意:一个用户请求可能产生多次上游请求。因此还应分别监控:
- 用户请求速率;
- 逻辑模型调用次数;
- 实际上游 HTTP 请求次数;
- 重试次数和重试放大倍数。
13. 为什么不能只依赖模型厂商的 429
等上游返回 429 才降速,说明请求已经离开自己的系统并消耗网络和连接资源。
自己的限流器可以保护:
- FastAPI 工作进程;
- HTTP 连接池;
- 用户或部门额度;
- 模型费用预算;
- 下游 CRM、数据库等其他资源。
上游 429 和本地限流不是二选一,两者保护的边界不同。
14. 单进程限流器的局限
本文令牌桶只存在于当前 Python 进程:
进程 A:自己的 10 个令牌
进程 B:自己的 10 个令牌
进程 C:自己的 10 个令牌
启动三个进程后,实际总容量变成三倍。多实例生产系统通常需要:
- API Gateway(API 网关)统一限流;
- Redis + Lua 实现原子分布式令牌桶;
- 按用户、租户、部门和模型分别设置配额;
- 对高成本模型设置更严格限制。
本文实现用于理解算法和本地项目,不能直接当成分布式限流结论。
15. 熔断器应该统计哪些失败
适合计入熔断的通常是依赖可用性错误:
- 网络连接失败;
- 读取超时;
- 上游 500、502、503、504;
- 持续的 429 或服务过载。
不应该计入熔断的通常是调用方错误:
- 用户输入不合法;
- 请求体缺少字段;
- 业务权限不足;
- 调用了不存在的本地功能。
否则,某个用户持续发送错误参数,可能把所有正常用户的模型通道熔断。
16. 熔断和降级不是同一件事
熔断表示暂时停止访问故障依赖;降级表示熔断后系统向用户提供什么替代能力。
可选降级方案:
- 返回“模型服务繁忙,请稍后重试”;
- 返回经过审核的固定说明;
- 只提供关键词搜索,不生成答案;
- 从缓存中返回最近的相同问题答案;
- 在合规允许且能力匹配时切换备用模型;
- 将任务放入队列,恢复后异步处理。
不能把“自动换模型”当成默认降级,因为它可能改变数据区域、费用和输出行为。
17. 对抗性审查:当前代码仍有哪些风险
17.1 简化熔断器不等于成熟生产库
本文熔断器用于解释状态机。在高并发下,较早发出的请求可能在熔断器状态变化后才返回,其成功或失败会影响当前状态。成熟实现通常还会使用滑动窗口、调用代次、失败率阈值和更严格的并发控制。
17.2 连续失败不等于失败率
本文使用连续失败次数。真实系统可能需要“最近 100 次调用失败率超过 50%”等滑动窗口规则,以避免偶发成功不断重置计数。
17.3 重试可能重复计费
读取超时时,上游可能已经完成生成。重试可能再次产生费用。若模型服务支持幂等键、任务 ID 或结果查询,应使用官方机制。
17.4 本地限流没有用户维度
一个高频用户可能消耗整个进程的令牌,让其他用户无法调用。生产系统至少需要按用户或租户限流,并设置系统总限流。
17.5 没有 Bulkhead(舱壁隔离)
如果所有模型共享同一个连接池和并发队列,一个慢模型可能拖累其他模型。可以为不同模型设置独立连接池、信号量和任务队列。
17.6 内部状态接口不能公开
/internal/resilience 只用于演示。生产环境应放在内部网络并增加身份认证,避免泄露系统运行状态。
17.7 参数必须通过压测和监控确定
3 次失败、30 秒恢复、每秒 2 个令牌 都是示例值,不是通用最佳实践。真实参数必须根据流量、延迟、错误率、预算和服务商限制确定。
17.8 示例接口没有用户鉴权
当前 FastAPI 代码用于解释韧性机制,没有实现登录认证和租户权限。生产系统必须先识别用户,再执行用户级额度、部门级权限和系统总限流,不能把一个全局令牌桶当成完整的访问控制。
18. 需要监控哪些指标
至少记录:
- 用户请求总数;
- 本地限流拒绝数;
- 上游请求总数;
- 每次调用的尝试次数;
- 连接超时和读取超时次数;
- 上游 429、5xx 次数;
- 熔断器打开次数和持续时间;
- 端到端响应时间;
- 首 Token 时间;
- Token 用量和费用。
没有指标就无法判断重试是在提高成功率,还是在放大故障。
19. 本篇总结
本文实现了一条完整的系统级韧性链路:
参数校验
↓
进程内令牌桶限流
↓
熔断器状态判断
↓
45 秒整体截止时间
↓
最多 3 次有限重试
↓
每次 HTTP 连接与读取超时
↓
外部模型服务
关键结论:
- 超时限制等待时间;
- 重试只处理短暂故障;
- 限流保护入口资源和费用;
- 熔断在持续故障时快速失败;
- 重试必须受到整体截止时间约束;
- 本地限流和熔断不能直接替代分布式方案;
- 参数必须通过监控和压测确定,不能凭感觉照抄。
下一篇将记录模型 Token 用量、调用时延和接口成本,为限流、模型路由和业务评估提供数据基础。
20. 练习题
- 把令牌桶容量设为 2,连续发送 5 个请求并观察 429;
- 使用 Fake Operation 连续制造 3 次最终失败,观察熔断状态;
- 等待恢复时间后,验证半开状态只允许一个探测请求;
- 把整体超时设为 1 秒,验证重试链路是否被及时取消;
- 给每次重试增加结构化日志,但不要记录 API Key;
- 统计“用户请求数”和“实际上游请求数”的放大倍数;
- 设计按用户和系统总量两层限流方案;
- 思考备用模型切换需要满足哪些合规和能力条件。
更多推荐
所有评论(0)