如何监控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
相关产品推荐
相关产品推荐

