写Python异步代码时,你是否遇到过这样的警告:Task was destroyed but it is pending?或者更糟的情况:协程里抛了异常,主流程却毫无察觉,日志里一条记录都没有,直到数据缺了一块才发现问题。这类问题的根源几乎都指向同一处:任务的生命周期没有被正确管理。事件循环只负责调度,不会替你兜底处理游离状态的协程。本文从错误成因讲起,逐步给出任务分组与异常捕获的完整解法。

事件循环错误是怎么产生的
先看一段典型的错误代码。开发者在协程函数里创建了一个后台任务,但没有保存引用,也没有等待它结束:
import asyncio
async def save_log():
await asyncio.sleep(1)
# 模拟写库失败
raise RuntimeError("数据库连接超时")
async def main():
asyncio.create_task(save_log()) # 创建后直接丢弃引用
print("主流程继续执行")
asyncio.run(main())
这段代码有两个隐患。第一,create_task返回的Task对象没有被任何变量持有,垃圾回收器随时可能回收它,导致任务执行到一半就被销毁,这就是Task was destroyed but it is pending警告的来源。第二,save_log里抛出的异常发生在主流程结束之后,事件循环已经关闭,异常信息只会打印到stderr,如果你的日志系统只采集log模块的输出,这条错误就彻底消失了。
本质上,事件循环对任务是弱管理关系。它只保证被调度的任务获得执行机会,但不保证任务完成、不保证异常传递、不保证优雅退出。这三个不保证就是异步错误的温床。要根治,必须把任务纳入一个有明确生命周期的容器中,这就是任务分组要解决的问题。
用gather与TaskGroup做任务分组
任务分组的核心思想是:把一组并发任务交给一个统一的管理者,由它负责等待全部完成、收集结果、汇总异常。asyncio.gather是最经典的方式:
import asyncio
async def fetch(url: str, delay: float):
await asyncio.sleep(delay)
if "bad" in url:
raise ValueError(f"无效地址: {url}")
return f"{url} 的响应内容"
async def main():
tasks = [
fetch("https://ipipp.com/api/a", 1),
fetch("https://ipipp.com/bad-url", 2),
fetch("https://ipipp.com/api/c", 3),
]
results = await asyncio.gather(*tasks, return_exceptions=True)
for r in results:
if isinstance(r, Exception):
print(f"捕获到异常: {r!r}")
else:
print(f"正常结果: {r}")
asyncio.run(main())
return_exceptions=True是关键参数。不加它时,gather会在第一个异常出现时立即抛出,其余任务继续在后台运行但结果被丢弃;加上之后,异常对象会作为结果返回,你可以逐个检查。这种方式适合批量任务彼此独立、允许部分失败的场景,比如批量抓取网页。
Python 3.11之后更推荐使用TaskGroup。它的语义更严格:组内任何一个任务失败,其余任务会被自动取消,最后以ExceptionGroup的形式把所有异常打包抛出:
import asyncio
async def fetch(name: str, delay: float):
await asyncio.sleep(delay)
if name == "bad":
raise ValueError(f"任务 {name} 失败")
return name
async def main():
try:
async with asyncio.TaskGroup() as tg:
tg.create_task(fetch("a", 1), name="任务a")
tg.create_task(fetch("bad", 2), name="任务bad")
tg.create_task(fetch("c", 3), name="任务c")
except* ValueError as eg:
for exc in eg.exceptions:
print(f"处理值错误: {exc}")
except* TimeoutError as eg:
print("存在超时任务")
asyncio.run(main())
两种方式的取舍很明确:任务之间无依赖、允许部分失败,用gather加return_exceptions;任务是整体操作、一个失败全部作废(比如同时写主库和缓存),用TaskGroup,它保证退出async with块时组内没有任何遗留任务,从机制上杜绝了游离任务。注意except*是3.11引入的语法,专门用于拆解ExceptionGroup。
结构化异常捕获与资源清理
分组只解决了一半问题,另一半是异常发生后如何清理资源。协程被取消时会抛出CancelledError,很多人把它当普通异常捕获后吞掉,导致任务拒绝取消、程序卡死。正确做法是在finally块中做清理,并让取消异常继续传播:
import asyncio
async def worker(conn_name: str):
conn = None
try:
conn = await acquire_connection(conn_name)
await asyncio.sleep(10) # 模拟长耗时操作
except asyncio.CancelledError:
print(f"{conn_name} 被取消,记录中断状态")
raise # 必须重新抛出,否则任务无法真正取消
finally:
if conn:
await conn.close() # 确保连接释放
print(f"{conn_name} 连接已关闭")
async def acquire_connection(name):
await asyncio.sleep(0.1)
return type("Conn", (), {"close": lambda self: asyncio.sleep(0)})()
async def main():
task = asyncio.create_task(worker("订单同步"))
await asyncio.sleep(1)
task.cancel()
try:
await task
except asyncio.CancelledError:
print("主流程确认任务已终止")
asyncio.run(main())
除了手动取消,超时也是常见的异常来源。给任务加一层wait_for或timeout保护,可以避免单个慢任务拖垮整组并发:
import asyncio
async def slow_api():
await asyncio.sleep(30)
return "结果"
async def main():
try:
result = await asyncio.wait_for(slow_api(), timeout=3)
except asyncio.TimeoutError:
print("接口超时,走降级逻辑")
result = "默认兜底数据"
print(f"最终使用: {result}")
asyncio.run(main())
最后梳理几条实践原则:一是永远保存create_task的返回值,最简单的方式是立刻await它,或者放进任务集合;二是对外层入口统一包一层try except,把ExceptionGroup拆开逐条记录日志,而不是只打印repr;三是清理逻辑写在finally里,CancelledError捕获后必须raise回去;四是批量操作优先选TaskGroup,让结构化并发替你管理生命周期。把这几点落实到位,事件循环层面的报错基本可以从你的异步程序里绝迹。