如何在FastAPI中后台运行异步周期性任务
基于FastAPI实现异步周期性后台任务的正确方案
问题背景
需要实现一组可通过API控制启停、以异步协程方式周期性运行的后台任务,尝试使用FastAPI的BackgroundTasks时遇到以下问题:
- 任务无法并发执行,只能按顺序运行
- 启动任务后服务器无法处理其他请求(如访问
/docs会超时),后续得知BackgroundTasks并非为长期运行的后台任务场景设计
需求明确:
- 仅使用异步协程,不使用线程
- 支持应用启动后通过API启停任务
- 任务需能调用同一服务器的其他API
问题根源
BackgroundTasks的设计目标是处理请求结束后的短时间收尾任务,它绑定到单个请求的生命周期,所有添加的任务会在请求返回后依次执行。如果添加无限循环的长期任务,会直接阻塞事件循环,导致服务器无法处理后续请求。
解决方案:使用asyncio.create_task
直接利用Python的asyncio.create_task将长期异步任务加入事件循环,同时维护任务状态跟踪,确保可以通过API正确启停。
完整实现代码
from fastapi import FastAPI import asyncio from enum import Enum import logging # 配置日志 logging.basicConfig(level=logging.DEBUG) logger = logging.getLogger(__name__) app = FastAPI() class BaseTask: def __init__(self, name, wait_duration: int): self.name = name self.running = False self.wait_duration = wait_duration self._task = None # 用于跟踪已启动的asyncio任务 async def _run_loop(self): while self.running: try: await self.run_task() logger.debug(f"{self.name} completed, waiting for {self.wait_duration} seconds...") await asyncio.sleep(self.wait_duration) except Exception as e: logger.error(f"{self.name} encountered error: {str(e)}", exc_info=True) # 出错后等待一段时间再继续,避免无限崩溃循环 await asyncio.sleep(self.wait_duration) async def run_task(self): raise NotImplementedError def start(self): if not self.running and self._task is None: self.running = True # 创建并启动异步任务 self._task = asyncio.create_task(self._run_loop()) logger.debug(f"{self.name} started") def stop(self): if self.running: self.running = False logger.debug(f"{self.name} stopping...") # 等待任务自然结束(因为是循环,设置running=False后会在下一轮退出) self._task = None class TaskB(BaseTask): async def run_task(self): logger.debug(f"{self.name} running...") await asyncio.sleep(2) # 示例:调用同一服务器的其他API # async with httpx.AsyncClient() as client: # response = await client.get("http://localhost:8085/health") # logger.debug(f"{self.name} called health API: {response.status_code}") class TaskC(BaseTask): async def run_task(self): logger.debug(f"{self.name} running...") await asyncio.sleep(5) period_in_seconds = 10 class TaskList(str, Enum): TASK_B = "TASK_B" TASK_C = "TASK_C" ALL = "ALL" # 初始化任务列表 tasks = { TaskList.TASK_B.value: TaskB(TaskList.TASK_B.value, period_in_seconds), TaskList.TASK_C.value: TaskC(TaskList.TASK_C.value, period_in_seconds), } @app.post("/start-tasks/{task_name}") async def start_task(task_name: TaskList): target_tasks = [] if task_name == TaskList.ALL: target_tasks = list(tasks.values()) else: target_tasks = [tasks[task_name.value]] for task in target_tasks: task.start() return {"message": f"Started tasks: {[t.name for t in target_tasks]}"} @app.post("/stop-tasks/{task_name}") async def stop_task(task_name: TaskList): target_tasks = [] if task_name == TaskList.ALL: target_tasks = list(tasks.values()) else: target_tasks = [tasks[task_name.value]] for task in target_tasks: task.stop() return {"message": f"Stopped tasks: {[t.name for t in target_tasks]}"} @app.get("/health") async def health_check(): return {"status": "healthy"} if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8085)
关键改动说明
替换BackgroundTasks为asyncio.create_task:
- 将任务的循环逻辑封装到
_run_loop方法,调用asyncio.create_task将其加入事件循环,实现真正的并发执行 - 每个任务维护自己的
_task引用,避免重复启动
- 将任务的循环逻辑封装到
任务状态管理:
- 新增
_task属性跟踪已启动的任务实例,防止重复启动同一任务 - 停止任务时仅设置
running=False,让任务循环自然退出,避免强制终止导致的资源泄漏
- 新增
异常处理:
- 在任务循环中添加异常捕获,单个任务出错不会影响其他任务运行,且出错后会自动等待并继续执行
API适配:
- 任务存储改为字典,更方便按名称查找
- 启停API逻辑简化,直接操作任务实例的start/stop方法
验证效果
- 启动服务后调用
POST /start-tasks/ALL,查看日志会看到TaskB和TaskC并发执行 - 同时访问
GET /health或/docs,服务器可以正常响应,不会出现超时 - 调用
POST /stop-tasks/ALL,任务会在当前周期结束后停止
内容的提问来源于stack exchange,提问作者Luca
相关产品推荐
相关产品推荐

