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

如何监控Python/FastAPI事件循环中的待处理任务数量?

监控FastAPI待处理任务数量的实用方法

FastAPI基于asyncio事件循环运行,待处理任务积压确实是性能变慢的典型信号。以下是几种直接有效的监控方式:

1. 快速实现临时监控端点

直接利用asyncio的API获取当前事件循环中的所有任务,过滤出未完成的待处理任务,暴露为API端点:

from fastapi import FastAPI
import asyncio

app = FastAPI()

@app.get("/pending-tasks")
async def get_pending_task_count():
    loop = asyncio.get_running_loop()
    all_tasks = asyncio.all_tasks(loop)
    current_task = asyncio.current_task()
    # 排除当前请求任务本身,避免统计误差
    pending_tasks = [task for task in all_tasks if task != current_task and not task.done()]
    return {"pending_task_count": len(pending_tasks)}

适用场景:单进程部署下的临时排查,快速验证待处理任务是否积压。

2. 结合Prometheus实现长期监控

如果需要持续监控并配置告警,可以用Prometheus自定义指标,定期采集待处理任务数:

from fastapi import FastAPI
from prometheus_client import Gauge, start_http_server
import asyncio

app = FastAPI()
# 定义Prometheus指标
PENDING_TASKS_GAUGE = Gauge(
    "fastapi_pending_tasks_total",
    "Total number of pending async tasks in the event loop"
)

async def task_monitor():
    loop = asyncio.get_running_loop()
    while True:
        all_tasks = asyncio.all_tasks(loop)
        current_task = asyncio.current_task()
        pending_count = len([t for t in all_tasks if t != current_task and not t.done()])
        PENDING_TASKS_GAUGE.set(pending_count)
        await asyncio.sleep(1)  # 每秒采集一次

@app.on_event("startup")
async def startup():
    # 启动Prometheus metrics服务,端口8001
    start_http_server(8001)
    # 后台启动监控任务
    asyncio.create_task(task_monitor())

适用场景:生产环境长期监控,可结合Grafana配置阈值告警(比如任务数超过100时触发告警)。

3. 用事件循环钩子精准追踪任务生命周期

通过自定义任务工厂,监听任务的创建和完成事件,实时统计待处理任务数,性能更优:

from fastapi import FastAPI
import asyncio

app = FastAPI()
pending_task_counter = 0

@app.on_event("startup")
async def setup_task_tracking():
    loop = asyncio.get_running_loop()
    original_factory = loop.get_task_factory()

    def custom_task_factory(loop, coro):
        global pending_task_counter
        # 创建任务
        task = original_factory(loop, coro) if original_factory else asyncio.Task(coro, loop=loop)
        pending_task_counter += 1
        # 任务完成时递减计数器
        def on_task_done(_):
            global pending_task_counter
            pending_task_counter -= 1
        task.add_done_callback(on_task_done)
        return task

    loop.set_task_factory(custom_task_factory)

@app.get("/pending-tasks")
async def get_pending_count():
    return {"pending_task_count": pending_task_counter}

适用场景:高并发场景下的精准监控,避免遍历所有任务带来的性能开销。


注意事项

  • 多worker部署(比如uvicorn --workers 4)时,每个worker是独立进程,需分别监控每个worker的任务数,不能直接汇总。
  • 所有方法都要排除监控自身的任务(比如当前请求任务、定时采集任务),避免统计偏差。
  • 第三方异步库(如aiohttp、motor)创建的任务也会被统计,这属于正常的事件循环负载,无需过滤。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 05:53:17