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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 13:35:20