如何在FastAPI中不使用BackgroundTasks异步处理SQS消息?
FastAPI 实现独立后台轮询SQS任务的方案
直接在__init__.py中运行同步轮询任务会阻塞FastAPI主线程,导致无法处理HTTP请求——因为主线程被循环占用,无法启动ASGI服务。以下是几种在FastAPI项目内实现独立后台轮询的可行方案:
方案1:后台线程(简单快速)
利用Python标准库threading,在FastAPI启动时启动一个守护线程处理SQS轮询,不阻塞主线程。
from fastapi import FastAPI import boto3 import threading import time app = FastAPI() # 初始化SQS客户端 sqs = boto3.client('sqs', region_name='your-region') QUEUE_URL = 'your-sqs-queue-url' def poll_and_process_sqs(): while True: # 长轮询拉取消息(减少空请求频率) response = sqs.receive_message( QueueUrl=QUEUE_URL, MaxNumberOfMessages=10, WaitTimeSeconds=20 ) if 'Messages' in response: for msg in response['Messages']: # 替换为你的消息处理逻辑 handle_message(msg) # 删除已处理完成的消息 sqs.delete_message( QueueUrl=QUEUE_URL, ReceiptHandle=msg['ReceiptHandle'] ) # 短间隔避免CPU空转(长轮询已设置等待时间,可根据需求调整) time.sleep(1) def handle_message(message): try: print(f"Processing event: {message['Body']}") # 这里写实际的业务处理代码 except Exception as e: print(f"Failed to process message: {str(e)}") # 可选:将失败消息转入死信队列或重新入队 # 在FastAPI启动时启动后台线程 @app.on_event("startup") def start_sqs_poller(): # 守护线程:主进程退出时自动终止 poll_thread = threading.Thread(target=poll_and_process_sqs, daemon=True) poll_thread.start() # 你的事件接收接口 @app.post("/api/event") async def accept_event(): # 这里添加请求验证逻辑 # 验证通过后发送消息到SQS sqs.send_message( QueueUrl=QUEUE_URL, MessageBody="your-event-payload" ) return {"status": "accepted"}
方案2:Lifespan上下文管理器(优雅管理生命周期)
使用FastAPI的Lifespan特性,可以更优雅地控制后台任务的启动与停止,适合需要在服务关闭时做清理的场景。
from fastapi import FastAPI import boto3 import threading import time from contextlib import asynccontextmanager # 控制轮询线程运行状态 _is_running = True @asynccontextmanager async def app_lifespan(app: FastAPI): # 启动阶段:初始化资源、启动轮询线程 sqs = boto3.client('sqs', region_name='your-region') QUEUE_URL = 'your-sqs-queue-url' def poll_and_process_sqs(): while _is_running: response = sqs.receive_message( QueueUrl=QUEUE_URL, MaxNumberOfMessages=10, WaitTimeSeconds=20 ) if 'Messages' in response: for msg in response['Messages']: handle_message(msg) sqs.delete_message( QueueUrl=QUEUE_URL, ReceiptHandle=msg['ReceiptHandle'] ) time.sleep(1) poll_thread = threading.Thread(target=poll_and_process_sqs, daemon=True) poll_thread.start() yield # 服务运行阶段 # 关闭阶段:停止轮询 global _is_running _is_running = False app = FastAPI(lifespan=app_lifespan) def handle_message(message): try: print(f"Processing event: {message['Body']}") # 业务处理逻辑 except Exception as e: print(f"Message processing failed: {str(e)}") @app.post("/api/event") async def accept_event(): # 请求验证逻辑 # 发送消息到SQS sqs.send_message( QueueUrl=QUEUE_URL, MessageBody="your-event-payload" ) return {"status": "accepted"}
关键注意事项
- 守护线程特性:守护线程会随主进程终止而结束,如果需要确保正在处理的消息完成,可在关闭时添加
poll_thread.join(timeout=30)等逻辑,但注意不要阻塞服务关闭流程过久。 - SQS长轮询:设置
WaitTimeSeconds(建议10-20秒)可以减少空请求次数,降低API调用成本并提升性能。 - 异常处理:必须在消息处理逻辑中添加异常捕获,避免单个消息处理失败导致整个轮询线程崩溃。
内容的提问来源于stack exchange,提问作者offensivedev
相关产品推荐
相关产品推荐

