无需消息队列实现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
相关产品推荐
相关产品推荐

