大模型 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. 练习题

  1. 把令牌桶容量设为 2,连续发送 5 个请求并观察 429;
  2. 使用 Fake Operation 连续制造 3 次最终失败,观察熔断状态;
  3. 等待恢复时间后,验证半开状态只允许一个探测请求;
  4. 把整体超时设为 1 秒,验证重试链路是否被及时取消;
  5. 给每次重试增加结构化日志,但不要记录 API Key;
  6. 统计“用户请求数”和“实际上游请求数”的放大倍数;
  7. 设计按用户和系统总量两层限流方案;
  8. 思考备用模型切换需要满足哪些合规和能力条件。

更多推荐