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

如何在FastAPI应用中部署非阻塞SQS监听器(无需Celery)

解决FastAPI中SQS监听器阻塞主事件循环的问题

核心问题分析

你的代码存在两个关键阻塞点:

  1. 如果receive_messages是同步调用(比如基于boto3的同步方法),会直接卡住事件循环,完全占用主线程,导致API端点无法响应。
  2. 无限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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 02:36:12