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

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

关键改动说明

  1. 替换BackgroundTasks为asyncio.create_task:

    • 将任务的循环逻辑封装到_run_loop方法,调用asyncio.create_task将其加入事件循环,实现真正的并发执行
    • 每个任务维护自己的_task引用,避免重复启动
  2. 任务状态管理:

    • 新增_task属性跟踪已启动的任务实例,防止重复启动同一任务
    • 停止任务时仅设置running=False,让任务循环自然退出,避免强制终止导致的资源泄漏
  3. 异常处理:

    • 在任务循环中添加异常捕获,单个任务出错不会影响其他任务运行,且出错后会自动等待并继续执行
  4. API适配:

    • 任务存储改为字典,更方便按名称查找
    • 启停API逻辑简化,直接操作任务实例的start/stop方法

验证效果

  • 启动服务后调用POST /start-tasks/ALL,查看日志会看到TaskB和TaskC并发执行
  • 同时访问GET /health或/docs,服务器可以正常响应,不会出现超时
  • 调用POST /stop-tasks/ALL,任务会在当前周期结束后停止

内容的提问来源于stack exchange,提问作者Luca

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 07:01:13