如何在FastAPI应用中部署非阻塞SQS监听器(无需Celery)
解决FastAPI中SQS监听器阻塞主事件循环的问题
核心问题分析
你的代码存在两个关键阻塞点:
- 如果
receive_messages是同步调用(比如基于boto3的同步方法),会直接卡住事件循环,完全占用主线程,导致API端点无法响应。 - 无限
while(True)循环没有给事件循环让出执行权的逻辑,即便用了异步调用,也会持续占用事件循环资源,影响API请求处理。
无需Celery的解决方案
1. 替换为异步SQS客户端
使用aiobotocore替代同步的boto3,确保所有SQS操作都是非阻塞的异步调用。先安装依赖:
pip install aiobotocore
2. 改造监听器代码
修改监听器逻辑,保证所有操作异步化,同时在循环中加入异步睡眠,主动让出事件循环:
import asyncio import logging from aiobotocore.session import get_session logger = logging.getLogger(__name__) async def create_queue(queue_name, region="your-region"): session = get_session() async with session.create_client('sqs', region_name=region) as client: resp = await client.create_queue(QueueName=queue_name) return resp['QueueUrl'] async def receive_messages(queue_url, max_num=5, wait_time=10, region="your-region"): session = get_session() async with session.create_client('sqs', region_name=region) as client: resp = await client.receive_message( QueueUrl=queue_url, MaxNumberOfMessages=max_num, WaitTimeSeconds=wait_time # 长轮询,减少空轮询次数 ) return resp.get('Messages', []) async def process_message(queue_url, message, region="your-region"): # 这里替换为你的业务逻辑,比如调用原有的search功能 logger.info(f"Processing message content: {message['Body']}") # 处理完成后删除队列中的消息 session = get_session() async with session.create_client('sqs', region_name=region) as client: await client.delete_message( QueueUrl=queue_url, ReceiptHandle=message['ReceiptHandle'] ) async def sqs_listener(): logger.info("SQS Listener started") queue_url = await create_queue("testQueue") logger.info(f"Listening on queue: {queue_url}") while True: try: messages = await receive_messages(queue_url) if messages: # 并行处理多条消息(可根据业务调整为串行) tasks = [process_message(queue_url, msg) for msg in messages] await asyncio.gather(*tasks) # 每次轮询后短暂睡眠,让出事件循环给API请求 await asyncio.sleep(0.1) except Exception as e: logger.error(f"Listener error occurred: {str(e)}", exc_info=True) # 出错后延迟重试,避免频繁报错 await asyncio.sleep(5)
3. 正确配置Lifespan上下文
在lifespan中创建监听器任务,并在应用关闭时优雅取消任务:
from fastapi import FastAPI import asyncio app = FastAPI(lifespan=lifespan) @asynccontextmanager async def lifespan(app: FastAPI): # 创建并启动监听器任务 listener_task = asyncio.create_task(sqs_listener()) yield # 应用关闭时取消任务并等待结束 listener_task.cancel() try: await listener_task except asyncio.CancelledError: logger.info("SQS Listener stopped gracefully") @app.get("/health") async def health_check(): return {"status": "healthy"} @app.post("/search") async def search(query: str): # 保留原有的搜索测试端点 return {"result": f"Search results for query: {query}"}
关键优化说明
- 异步客户端:
aiobotocore提供完全异步的AWS SDK调用,不会阻塞事件循环。 - 长轮询:
WaitTimeSeconds=10让SQS在有消息时立即返回,无消息时最多等待10秒,减少不必要的空轮询。 - 主动出让事件循环:
await asyncio.sleep(0.1)确保API端点能及时获取事件循环资源处理请求。 - 异常容错:监听器内置异常捕获和重试逻辑,避免单个错误导致整个监听器崩溃。
- 优雅停止:在lifespan结束时取消任务,保证应用关闭时监听器能正常终止。
内容的提问来源于stack exchange,提问作者Vikas Palakkat
相关产品推荐
相关产品推荐

