AI服务生产部署实战:异步编程与FastAPI高并发架构设计
1. 从“玩具”到“产品”:为什么AI服务部署是道坎
最近和几个做AI应用的朋友聊天,发现一个挺普遍的现象:大家花大量时间在模型选型、Prompt调优、Agent流程设计上,搞出来的Demo在本地跑得飞快,逻辑也堪称精妙。但一旦说到“上线给用户用”,气氛就微妙起来了。要么是接口响应慢得像在挤牙膏,用户等个答案要十几秒;要么是并发一上来,服务直接挂掉,返回一堆“服务器内部错误”;更头疼的是,LLM(大语言模型)的API调用本身就不稳定,偶尔来个超时或者限流,整个服务链条就断了。
这其实就是典型的“玩具”与“产品”的差距。我们之前章节讨论的,无论是用LangChain搭流程,还是用Dify搞低代码,大多是在单线程、低并发的理想环境下验证逻辑可行性。而 生产部署 要解决的,是把这套逻辑变成一个7x24小时稳定、高效、能扛住真实用户流量的在线服务。这里面核心就两个词: 异步 和 部署 。
“异步”不是简单的技术选型,它关乎用户体验的底线。想象一下,用户在你的AI客服对话框里输入问题,前端转圈圈转了半分钟,这体验足以让用户关掉页面。而“部署”则决定了服务的天花板,涉及资源管理、弹性伸缩、故障恢复等一系列工程问题。很多人觉得FastAPI写个
async def
就叫异步了,或者用Docker打个包扔服务器就叫生产部署了,这中间差的火候,正是本章要掰开揉碎讲清楚的东西。
我会结合最近的热点,比如如何应对LLM API的429限流错误、如何设计健壮的异步任务流、以及如何利用FastAPI等现代框架构建真正面向生产的环境,把这条从开发到上线的路铺实。
2. 深入异步:超越
async/await
的效能实战
一提到Python异步,很多人第一反应是
asyncio
和
async/await
语法。这没错,但如果我们止步于此,就像只学了汽车方向盘却不懂变速箱。对于AI服务,尤其是重度依赖外部API(如OpenAI、通义千问等)的服务,异步的核心价值在于
高效处理I/O等待
。
2.1 理解AI服务中的I/O瓶颈:LLM API调用是主因
一个典型的AI服务处理流程,CPU密集的计算其实并不多。时间主要消耗在:
- 网络I/O :向远程LLM API发送请求并等待响应。这个延迟通常在几百毫秒到数秒不等,且极不稳定。
- 磁盘I/O :读取向量数据库(如Chroma、Milvus)中的知识库文档。
- 其他外部服务 :调用搜索引擎、数据库、或其他微服务。
如果使用传统的同步方式,服务器在等待LLM响应的这几秒钟内,当前工作线程会被完全阻塞,什么也干不了。它不能去处理下一个用户的请求,只能空等。这就是为什么同步服务并发能力极差,资源利用率低下的原因。
异步编程通过 事件循环 机制解决了这个问题。当一个异步任务(例如,发起一个LLM API调用)需要等待时,它会主动告知事件循环:“我先歇会儿,等有结果了再叫我”。事件循环就会立刻去执行其他已经就绪的任务。等网络响应返回,事件循环再回来唤醒这个任务继续执行。这样,单个线程就能并发处理成百上千个网络连接,极大地提升了吞吐量。
2.2 FastAPI的异步实践:从路由到依赖注入
FastAPI天生对异步支持友好,但这不代表用了FastAPI就自动获得了高性能。
首先,正确声明异步路由:
from fastapi import FastAPI, BackgroundTasks
import httpx
app = FastAPI()
# 正确:处理函数是异步的,内部执行了异步I/O操作
@app.post("/chat/")
async def chat_completion(question: str):
async with httpx.AsyncClient() as client:
# 假设调用一个LLM API
response = await client.post(
"https://api.llm-provider.com/v1/chat",
json={"message": question},
timeout=30.0
)
return response.json()
# 错误示例:在异步函数内调用同步的、阻塞的LLM客户端库
# async def bad_example(question: str):
# # 某些旧的或设计不佳的SDK可能是同步的
# result = some_sync_llm_client.generate(question) # 这会阻塞事件循环!
# return result
关键点在于,你使用的所有下游客户端(HTTP客户端、数据库驱动、LLM SDK)都必须是
异步兼容
的。对于HTTP请求,推荐使用
httpx
或
aiohttp
。对于数据库,比如PostgreSQL,要用
asyncpg
而不是
psycopg2
。
其次,善用后台任务(BackgroundTasks)处理非即时需求: 不是所有操作都需要即时响应给用户。例如,将对话记录存入数据库、发送异步通知、或触发一个耗时的数据分析任务。
from fastapi import BackgroundTasks
from pydantic import BaseModel
class ChatLog(BaseModel):
user_id: str
question: str
answer: str
def write_log_to_db(chat_log: ChatLog):
# 这是一个同步的、可能较慢的数据库写入操作
# 注意:这里为了演示用了同步函数,实际生产环境应用异步ORM如Tortoise-ORM或SQLAlchemy 1.4+异步模式
time.sleep(0.5) # 模拟耗时
print(f"Log saved for user: {chat_log.user_id}")
@app.post("/chat-with-log/")
async def chat_with_log(
question: str,
background_tasks: BackgroundTasks
):
# 1. 先处理核心的聊天请求
async with httpx.AsyncClient() as client:
llm_response = await client.post(LLM_API, json={"message": question})
answer = llm_response.json()["choices"][0]["message"]["content"]
# 2. 将日志记录任务放入后台,主流程无需等待其完成
log_entry = ChatLog(user_id="user123", question=question, answer=answer)
background_tasks.add_task(write_log_to_db, log_entry)
# 3. 立即返回响应给用户
return {"answer": answer}
这样,用户能快速拿到AI的回复,而日志写入这种不影响主流程的操作在后台慢慢进行,实现了请求响应时间的优化。
2.3 应对LLM API的不稳定性:重试、降级与熔断
LLM服务商(如OpenAI)的API限流(429错误)和间歇性故障是生产环境中的常态。一个健壮的异步服务必须能处理这些故障。
策略一:指数退避重试 直接失败或固定间隔重试会给下游API带来脉冲压力。指数退避能在失败后逐渐增加重试间隔。
import asyncio
import httpx
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception
# 定义一个判断是否需要重试的异常类型
def is_retryable_error(e):
return isinstance(e, (httpx.HTTPStatusError, httpx.RequestError))
@retry(
stop=stop_after_attempt(4), # 最多重试4次(即初始1次+3次重试)
wait=wait_exponential(multiplier=1, min=1, max=10), # 指数退避:1s, 2s, 4s, ... 最大10s
retry=retry_if_exception(is_retryable_error)
)
async def call_llm_api_with_retry(client: httpx.AsyncClient, prompt: str):
try:
resp = await client.post(LLM_API_URL, json={"prompt": prompt}, timeout=30.0)
resp.raise_for_status() # 如果状态码不是2xx,抛出HTTPStatusError
return resp.json()
except httpx.HTTPStatusError as e:
if e.response.status_code == 429:
print(f"Rate limited, will retry. Headers: {e.response.headers}")
raise # 触发重试
elif 500 <= e.response.status_code < 600:
print(f"Server error {e.response.status_code}, will retry.")
raise # 触发重试
else:
# 4xx客户端错误,如认证失败,不应重试
raise
这里使用了
tenacity
库,它让重试逻辑变得非常清晰。注意,我们只对可重试的错误(429限流、5xx服务器错误、网络问题)进行重试。
策略二:服务降级 当主要LLM服务持续不可用时,应有备选方案。例如,可以降级到一个更简单、更稳定的模型,或者返回一个预定义的缓存响应。
async def get_ai_response(user_input: str):
primary_provider = "openai"
fallback_provider = "local_llm" # 或另一个备用API
try:
return await call_primary_llm(user_input, primary_provider)
except (httpx.HTTPStatusError, httpx.RequestError, asyncio.TimeoutError) as e:
logging.warning(f"Primary provider {primary_provider} failed: {e}. Switching to fallback.")
# 触发降级逻辑
return await call_fallback_llm(user_input, fallback_provider)
策略三:熔断器模式(Circuit Breaker)
防止在下游服务故障时,持续不断的请求将其压垮,也避免自身资源被耗尽。可以使用
aiocircuitbreaker
库。
from aiocircuitbreaker import circuit
@circuit(failure_threshold=5, recovery_timeout=30)
async def call_llm_api_circuit(prompt: str):
return await call_llm_api_with_retry(prompt)
当
call_llm_api_circuit
在短时间内失败超过5次,熔断器会“打开”,后续30秒内所有对该函数的调用会立即失败(抛出
CircuitBreakerError
),而不会真正去请求下游API。30秒后,熔断器进入“半开”状态,允许一个试探请求通过,如果成功则关闭熔断器,恢复服务;如果失败,则重新打开。这给了下游服务恢复的时间。
3. 构建健壮的生产级异步任务流
对于超过HTTP请求超时时间(比如超过30秒)的AI长任务,或者需要多步骤编排的复杂AI Agent工作流,我们不能让用户在前端一直等待。这时就需要引入 异步任务队列 。这不仅是“异步”,更是“解耦”和“持久化”。
3.1 任务队列选型:Celery vs RQ vs Dramatiq vs Arq
这是架构决策的关键一步。每个方案都有其适用场景。
| 特性 | Celery | RQ (Redis Queue) | Dramatiq | Arq |
|---|---|---|---|---|
| 成熟度 | 极高 ,行业标准 | 高,简单直接 | 中等,现代化 | 中等,专为asyncio设计 |
| 复杂度 | 高,功能多配置繁 | 低,上手快 | 中等 | 低 |
| Broker支持 | RabbitMQ, Redis, 等 | 仅Redis | RabbitMQ, Redis | 仅Redis |
| 异步支持 | 原生一般,需搭配gevent/eventlet | 同步 | 同步 | 原生asyncio |
| 性能 | 优秀 | 良好 | 优秀 | 优秀(异步) |
| 适用场景 | 大型、复杂、需要多种broker和复杂路由的分布式系统 | 轻量级、快速上手的项目,团队熟悉Redis | 需要高性能和中间件支持,但比Celery简洁 | Python异步生态项目,任务本身是异步的 |
对于现代AI服务(基于FastAPI,大量async/await),我的建议是:
-
首选Arq :如果你的任务逻辑本身就是异步的(例如,任务内需要调用异步的LLM API、异步数据库),Arq是绝配。它用起来非常直观,Worker直接使用asyncio事件循环。
# tasks.py import asyncio from arq import create_pool, cron from arq.connections import RedisSettings async def long_ai_task(ctx, prompt: str): # ctx['redis'] 可以获取Redis连接 await asyncio.sleep(5) # 模拟长时间AI处理 result = f"Processed: {prompt}" return result # Worker配置 class WorkerSettings: redis_settings = RedisSettings(host="localhost") functions = [long_ai_task] cron_jobs = [] # 可以配置定时任务# 在FastAPI中触发任务 from arq import create_pool from .tasks import RedisSettings @app.on_event("startup") async def startup_event(): app.state.arq_pool = await create_pool(RedisSettings(host="redis")) @app.post("/submit-task/") async def submit_task(prompt: str): job = await app.state.arq_pool.enqueue_job("long_ai_task", prompt) return {"job_id": job.job_id} -
次选Dramatiq :如果你需要比RQ更强的功能(如中间件、速率限制),但又觉得Celery太重,且任务以同步CPU计算为主,Dramatiq是很好的选择。它通过多进程和Actor模型实现高性能。
-
慎用Celery :除非你的项目已经用了Celery,或者需要其非常高级的特性(如复杂路由、多个队列、Chord/Group等工作流),否则对于新的AI项目,Celery的配置复杂度和与异步代码的整合成本可能过高。
3.2 任务状态管理与结果回传
任务提交到队列后,我们需要让用户能查询进度和结果。一个常见的模式是使用Redis同时作为Broker和结果后端。
流程设计:
-
用户请求触发一个长任务,API立即返回一个唯一的
task_id。 -
API将任务放入队列(如Arq),并将
task_id与一个初始状态(如PENDING)存入Redis。 -
Worker从队列取出任务并执行,在执行过程中,通过
task_id更新Redis中的状态(如PROCESSING、PROGRESS: 50%)。 -
任务完成后,Worker将最终结果或错误信息存入Redis(状态更新为
SUCCESS或FAILED)。 -
用户通过另一个API端点,凭
task_id轮询查询任务状态和结果。
在FastAPI中的实现示例:
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
from enum import Enum
import uuid
import aioredis
app = FastAPI()
# 连接Redis
redis = aioredis.from_url("redis://localhost", decode_responses=True)
class TaskStatus(str, Enum):
PENDING = "pending"
PROCESSING = "processing"
SUCCESS = "success"
FAILED = "failed"
class TaskResponse(BaseModel):
task_id: str
status: TaskStatus
result: Optional[str] = None
error: Optional[str] = None
progress: Optional[int] = None
@app.post("/start-analysis/", response_model=TaskResponse)
async def start_analysis(data: dict):
task_id = str(uuid.uuid4())
# 1. 初始状态存入Redis,设置过期时间(如1小时)
initial_state = {
"status": TaskStatus.PENDING,
"result": None,
"error": None,
"progress": 0
}
await redis.hset(f"task:{task_id}", mapping=initial_state)
await redis.expire(f"task:{task_id}", 3600)
# 2. 将任务放入异步队列(这里以Arq为例)
# 假设app.state.arq_pool已在startup中创建
job = await app.state.arq_pool.enqueue_job("analyze_data_task", task_id, data)
# 3. 将Arq的job_id也关联存储,方便管理(可选)
await redis.hset(f"task:{task_id}", "arq_job_id", job.job_id)
return TaskResponse(task_id=task_id, status=TaskStatus.PENDING)
@app.get("/task/{task_id}", response_model=TaskResponse)
async def get_task_status(task_id: str):
# 从Redis中获取任务状态
task_data = await redis.hgetall(f"task:{task_id}")
if not task_data:
raise HTTPException(status_code=404, detail="Task not found")
return TaskResponse(
task_id=task_id,
status=task_data.get("status", TaskStatus.PENDING),
result=task_data.get("result"),
error=task_data.get("error"),
progress=int(task_data.get("progress", 0))
)
而在Worker端(Arq任务函数中),需要更新这个状态:
# 在tasks.py的long_ai_task中
async def analyze_data_task(ctx, task_id: str, data: dict):
redis = ctx["redis"]
try:
# 更新状态为处理中
await redis.hset(f"task:{task_id}", "status", "processing")
await redis.hset(f"task:{task_id}", "progress", 10)
# 模拟处理步骤1
await asyncio.sleep(2)
await redis.hset(f"task:{task_id}", "progress", 50)
# 模拟处理步骤2(调用AI模型等)
result = await call_llm_api(data["query"])
await redis.hset(f"task:{task_id}", "progress", 90)
# 处理完成,存储结果
await redis.hset(f"task:{task_id}", "status", "success")
await redis.hset(f"task:{task_id}", "result", result)
await redis.hset(f"task:{task_id}", "progress", 100)
except Exception as e:
# 处理失败,存储错误信息
await redis.hset(f"task:{task_id}", "status", "failed")
await redis.hset(f"task:{task_id}", "error", str(e))
raise # 让Arq也知道任务失败了
这样,一个完整的、可查询的异步任务流程就搭建起来了。前端可以通过轮询
/task/{task_id}
接口,或者更好的方式,使用WebSocket来接收实时状态更新。
4. 生产环境部署:从单机到可扩展集群
将开发好的FastAPI应用部署出去,并确保其稳定运行,需要一整套的考量。我们不再是用
uvicorn main:app --reload
这种开发命令了。
4.1 服务进程管理:Gunicorn with Uvicorn Workers
对于生产环境,我们需要一个更健壮的ASGI服务器。
uvicorn
本身是轻量级的,建议搭配
gunicorn
作为进程管理器,利用其成熟的热重启、负载均衡、进程管理功能。
为什么是Gunicorn + Uvicorn?
- Gunicorn :是一个WSGI/ASGI的进程管理器。它负责管理多个工作进程(Worker),处理请求分发、进程守护、优雅重启等。
-
Uvicorn Worker
:Gunicorn本身处理ASGI协议效率不高,我们需要使用
uvicorn.workers.UvicornWorker。这样,每个Gunicorn工作进程内部运行的是一个Uvicorn服务器实例,专门处理异步请求。
部署命令示例:
gunicorn main:app \
--workers 4 \ # 工作进程数,通常建议为 (CPU核心数 * 2) + 1
--worker-class uvicorn.workers.UvicornWorker \ # 关键:使用Uvicorn Worker
--bind 0.0.0.0:8000 \
--timeout 120 \ # 请求超时时间,对于长AI任务可以设长一些
--keep-alive 5 \
--access-logfile - \ # 访问日志输出到标准输出,方便容器收集
--error-logfile - \
--capture-output \
--log-level info
关键参数解析:
-
--workers:进程数。异步应用的特点是I/O密集型,所以可以设置比CPU核心数更多的Worker,以充分利用I/O等待时间。但也不是越多越好,需要根据实际负载测试。 -
--timeout:非常重要!默认是30秒。如果你的AI任务平均响应时间超过30秒,必须调大此值,否则Gunicorn会认为Worker僵死并将其杀掉。 -
--access-logfile和--error-logfile:设置为-表示输出到标准输出/错误,这是容器化部署的最佳实践,方便Docker或K8s收集日志。
4.2 容器化部署:Docker与最佳实践
容器化是现代化部署的标配。它能确保环境一致性,简化依赖管理。
一个生产可用的Dockerfile示例:
# 使用官方Python slim镜像作为基础,减少镜像体积
FROM python:3.11-slim as builder
# 安装编译依赖(如果需要编译某些Python包)
RUN apt-get update && apt-get install -y \
gcc \
g++ \
--no-install-recommends && \
rm -rf /var/lib/apt/lists/*
# 设置工作目录
WORKDIR /app
# 先复制依赖声明文件,利用Docker层缓存
COPY requirements.txt .
# 安装Python依赖(使用清华PyPI镜像加速)
RUN pip install --no-cache-dir -i https://pypi.tuna.tsinghua.edu.cn/simple -r requirements.txt
# 第二阶段:运行阶段
FROM python:3.11-slim
# 安装运行时可能需要的系统库(如SSL库)
RUN apt-get update && apt-get install -y \
curl \
--no-install-recommends && \
rm -rf /var/lib/apt/lists/*
# 创建非root用户运行应用,增强安全性
RUN useradd --create-home --shell /bin/bash appuser
USER appuser
WORKDIR /home/appuser/app
# 从构建阶段复制已安装的Python包
COPY --from=builder /usr/local/lib/python3.11/site-packages /usr/local/lib/python3.11/site-packages
COPY --from=builder /usr/local/bin /usr/local/bin
# 复制应用代码
COPY --chown=appuser:appuser . .
# 暴露端口
EXPOSE 8000
# 健康检查
HEALTHCHECK --interval=30s --timeout=3s --start-period=5s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1
# 使用Gunicorn启动应用
CMD ["gunicorn", "main:app", \
"--workers", "4", \
"--worker-class", "uvicorn.workers.UvicornWorker", \
"--bind", "0.0.0.0:8000", \
"--timeout", "120", \
"--access-logfile", "-", \
"--error-logfile", "-"]
最佳实践要点:
- 多阶段构建 :第一阶段安装编译依赖和Python包,第二阶段只复制运行所需的最小文件,大幅减小最终镜像体积。
- 使用非root用户 :避免以root权限运行容器,减少安全风险。
- 设置健康检查 :让容器编排平台(如K8s)能感知应用是否存活、是否就绪。
- 日志输出到标准流 :方便统一的日志收集系统(如ELK、Loki)进行处理。
-
使用
.dockerignore文件 :排除__pycache__、.git、虚拟环境目录等不必要的文件,加速构建。
4.3 配置管理与环境变量
生产环境的配置(如数据库连接串、LLM API密钥、第三方服务地址)绝不能硬编码在代码中。必须使用环境变量。
推荐使用Pydantic的
BaseSettings
进行配置管理:
# config.py
from pydantic_settings import BaseSettings
from typing import Optional
class Settings(BaseSettings):
# 应用配置
app_name: str = "My AI Service"
debug: bool = False
# Redis配置(用于缓存、任务队列)
redis_url: str = "redis://localhost:6379/0"
# 数据库配置
database_url: str
# LLM API配置
openai_api_key: Optional[str] = None
openai_base_url: Optional[str] = "https://api.openai.com/v1"
anthropic_api_key: Optional[str] = None
# 其他第三方服务
sentry_dsn: Optional[str] = None
# 从 `.env` 文件加载变量
class Config:
env_file = ".env"
env_file_encoding = 'utf-8'
case_sensitive = False # 环境变量不区分大小写
settings = Settings()
在代码中,通过
from config import settings
来使用配置,如
settings.redis_url
。环境变量可以来自系统的环境变量,也可以来自项目根目录的
.env
文件(开发环境使用,生产环境不应提交此文件)。
生产环境注入环境变量的方式:
-
Docker
:在
docker run命令中使用-e参数,或在docker-compose.yml的environment部分定义。 -
Kubernetes
:在Deployment的
env字段或使用ConfigMap/Secret。 - 云平台 :如AWS ECS、Google Cloud Run等都提供了便捷的环境变量配置界面。
4.4 监控、日志与告警
服务上线后,必须要有眼睛盯着它。
1. 结构化日志:
不要再用简单的
print
了。使用
structlog
或标准的
logging
模块配置JSON格式的日志,方便日志分析系统(如ELK、Loki+Grafana)进行解析和查询。
# logging_config.py
import logging
import sys
from pythonjsonlogger import jsonlogger
# 配置JSON格式的日志处理器
handler = logging.StreamHandler(sys.stdout)
formatter = jsonlogger.JsonFormatter(
'%(asctime)s %(name)s %(levelname)s %(message)s %(module)s %(funcName)s'
)
handler.setFormatter(formatter)
# 获取根日志记录器并配置
root_logger = logging.getLogger()
root_logger.addHandler(handler)
root_logger.setLevel(logging.INFO)
# 在你的应用代码中
import logging
logger = logging.getLogger(__name__)
async def some_async_function():
try:
# ... 业务逻辑
logger.info("LLM API call succeeded", extra={"model": "gpt-4", "duration_ms": 1200})
except Exception as e:
logger.error("LLM API call failed",
exc_info=True, # 自动记录异常堆栈
extra={"error_type": type(e).__name__, "prompt_preview": prompt[:100]})
2. 应用性能监控: 集成像 Sentry 这样的错误追踪工具,它能自动捕获未处理的异常,并附带丰富的上下文信息(如请求参数、用户信息、环境变量),极大加速线上问题的排查。 对于性能指标(如接口响应时间、LLM API调用延迟、队列长度),可以使用 Prometheus 客户端库暴露指标,然后通过Grafana进行可视化。
3. 健康检查端点:
为你的FastAPI应用添加一个
/health
端点,用于检查应用本身及其关键依赖(如数据库、Redis、外部API)的状态。这被容器编排平台和负载均衡器广泛使用。
from fastapi import Depends
from sqlalchemy.ext.asyncio import AsyncSession
from redis import asyncio as aioredis
import httpx
@app.get("/health")
async def health_check(
db: AsyncSession = Depends(get_db),
redis: aioredis.Redis = Depends(get_redis)
):
checks = {}
# 检查数据库
try:
await db.execute("SELECT 1")
checks["database"] = "healthy"
except Exception as e:
checks["database"] = f"unhealthy: {e}"
# 检查Redis
try:
await redis.ping()
checks["redis"] = "healthy"
except Exception as e:
checks["redis"] = f"unhealthy: {e}"
# 检查关键外部API(可选,注意频率)
# async with httpx.AsyncClient() as client:
# try:
# resp = await client.get("https://api.openai.com/v1/models", timeout=5.0)
# checks["openai_api"] = "healthy" if resp.status_code == 200 else f"unhealthy: {resp.status_code}"
# except Exception as e:
# checks["openai_api"] = f"unhealthy: {e}"
overall_status = "healthy" if all(v == "healthy" for v in checks.values()) else "unhealthy"
return {"status": overall_status, "details": checks}
5. 进阶话题:应对高并发与LLM API限流
当你的服务用户量增长,或者遇到LLM服务商严格的速率限制时,简单的重试和降级可能不够,需要更系统的策略。
5.1 请求排队与速率限制
如果LLM API的并发限制是每分钟N次,而你的用户请求可能超过这个数,你就需要在服务端实现一个 请求队列和速率限制器 。
方案:使用Redis实现令牌桶算法 令牌桶算法是一个经典且灵活的限流算法。我们可以为每个LLM API端点(或每个用户)维护一个“令牌桶”。
import asyncio
import time
import aioredis
class RateLimiter:
def __init__(self, redis_client, key_prefix, max_tokens, refill_rate):
"""
:param redis_client: aioredis客户端
:param key_prefix: 限流键前缀,如 `rate_limit:openai:chat`
:param max_tokens: 桶容量
:param refill_rate: 每秒补充的令牌数
"""
self.redis = redis_client
self.key_prefix = key_prefix
self.max_tokens = max_tokens
self.refill_rate = refill_rate
async def _get_bucket_key(self, identifier: str):
return f"{self.key_prefix}:{identifier}"
async def acquire(self, identifier: str, tokens=1, timeout=10):
"""
尝试获取令牌,如果成功返回True,否则等待或超时返回False
"""
bucket_key = await self._get_bucket_key(identifier)
lua_script = """
local key = KEYS[1]
local max_tokens = tonumber(ARGV[1])
local refill_rate = tonumber(ARGV[2])
local tokens_requested = tonumber(ARGV[3])
local now = tonumber(ARGV[4])
local bucket = redis.call('HMGET', key, 'tokens', 'last_refill')
local current_tokens = max_tokens
local last_refill = now
if bucket[1] then
current_tokens = tonumber(bucket[1])
last_refill = tonumber(bucket[2])
end
-- 计算需要补充的令牌
local time_passed = now - last_refill
local tokens_to_add = math.floor(time_passed * refill_rate)
local new_tokens = math.min(max_tokens, current_tokens + tokens_to_add)
if new_tokens >= tokens_requested then
-- 有足够令牌,消耗它们
new_tokens = new_tokens - tokens_requested
redis.call('HMSET', key, 'tokens', new_tokens, 'last_refill', now)
redis.call('EXPIRE', key, math.ceil(max_tokens / refill_rate) + 10) -- 设置合理的过期时间
return 1 -- 成功
else
-- 令牌不足,计算需要等待的时间
local tokens_needed = tokens_requested - new_tokens
local wait_time = tokens_needed / refill_rate
redis.call('HMSET', key, 'tokens', new_tokens, 'last_refill', now)
redis.call('EXPIRE', key, math.ceil(max_tokens / refill_rate) + 10)
return wait_time -- 返回需要等待的秒数
end
"""
start_time = time.time()
while time.time() - start_time < timeout:
now = time.time()
# 使用Lua脚本保证原子性操作
result = await self.redis.eval(
lua_script, 1, bucket_key,
self.max_tokens, self.refill_rate, tokens, now
)
if result == 1:
return True # 成功获取令牌
else:
# result是需要等待的秒数
wait_time = float(result)
await asyncio.sleep(min(wait_time, 0.1)) # 短暂休眠后重试
return False # 超时未获取到
# 使用示例:限制每个用户每分钟最多调用10次ChatGPT
limiter = RateLimiter(redis, "rate_limit:openai:chat", max_tokens=10, refill_rate=10/60) # 每分钟10个,即每6秒1个
async def call_llm_with_rate_limit(user_id: str, prompt: str):
if not await limiter.acquire(user_id, tokens=1, timeout=30):
raise Exception("Rate limit exceeded. Please try again later.")
# 调用真正的LLM API
return await call_openai_chat(prompt)
这个方案将限流逻辑放在你的服务端,而不是每个客户端,实现了集中控制。你可以根据API Key、用户ID、IP地址等不同维度进行限流。
5.2 异步流式响应:提升长文本生成体验
对于需要生成长篇内容的AI服务(如写报告、生成代码),等待全部内容生成完再一次性返回,用户体验很差。FastAPI支持 流式响应 ,可以逐块(chunk)地将内容推送给客户端。
实现SSE(Server-Sent Events)流式响应:
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import asyncio
import json
app = FastAPI()
async def fake_llm_stream_generator(prompt: str):
"""
模拟一个流式返回的LLM。
实际中,这里应该调用支持流式响应的LLM API(如OpenAI的stream=True)。
"""
# 模拟分块生成
simulated_chunks = [
f"思考用户的问题:'{prompt}'...\n\n",
"首先,我们需要理解这个问题的核心。",
"它涉及到几个关键点。",
"第一点,...",
"第二点,...",
"\n\n以上就是我的分析。"
]
for chunk in simulated_chunks:
# 模拟每块生成需要一点时间
await asyncio.sleep(0.5)
# SSE格式要求:`data: <content>\n\n`
yield f"data: {json.dumps({'content': chunk})}\n\n"
@app.post("/chat/stream")
async def chat_stream(prompt: str):
generator = fake_llm_stream_generator(prompt)
return StreamingResponse(
generator,
media_type="text/event-stream", # SSE的媒体类型
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no", # 禁用Nginx等代理的缓冲
}
)
前端可以使用
EventSource
API来接收这个流:
const eventSource = new EventSource(`/chat/stream?prompt=${encodeURIComponent(userInput)}`);
eventSource.onmessage = (event) => {
const data = JSON.parse(event.data);
// 将data.content逐步追加到页面上
document.getElementById('output').innerHTML += data.content;
};
eventSource.onerror = (error) => {
console.error('Stream error:', error);
eventSource.close();
};
流式响应不仅能极大提升用户体验(感觉响应更快),还能在生成过程中就发现错误并中断,避免用户长时间等待后得到一个失败结果。
5.3 负载均衡与水平扩展
当单台服务器无法承受流量时,就需要水平扩展。对于无状态的FastAPI应用,这相对简单。
架构要点:
- 多副本部署 :在Kubernetes或Docker Swarm中,可以轻松启动多个应用副本(Pod/容器)。
- 负载均衡器 :使用Nginx、HAProxy或云服务商(如AWS ALB、GCP Cloud Load Balancing)的负载均衡器,将流量分发到各个副本。
- 共享状态外置 :确保应用本身是无状态的。所有需要共享的数据(如Session、缓存、任务队列)必须存储在外部的中心化服务中,如Redis、PostgreSQL。 绝对不能 存在本地内存的状态。
-
健康检查
:负载均衡器需要配置健康检查端点(如我们之前实现的
/health),自动将不健康的实例从流量池中剔除。
一个简单的Nginx配置示例:
upstream ai_backend {
# 假设你的FastAPI容器在8000端口,且运行了3个副本
server host1:8000;
server host2:8000;
server host3:8000;
# 可以配置负载均衡策略,如least_conn;
}
server {
listen 80;
server_name your-ai-service.com;
location / {
proxy_pass http://ai_backend;
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
# 重要:对于流式响应,需要禁用代理缓冲
proxy_buffering off;
proxy_cache off;
}
# 可选:静态文件服务
location /static/ {
alias /path/to/your/static/files/;
}
}
通过这套组合拳——异步处理、任务队列、容器化、监控告警、限流排队、流式响应和水平扩展——你的AI服务就具备了面向真实生产环境挑战的能力。这不再是那个在本地跑得欢的“玩具”,而是一个真正能服务用户、稳定可靠的“产品”了。
更多推荐
所有评论(0)