
本文详解如何使用 asyncio.all_tasks() 实时获取并分析当前事件循环中所有未完成的异步任务,帮助定位 fastapi 应用性能瓶颈、任务积压或长时阻塞问题。
本文详解如何使用 asyncio.all_tasks() 实时获取并分析当前事件循环中所有未完成的异步任务,帮助定位 fastapi 应用性能瓶颈、任务积压或长时阻塞问题。
在构建高并发 FastAPI 服务时,理解异步任务的生命周期至关重要。你提到的“每次 await 都会调度一个任务、形成 backlog”的直觉有一定合理性,但需澄清:await 本身并不直接创建新任务——它只是挂起当前协程,将控制权交还事件循环;真正生成可被调度的独立执行单元(即 asyncio.Task)的是显式调用 asyncio.create_task()、asyncio.ensure_future(),或由框架(如 FastAPI 的 BackgroundTasks、Depends 异步依赖)自动封装的协程。当事件循环因 I/O 阻塞、CPU 密集型操作未让出、或任务调度失衡而变慢时,未完成的 Task 数量确实可能持续增长,成为可观测的性能劣化信号。
✅ 获取所有待处理与运行中任务
Python 标准库提供了权威且轻量的接口:asyncio.all_tasks(loop=None)。它返回一个 set,包含当前运行事件循环中所有尚未完成(done() is False)的 Task 对象,涵盖三类状态:
- ✅ 已调度但尚未开始执行(pending)
- ✅ 正在运行(running),包括当前正在执行 all_tasks() 的协程自身(即当前 Task 也会被包含)
- ✅ 已暂停(suspended),例如在 await asyncio.sleep(1) 或 await httpx.AsyncClient().get(...) 期间
import asyncio
import logging
# 在任意 async 函数内调用(如 FastAPI 路由、中间件或健康检查端点)
async def debug_pending_tasks():
tasks = asyncio.all_tasks() # 自动获取当前运行的 event loop
logging.info(f"当前待处理/运行中任务总数: {len(tasks)}")
for task in tasks:
# 简洁日志:任务名、状态、协程位置
status = "pending" if not task.done() else "done"
coro_name = task.get_coro().__qualname__ if task.get_coro() else "unknown"
logging.debug(f"Task {task.get_name()} [{status}] → {coro_name}")
return {"total_pending": len(tasks), "tasks": [t.get_name() for t in tasks]}? 注意:asyncio.all_tasks() 不包含普通协程(coroutine 对象),只返回已通过 create_task() 等方式“升格”为 Task 的可调度实体。这是关键设计——只有 Task 才具备独立取消、命名、错误传播等能力。
? 深度诊断:定位卡顿与堆积根源
单纯统计数量不够,还需结合上下文分析。以下两个实用技巧可快速定位异常:
SkillSub Pro - Python 题解与代码注释双功能技能功能概述SkillSub Pro - Python 题解与代码注释双功能技能是一项面向实际任务的技能,主要用于SkillSub Pro 是一个 Python 题解生成与代码注释的 双功能合体技能 ,专为学生、算法学习者和开发者设计;✅ 一个技能,两种用途 :;核心要点📝 题解模式 :输入题目/题号,自动生成完整 Python 题解(含详细注释、解题思路、复杂度分析);💬 注释模式 :输入 Python 代码,自动添加详细中。它将相关步骤、
立即学习“Python免费学习笔记(深入)”;
1. 输出每个任务的完整调用栈(识别阻塞点)
for task in asyncio.all_tasks():
if not task.done() and not task.cancelled():
# 打印堆栈,精准定位 await 卡在何处(如某次数据库查询、第三方 API 调用)
task.print_stack(limit=5) # 仅显示最深 5 层,避免日志爆炸2. 构建实时监控端点(FastAPI 示例)
from fastapi import APIRouter, Depends
from starlette.responses import JSONResponse
router = APIRouter()
@router.get("/debug/tasks")
async def list_pending_tasks():
loop = asyncio.get_running_loop()
tasks = asyncio.all_tasks(loop)
task_info = []
for task in tasks:
info = {
"name": task.get_name(),
"state": "pending" if not task.done() else ("cancelled" if task.cancelled() else "done"),
"created_at": getattr(task, "_source_traceback", "N/A"), # 非标准属性,需配合调试器
}
# 可选:记录任务创建时间(需自定义 Task 类或使用 asyncio-monitor 等工具)
task_info.append(info)
return JSONResponse(content={
"total": len(tasks),
"pending": len([t for t in tasks if not t.done()]),
"details": task_info
})部署后访问 /debug/tasks 即可获得结构化任务快照,便于集成 Prometheus 或 Grafana 做趋势监控。
⚠️ 重要注意事项
- 不要在生产高频路由中频繁调用:all_tasks() 是轻量操作,但遍历 + print_stack() 会产生可观开销,建议仅用于诊断端点或低频健康检查。
- asyncio.run() 会创建根任务:主程序入口(如 asyncio.run(main()))本身也是一个 Task,始终包含在结果集中,勿误判为“异常堆积”。
- 区分“待处理”与“积压”:少量 pending 任务是异步常态;需关注长时间(如 >30s)未完成的任务,结合 task.get_coro() 和堆栈判断是否陷入死锁、未超时的网络请求或 CPU 密集循环。
-
FastAPI 的生命周期绑定:在 startup 事件中启动后台监控协程(如每 5 秒采样一次),比每次请求都查更高效:
@app.on_event("startup") async def start_task_monitor(): asyncio.create_task(_monitor_task_backlog())
掌握 asyncio.all_tasks(),你就拥有了透视 FastAPI 异步内核的“X 光机”。它不解决性能问题本身,但能让你在毫秒级延迟出现前,就看见那条正在缓慢膨胀的任务队列——而这,正是工程化可观测性的起点。

















