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

无需消息队列实现FastAPI请求排队及异步处理方案问询

使用asyncio.Queue实现FastAPI语音识别请求队列化处理

需求说明

  • 接收含语音文件的压缩包及结果回传URL后,立即返回200状态码告知请求已启动
  • 语音识别为耗时操作,需将新请求存入队列,按顺序依次处理
  • 处理完成后,将结果发送至指定回传URL
  • 不使用Kafka、RabbitMQ等消息中间件,基于asyncio.Queue实现

现有代码

from fastapi import FastAPI, File, UploadFile, BackgroundTasks, status
from fastapi.responses import JSONResponse
import logging
import time
import os
import rarfile
import zipfile

# 假设已定义以下变量和函数:UPLOADED_FILES_PATH, split_filename, is_archive_file, save_file_to_uploads, unrar_files, unzip_files, recognition_wav, model

app = FastAPI()

@app.post("/uprecognize", tags=["Upload and recognize"], status_code=status.HTTP_200_OK)
async def upload_recognize(
        url_for_request: str,
        background_tasks: BackgroundTasks,
        file: UploadFile = File(...),
                        ):
    logger = logging.getLogger(__name__)
    full_name = split_filename(file)
    if not is_archive_file(file):
        logger.error("File must be RAR or ZIP format")
        return JSONResponse(content={'msg': 'File must be RAR or ZIP format'}, status_code=status.HTTP_400_BAD_REQUEST)
    else:
        start = time.time()
        await save_file_to_uploads(file, full_name)
        end = time.time()
    if not os.path.exists(UPLOADED_FILES_PATH + '/' + os.path.splitext(full_name)[0]):
        os.mkdir(UPLOADED_FILES_PATH + '/' + os.path.splitext(full_name)[0])
    if os.path.exists(UPLOADED_FILES_PATH + '/' + full_name) and rarfile.is_rarfile(UPLOADED_FILES_PATH + '/' + full_name):
        unrar_files(UPLOADED_FILES_PATH + '/' + full_name)
    elif os.path.exists(UPLOADED_FILES_PATH + '/' + full_name) and zipfile.is_zipfile(UPLOADED_FILES_PATH + '/' + full_name):
        unzip_files(UPLOADED_FILES_PATH + '/' + full_name)
    else:
        logger.error("File not found")
        return JSONResponse(content={'msg': 'File not found'}, status_code=status.HTTP_404_NOT_FOUND)
    background_tasks.add_task(recognition_wav, full_name, logger, model, url_for_request)
    return JSONResponse(content={'msg':'Start recognition'},
                        status_code=status.HTTP_200_OK,
                        background=background_tasks)

实现方案

1. 初始化队列并启动消费者协程

在FastAPI启动时创建asyncio.Queue,并启动一个长期运行的消费者协程,负责从队列中取任务并依次处理。

import asyncio
from contextlib import asynccontextmanager

# 初始化队列,设置最大容量防止内存溢出
task_queue = asyncio.Queue(maxsize=100)

@asynccontextmanager
async def lifespan(app: FastAPI):
    # 启动消费者协程
    consumer_task = asyncio.create_task(consumer())
    yield
    # 关闭应用时取消消费者协程
    consumer_task.cancel()
    await consumer_task

app = FastAPI(lifespan=lifespan)

2. 修改接口逻辑:将任务加入队列后立即返回

调整原接口逻辑,仅完成文件校验和保存,将任务参数加入队列后立即返回响应,避免阻塞请求。

@app.post("/uprecognize", tags=["Upload and recognize"], status_code=status.HTTP_200_OK)
async def upload_recognize(
        url_for_request: str,
        file: UploadFile = File(...),
                        ):
    logger = logging.getLogger(__name__)
    full_name = split_filename(file)
    
    # 校验文件格式
    if not is_archive_file(file):
        logger.error("File must be RAR or ZIP format")
        return JSONResponse(content={'msg': 'File must be RAR or ZIP format'}, status_code=status.HTTP_400_BAD_REQUEST)
    
    # 保存文件到本地
    try:
        await save_file_to_uploads(file, full_name)
    except Exception as e:
        logger.error(f"Failed to save file: {str(e)}")
        return JSONResponse(content={'msg': 'Failed to save file'}, status_code=status.HTTP_500_INTERNAL_SERVER_ERROR)
    
    # 将任务加入队列,立即返回响应
    await task_queue.put({
        "full_name": full_name,
        "url_for_request": url_for_request,
        "logger": logger
    })
    
    return JSONResponse(content={'msg': 'Recognition task added to queue'}, status_code=status.HTTP_200_OK)

3. 实现消费者协程:依次处理队列任务

消费者协程循环从队列取任务,执行解压、语音识别、结果回传操作。注意将同步阻塞操作转为异步执行,避免阻塞事件循环。

async def consumer():
    logger = logging.getLogger(__name__)
    while True:
        try:
            # 从队列获取任务,队列为空时阻塞等待
            task = await task_queue.get()
            full_name = task["full_name"]
            url_for_request = task["url_for_request"]
            task_logger = task["logger"]
            
            try:
                # 创建解压目录
                extract_dir = os.path.join(UPLOADED_FILES_PATH, os.path.splitext(full_name)[0])
                if not os.path.exists(extract_dir):
                    os.makedirs(extract_dir)
                
                # 解压文件:同步操作转异步执行
                file_path = os.path.join(UPLOADED_FILES_PATH, full_name)
                if rarfile.is_rarfile(file_path):
                    await asyncio.to_thread(unrar_files, file_path)
                elif zipfile.is_zipfile(file_path):
                    await asyncio.to_thread(unzip_files, file_path)
                else:
                    task_logger.error(f"Invalid archive file: {full_name}")
                    continue
                
                # 执行语音识别:同步操作转异步执行
                await asyncio.to_thread(recognition_wav, full_name, task_logger, model, url_for_request)
                
            except Exception as e:
                task_logger.error(f"Failed to process task {full_name}: {str(e)}")
            finally:
                # 标记任务完成,更新队列内部计数
                task_queue.task_done()
                
        except asyncio.CancelledError:
            logger.info("Consumer task cancelled")
            break
        except Exception as e:
            logger.error(f"Consumer error: {str(e)}")

关键注意事项

  • 同步操作异步化:解压、语音识别等阻塞IO操作必须用asyncio.to_thread包装,否则会阻塞整个FastAPI事件循环,导致所有请求无法响应。
  • 队列容量限制:初始化asyncio.Queue时设置maxsize,防止大量请求堆积导致内存耗尽。队列满时await task_queue.put()会阻塞,可按需添加超时处理或返回"队列已满"提示。
  • 异常处理:消费者协程需捕获异常,避免单个任务失败导致整个消费者崩溃;接口中也要处理文件保存失败的情况。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 09:15:47