如何在Python FastAPI Web应用中实现生产级API请求排队
生产级FastAPI请求排队机制实现方案
问题背景
我们有一个FastAPI应用,需接收远程服务器数据并升级为生产级应用,核心需求如下:
- 并发入站API请求峰值约20万
- 使用SQLAlchemy连接Postgres/Oracle,数据库并发连接数有限制
- 需支持根据负载动态扩展
当前应用部署在Windows Server 2019(无Docker),尝试用asyncio.Queue实现排队后,仅100个并发请求就导致应用冻结,调用方出现60秒超时。现有代码见下文。
现有实现的核心问题
- 单worker瓶颈:仅启动1个worker,所有请求串行处理,完全无法应对高并发,直接导致队列积压、请求超时。
- 同步SQLAlchemy阻塞事件循环:
retrieve_device、update_device_by_id等操作是同步阻塞的,会卡住asyncio事件循环,让整个应用失去响应。 - 无边界队列导致内存溢出:
asyncio.Queue默认无上限,大量请求入队会耗尽服务器内存,最终引发应用崩溃。 - DB会话与请求上下文冲突:路由中通过
Depends(get_db)获取的DB会话绑定到当前请求上下文,在队列worker中复用会引发线程安全问题(同步SQLAlchemy会话并非协程安全)。 - 响应对象闭包引用问题:
create_device_task中直接修改路由的response对象,跨协程操作会导致不可预期的行为。
生产级实现方案
1. 启动多Worker池,匹配数据库连接数
根据数据库最大并发连接数(比如Postgres默认100),启动对应数量的worker,确保DB连接不超限。同时设置队列最大长度,触发背压机制。
2. 改用异步SQLAlchemy
同步SQLAlchemy会阻塞事件循环,必须切换到异步模式,使用asyncpg(Postgres)或cx_Oracle的异步适配,确保DB操作不阻塞协程。
3. 实现队列限流与背压
设置队列最大容量,当队列满时,返回503 Service Unavailable给新请求,避免服务器资源耗尽。
4. 动态调整Worker数量
监控队列长度和系统负载,动态增减worker数量(比如队列长度超过阈值时临时扩容,空闲时缩容)。
5. 请求超时与重试机制
为队列中的任务设置超时时间,超时的任务自动重试或返回错误,避免请求无限等待。
6. Windows下Asyncio优化
Windows的asyncio默认使用SelectorEventLoop,性能较差,建议切换为ProactorEventLoop(Python 3.8+支持),提升协程处理效率。
完整代码实现
queue_module.py
import asyncio from typing import Callable, Awaitable from functools import partial # 配置参数:根据DB最大连接数调整 MAX_WORKERS = 50 MAX_QUEUE_SIZE = 10000 fifo_queue = asyncio.Queue(maxsize=MAX_QUEUE_SIZE) worker_tasks = [] async def db_worker(worker_id: int, get_async_db_session: Callable): print(f"Starting DB Worker {worker_id}") while True: task_func, result_future, timeout = await fifo_queue.get() try: # 为每个任务创建独立的异步DB会话 async with get_async_db_session() as session: # 绑定会话到任务函数 task = partial(task_func, db=session) # 设置任务超时 result = await asyncio.wait_for(task(), timeout=timeout) result_future.set_result(result) except asyncio.TimeoutError: result_future.set_exception(Exception("Task timed out")) except Exception as e: result_future.set_exception(e) finally: fifo_queue.task_done() def start_workers(get_async_db_session: Callable): """启动worker池""" global worker_tasks for i in range(MAX_WORKERS): task = asyncio.create_task(db_worker(i, get_async_db_session)) worker_tasks.append(task) def stop_workers(): """停止所有worker""" for task in worker_tasks: task.cancel()
r_device.py 路由改造
from fastapi import status, Depends, Request, Response, HTTPException from sqlalchemy.ext.asyncio import AsyncSession from datetime import datetime import asyncio # 假设已配置异步DB会话工厂 async_session_factory = ... async def get_async_db() -> AsyncSession: async with async_session_factory() as session: yield session @router.post("/", response_model=Show_Device) async def create_device( request: Create_Device, log_request: Request, current_user: User = Depends(get_current_user) ): async def create_device_task(db: AsyncSession): api_logger.info( f"Create Device Request Received", extra={"ids": {"id": request.id}, "action": log_actions["0"], "caller": await caller_details(log_request)} ) # 异步查询 existing_device = await retrieve_device_async(id=request.id, db=db, logger=api_logger) if existing_device: api_logger.info( f"Device Present in Database, Updating Device", extra={"action": log_actions["2"], "ids": {"id": request.id}} ) updated_device = await update_device_by_id_async( id=request.id, device=request, db=db, logger=api_logger ) api_logger.info( f"Device Updated, Returning Updated Device", extra={"ids": {"id": request.id}, "action": log_actions["3"]} ) return {"data": updated_device, "status_code": status.HTTP_200_OK} else: request.create_date = datetime.now().strftime('%Y-%m-%d %H:%M:%S') request.last_detected_date = datetime.now().strftime('%Y-%m-%d %H:%M:%S') api_logger.info( f"Device Not Present in Database Creating New Device", extra={"action": log_actions["1"], "ids": {"id": request.id}} ) device = await create_new_device_async( device=request, db=db, logger=api_logger ) api_logger.info( f"Device Created, Returning New Device", extra={"ids": {"id": request.id}, "action": log_actions["3"]} ) return {"data": device, "status_code": status.HTTP_201_CREATED} result_future = asyncio.Future() try: # 入队时设置超时时间(比如5秒) await asyncio.wait_for( fifo_queue.put((create_device_task, result_future, 30)), timeout=5 ) except asyncio.TimeoutError: raise HTTPException( status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail="Service busy, please try again later" ) try: result = await asyncio.wait_for(result_future, timeout=60) return result["data"] except asyncio.TimeoutError: raise HTTPException( status_code=status.HTTP_504_GATEWAY_TIMEOUT, detail="Request processing timed out" )
启动时初始化Worker
在FastAPI应用启动事件中启动worker:
from fastapi import FastAPI from queue_module import start_workers, get_async_db app = FastAPI() @app.on_event("startup") async def startup_event(): # Windows下切换为ProactorEventLoop asyncio.set_event_loop_policy(asyncio.WindowsProactorEventLoopPolicy()) start_workers(get_async_db) @app.on_event("shutdown") async def shutdown_event(): stop_workers()
部署与监控建议
- Windows服务运行:用
pywin32将FastAPI应用注册为Windows服务,确保进程稳定运行。 - 队列监控:定期记录队列长度、worker状态,当队列长度超过阈值时触发告警。
- DB连接监控:监控数据库连接数,确保不超过上限,避免DB拒绝连接。
- 日志优化:为每个请求生成唯一ID,关联队列任务日志,便于排查问题。
- 负载测试:用
locust或k6进行高并发负载测试,验证排队机制的稳定性。
内容的提问来源于stack exchange,提问作者inbhupesh
相关产品推荐
相关产品推荐

