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

如何获取FastAPI后台任务运行数?验证K8s关闭时任务全完成

问题:如何获取FastAPI中当前运行的后台任务数量?

以下是我实现的FastAPI路由,每次请求/run-tasks端点时会生成一个后台任务:

import asyncio
import time

from fastapi import APIRouter
from starlette.background import BackgroundTasks


router = APIRouter()


async def background_task(sleep_time: int, task_id: int):
    print(f"Task {task_id} started")
    await asyncio.sleep(sleep_time)  # 模拟长时间运行的任务
    print(f"Task {task_id} completed")


@router.post("/run-tasks")
async def run_tasks(background_tasks: BackgroundTasks, sleep_time: int):
    background_tasks.add_task(background_task, sleep_time, int(time.time()))
    print(len(background_tasks.tasks))
    return {"message": "NEW SLEEP query arg background task started"}

当向应用进程发送SIGTERM信号时,日志会按以下顺序输出:

  • 先出现Waiting for background tasks to complete
  • 随后输出类似Task {some_id} completed的消息
  • 最后显示Waiting for application shutdown并关闭应用

我希望在每个任务完成后(即print(f"Task {task_id} completed")之后)打印剩余运行的后台任务数量,最终能看到数字0,以此确认应用关闭前所有任务都已完成。这对我在Kubernetes Pod中调优terminationGracePeriodSeconds参数至关重要。


解决方案

Starlette的BackgroundTasks是请求级别的实例,无法直接获取全局的后台任务计数,因此需要自己维护一个线程安全的全局计数器,来追踪所有活跃的后台任务:

修改后的代码如下:

import asyncio
import time
import threading
from fastapi import APIRouter
from starlette.background import BackgroundTasks

router = APIRouter()

# 全局活跃任务计数器 + 线程锁,保证多请求/多worker环境下的计数安全
active_tasks = 0
task_counter_lock = threading.Lock()


async def background_task(sleep_time: int, task_id: int):
    global active_tasks
    print(f"Task {task_id} started")
    try:
        await asyncio.sleep(sleep_time)  # 模拟耗时任务
    finally:
        print(f"Task {task_id} completed")
        # 任务完成后原子减少计数并打印剩余数量
        with task_counter_lock:
            active_tasks -= 1
            print(f"Remaining active background tasks: {active_tasks}")


@router.post("/run-tasks")
async def run_tasks(background_tasks: BackgroundTasks, sleep_time: int):
    global active_tasks
    task_id = int(time.time())
    # 添加任务前先原子增加计数
    with task_counter_lock:
        active_tasks += 1
    background_tasks.add_task(background_task, sleep_time, task_id)
    print(f"Added task {task_id}, current active tasks: {active_tasks}")
    return {"message": "NEW SLEEP query arg background task started"}

关键说明:

  1. 线程安全计数:用threading.Lock保护计数器的增减操作,避免多线程/多请求场景下的计数混乱(FastAPI默认用多线程模式运行,多worker部署时也能保证计数安全)。
  2. 任务生命周期追踪:添加任务时计数器+1,任务完成(无论成功/失败)时计数器-1,确保计数准确。
  3. Kubernetes适配:当Pod收到SIGTERM后,FastAPI会等待所有后台任务执行完毕,最后一个任务完成时会打印Remaining active background tasks: 0,以此确认所有任务处理完成,你可以根据这个日志调整terminationGracePeriodSeconds的取值,避免Pod被强制杀死时还有未完成的任务。

内容的提问来源于stack exchange,提问作者rokpoto.com

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 23:37:02