You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Python FastAPI Web应用中实现生产级API请求排队

生产级FastAPI请求排队机制实现方案

问题背景

我们有一个FastAPI应用,需接收远程服务器数据并升级为生产级应用,核心需求如下:

  • 并发入站API请求峰值约20万
  • 使用SQLAlchemy连接Postgres/Oracle,数据库并发连接数有限制
  • 需支持根据负载动态扩展

当前应用部署在Windows Server 2019(无Docker),尝试用asyncio.Queue实现排队后,仅100个并发请求就导致应用冻结,调用方出现60秒超时。现有代码见下文。

现有实现的核心问题

  1. 单worker瓶颈:仅启动1个worker,所有请求串行处理,完全无法应对高并发,直接导致队列积压、请求超时。
  2. 同步SQLAlchemy阻塞事件循环:retrieve_device、update_device_by_id等操作是同步阻塞的,会卡住asyncio事件循环,让整个应用失去响应。
  3. 无边界队列导致内存溢出:asyncio.Queue默认无上限,大量请求入队会耗尽服务器内存,最终引发应用崩溃。
  4. DB会话与请求上下文冲突:路由中通过Depends(get_db)获取的DB会话绑定到当前请求上下文,在队列worker中复用会引发线程安全问题(同步SQLAlchemy会话并非协程安全)。
  5. 响应对象闭包引用问题: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()

部署与监控建议

  1. Windows服务运行:用pywin32将FastAPI应用注册为Windows服务,确保进程稳定运行。
  2. 队列监控:定期记录队列长度、worker状态,当队列长度超过阈值时触发告警。
  3. DB连接监控:监控数据库连接数,确保不超过上限,避免DB拒绝连接。
  4. 日志优化:为每个请求生成唯一ID,关联队列任务日志,便于排查问题。
  5. 负载测试:用locust或k6进行高并发负载测试,验证排队机制的稳定性。

内容的提问来源于stack exchange,提问作者inbhupesh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.03 06:09:55