如何用Python、FastAPI、SQLAlchemy异步处理REST请求及任务状态更新?
FastAPI 异步后台任务实现方案
核心问题分析
你之前的代码仅实现部分异步,本质原因是阻塞操作(比如time.sleep或CPU密集计算)占用了FastAPI的事件循环,导致后续请求需要等待当前阻塞任务释放资源。要实现真正的异步,必须把耗时/阻塞的计算任务从事件循环中剥离,放到单独的线程、进程或者任务队列中执行。
实现步骤与代码示例
1. 数据库模型定义(SQLAlchemy)
首先定义存储任务状态的模型,包含必要的状态字段:
from sqlalchemy import Column, Integer, String, Float, DateTime, Enum from sqlalchemy.ext.declarative import declarative_base from datetime import datetime import enum Base = declarative_base() class TaskStatus(str, enum.Enum): PENDING = "PENDING" RUNNING = "RUNNING" COMPLETED = "COMPLETED" FAILED = "FAILED" class CalculationTask(Base): __tablename__ = "calculation_tasks" id = Column(Integer, primary_key=True, index=True) base = Column(Float, nullable=False) exponent = Column(Float, nullable=False) status = Column(Enum(TaskStatus), default=TaskStatus.PENDING) result = Column(Float, nullable=True) created_at = Column(DateTime, default=datetime.utcnow) updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow)
2. 异步任务处理逻辑
使用concurrent.futures.ThreadPoolExecutor把阻塞计算任务放到线程池执行,避免阻塞FastAPI的事件循环。同时实时更新数据库中的任务状态:
import time from sqlalchemy.orm import Session from concurrent.futures import ThreadPoolExecutor # 全局线程池,可根据服务器配置调整大小 executor = ThreadPoolExecutor(max_workers=4) def calculate_and_update_task(db: Session, task_id: int, base: float, exponent: float): # 更新任务为运行中 task = db.query(CalculationTask).filter(CalculationTask.id == task_id).first() task.status = TaskStatus.RUNNING db.commit() try: # 模拟随指数变化的延迟:指数越大,延迟越长 delay = exponent * 0.5 time.sleep(delay) # 执行计算 result = base ** exponent # 更新任务结果与状态 task.result = result task.status = TaskStatus.COMPLETED db.commit() except Exception as e: task.status = TaskStatus.FAILED db.commit()
3. FastAPI接口实现
POST接口负责创建任务并提交到后台线程,GET接口查询任务状态:
from fastapi import FastAPI, Depends, BackgroundTasks from sqlalchemy.orm import Session from sqlalchemy import create_engine from sqlalchemy.orm import sessionmaker # 数据库连接配置 SQLALCHEMY_DATABASE_URL = "sqlite:///./test.db" engine = create_engine(SQLALCHEMY_DATABASE_URL, connect_args={"check_same_thread": False}) SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine) # 创建数据库表 Base.metadata.create_all(bind=engine) app = FastAPI() # 获取数据库会话 def get_db(): db = SessionLocal() try: yield db finally: db.close() @app.post("/calculate") async def create_calculation_task(base: float, exponent: float, background_tasks: BackgroundTasks, db: Session = Depends(get_db)): # 创建初始任务记录(状态为Pending) task = CalculationTask(base=base, exponent=exponent) db.add(task) db.commit() db.refresh(task) # 提交任务到线程池执行 background_tasks.add_task(calculate_and_update_task, db, task.id, base, exponent) return {"task_id": task.id, "status": task.status} @app.get("/task/{task_id}") async def get_task_status(task_id: int, db: Session = Depends(get_db)): task = db.query(CalculationTask).filter(CalculationTask.id == task_id).first() if not task: return {"error": "Task not found"} return { "task_id": task.id, "base": task.base, "exponent": task.exponent, "status": task.status, "result": task.result, "updated_at": task.updated_at }
关键注意事项
- 避免阻塞事件循环:所有耗时操作(sleep、CPU计算)必须放到线程/进程池,不能直接在async函数中调用阻塞方法。如果是CPU密集型任务,建议用
ProcessPoolExecutor替代线程池。 - 数据库会话安全:每个后台任务要确保使用独立的数据库会话,或者在任务中重新获取会话(上述示例中依赖FastAPI的会话管理,长期运行的任务建议在内部创建会话)。
- 状态实时刷新:任务执行过程中每次更新状态都要提交数据库事务,确保GET接口能立即查询到最新状态。
进阶方案(分布式场景)
如果需要支持多实例部署或更复杂的任务调度,可以用Celery+Redis/RabbitMQ替代线程池:
- Celery负责任务的分发、执行和状态管理
- 数据库依然存储最终的任务结果与状态
- FastAPI只需要提交任务到Celery,无需管理线程池
内容的提问来源于stack exchange,提问作者Weronika Materkowska
相关产品推荐
相关产品推荐

