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等代理的缓冲 } )前端可以使用EventSourceAPI来接收这个流:
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服务就具备了面向真实生产环境挑战的能力。这不再是那个在本地跑得欢的“玩具”,而是一个真正能服务用户、稳定可靠的“产品”了。