asyncio 的任务会泄漏吗?怎么管理任务的生命周期?
简化版
「任务泄漏」指的是:Task 被创建出来后永远不会结束、也没人负责它**,于是它一直占着内存、连接和事件循环的调度槽位。最典型的五种形态:① while True 的 worker 协程在服务关闭时没被 cancel();② 等待一个永远不会完成的 Future(队列没人 put、锁没人释放、对端不响应又没设超时);③ 每个请求都 create_task 但任务没有结束条件;④ 忘记保存引用导致任务被 GC 掉(这是「任务消失」,比泄漏更隐蔽);⑤ 任务抛异常但异常被吞,上层以为它还在跑。检测手段很直接:定期打点 len(asyncio.all_tasks())——这个数字单调增长就是泄漏的铁证;配合给任务命名(create_task(coro, name=...))和 task.print_stack(),可以直接看出「是哪类任务在堆积、它们挂在哪一行」。根治的办法是「结构化并发」:Python 3.11+ 的 asyncio.TaskGroup 保证「退出 async with 块时,块内创建的所有任务一定都已结束」——要么正常完成、要么被取消并收尾完毕,不可能有任务逃逸到块外面;相比之下 gather 在异常时不会取消其余任务**,而裸 create_task 更是完全没有生命周期约束。优雅关闭的标准流程是:cancel() 所有剩余任务 → await asyncio.gather(*tasks, return_exceptions=True) 等它们真正收尾 → 再关闭 loop;跳过第二步就会看到 Task was destroyed but it is pending!,意味着清理逻辑没执行完。核心记忆:每个任务都必须有明确的所有者;用 TaskGroup 获得生命周期保证;用 all_tasks() 计数监控泄漏。
详细版
任务泄漏的五种形态:
| 形态 | 现象 | 修法 |
|---|---|---|
while True worker 没 cancel | 关闭时 hang / 任务数不降 | 退出时 cancel + await |
| 等待永不完成的 Future | 任务数只增不减、CPU 为 0 | 所有 await 加超时 |
| 每请求 create_task 无结束条件 | 任务数随 QPS 线性增长 | 加超时 / 用 TaskGroup |
| 忘记保存引用 | 任务「凭空消失」 | 存进 set + done_callback |
| 异常被吞后循环退出 | 功能静默失效 | done_callback 记录异常 |
import asyncio, logging, time
# ① ★统一的任务注册表(防 GC + 可观测)★
_tasks: set[asyncio.Task] = set()
def spawn(coro, *, name: str | None = None) -> asyncio.Task:
t = asyncio.create_task(coro, name=name)
_tasks.add(t) # ★强引用★
t.add_done_callback(_tasks.discard) # ★完成即移除,防止集合无限增长★
t.add_done_callback(_log_exc)
return t
def _log_exc(t: asyncio.Task) -> None:
if t.cancelled():
return
if (e := t.exception()) is not None:
logging.error("任务 %s 失败", t.get_name(), exc_info=e)
# ② ★泄漏检测:定期打点任务数★
async def leak_monitor(interval=30):
while True:
tasks = asyncio.all_tasks()
logging.info("当前任务数 %d", len(tasks))
metrics.gauge("asyncio.tasks", len(tasks)) # ★单调增长 = 泄漏★
await asyncio.sleep(interval)
# 需要细看时按名字分组统计
from collections import Counter
def task_histogram():
return Counter(t.get_name().rsplit("-", 1)[0] for t in asyncio.all_tasks())
# ③ ★TaskGroup:结构化并发的生命周期保证(3.11+)★
async def handle():
async with asyncio.TaskGroup() as tg:
tg.create_task(fetch_a(), name="fetch-a")
tg.create_task(fetch_b(), name="fetch-b")
# ★到这一行时,两个任务一定都结束了(正常完成或被取消并收尾完毕)★
# ★不可能有任务逃逸到 with 块外面★
# ④ ★优雅关闭的标准流程★
async def shutdown(timeout=10):
current = asyncio.current_task()
tasks = [t for t in asyncio.all_tasks() if t is not current]
if not tasks:
return
logging.info("取消 %d 个任务", len(tasks))
for t in tasks:
t.cancel()
done, pending = await asyncio.wait(tasks, timeout=timeout)
for t in pending: # ★超时仍未结束的★
logging.warning("任务 %s 未能在 %ds 内结束", t.get_name(), timeout)
# ⑤ ★五种泄漏的反面教材★
# (a) worker 没 cancel
async def bad_worker():
while True: # ★没有退出条件★
item = await q.get()
await handle(item)
# → 服务关闭时它还在 await q.get(),永远不结束
# (b) 等待永不完成的 Future
fut = asyncio.get_running_loop().create_future()
await fut # ★没人 set_result → 永久挂起★
# (c) 每请求 create_task 且无超时
async def on_request(req):
asyncio.create_task(slow_upstream(req)) # ★没引用、没超时、没人管★
# (d) 正确写法对照
async def good_worker(stop: asyncio.Event):
while not stop.is_set():
try:
item = await asyncio.wait_for(q.get(), timeout=1) # ★有超时★
except asyncio.TimeoutError:
continue # ★回到循环顶部检查 stop★
try:
await handle(item)
except Exception:
logging.exception("处理失败") # ★吞掉异常,循环继续★
finally:
q.task_done()
⚠️ 三个必须建立的认知:①
asyncio.all_tasks()的数量是判断泄漏最直接的指标——它返回当前事件循环中尚未完成的任务,正常服务里这个数字应该围绕某个基线波动(基线 ≈ 常驻后台任务数 + 在途请求数);如果它单调增长且不回落,几乎可以断定有泄漏。配合create_task(coro, name=...)给任务命名,再按名字前缀做直方图统计,能立刻看出是哪一类任务在堆积。② 「任务泄漏」和「任务消失」是两个相反的问题,但都源于「没有所有者」:泄漏是任务永远不结束(占资源),消失是没保存引用导致 Task 被 GC 回收(官方文档明确警告:事件循环只持有弱引用)——后者更隐蔽,表现为「这段代码好像没执行」。统一的spawn()工具同时解决两者。③TaskGroup提供的核心价值是「生命周期保证」而不只是「异常处理」:async with块退出时,块内创建的所有任务一定都已结束(正常完成,或者因为某个任务失败而被取消并完成收尾)——任务不可能逃逸到块外。这是gather给不了的:gather抛异常时其余任务仍在后台运行,而裸create_task更是完全没有约束。
完整版教学
一、Task 的生命周期与「所有者」
★ Task 的状态流转:
create_task(coro)
→ ★PENDING★(已创建,等待调度或正在运行)
├─ 协程正常返回 → ★FINISHED★(result 可取)
├─ 协程抛异常 → ★FINISHED★(exception 可取)
└─ 被 cancel() → 抛 CancelledError 进协程
├─ 协程未捕获/重新抛出 → ★CANCELLED★
└─ 协程捕获了且不重抛 → ★FINISHED★(★"取消失败"★)
查询:task.done() / task.cancelled() / task._state
★ 注意:cancel() 只是"请求取消",★不保证一定能取消★
协程可以在 except CancelledError 里不重新 raise → 取消被拒绝
★ 「所有者」这个概念(★本题的核心思想★):
每个 Task 都应该有一个明确的"负责人",负责:
① ★持有引用★(防止被 GC)
② ★等待它结束★(await / TaskGroup)
③ ★处理它的异常★(try/except 或 done_callback)
④ ★在关闭时取消它★
三种所有权模式:
┌────────────────┬──────────────────────────────────────┐
│ ★await 它★ │ 调用方就是所有者(最清晰) │
│ ★TaskGroup★ │ 块是所有者,★退出时保证全部结束★ │
│ ★注册表 + 回调★│ 模块级 set 是所有者(后台任务) │
└────────────────┴──────────────────────────────────────┘
✗ 裸 create_task 且不做任何后续处理 = ★没有所有者 = 泄漏或消失★
★ 事件循环持有什么:
★事件循环只持有 Task 的弱引用★(官方文档明确说明)
→ 正在等待 IO 的任务通过 selector 的回调间接被引用
→ 但这是★实现细节★,不能依赖
→ ★必须自己保存强引用★
★ 一个直观的判断方法:
写完 create_task(...) 后问自己三个问题:
① 谁会 await 它 / 谁持有它的引用?
② 它什么条件下会结束?(★有没有可能永远不结束?★)
③ 服务关闭时谁来取消它?
★ 任何一个答不上来 → 就是潜在的泄漏★
理解任务管理要先建立**「所有者」这个概念:每个 Task 都应该有一个明确的负责人,负责四件事——持有引用(防 GC)、等待它结束、处理它的异常、关闭时取消它。三种所有权模式:直接 await(调用方就是所有者,最清晰)、TaskGroup(块是所有者,退出时保证全部结束)、注册表 + 回调(模块级 set 是所有者,用于后台任务)。而裸 create_task 且不做任何后续处理 = 没有所有者 = 泄漏或消失**。这里要特别注意 cancel() 只是「请求取消」而不保证成功——如果协程在 except CancelledError 里捕获了却不重新 raise,取消就被拒绝了。还有个实用的自查方法:写完 create_task(...) 后问自己三个问题——谁持有它的引用?它什么条件下会结束(有没有可能永远不结束)?关闭时谁来取消它?任何一个答不上来就是潜在的泄漏。
二、五种泄漏形态
★ 形态一:while True 的 worker 没有退出机制★
async def worker():
while True:
item = await q.get() # ★永远阻塞在这里★
await handle(item)
spawn(worker())
→ 服务关闭时它还挂在 q.get() 上 → ★永远不结束★
→ 现象:优雅关闭 hang 住,或退出时一堆 "Task was destroyed but it is pending"
✓ 修法(三选一):
- 退出时显式 cancel(★最常用★)
- 用停止事件:while not stop.is_set(),配 wait_for(q.get(), timeout=1)
- 用哨兵值:往队列放 None,worker 见到就 return
★ 形态二:等待一个永远不会完成的 await★
await some_future # 没人 set_result
await lock.acquire() # 持有者崩了,没释放
await q.get() # 生产者已经死了
await sock.read() # 对端不响应,★又没设超时★
→ 现象:任务数只增不减、CPU 为 0、看起来"卡住"
✓ 修法:★所有可能永久等待的 await 都要有超时★
async with asyncio.timeout(5): # 3.11+
await risky()
# 或 await asyncio.wait_for(risky(), timeout=5)
★ 形态三:每请求创建任务但没有结束保证★
async def on_request(req):
asyncio.create_task(send_webhook(req)) # ★fire and forget★
return "ok"
→ 如果 webhook 的下游卡住(没超时)→ ★任务数随 QPS 线性增长★
→ 几小时后:内存耗尽 / 连接池打满
✓ 修法:加超时 + 限并发(Semaphore)+ 注册表管理
★ 形态四:忘记保存引用(★"消失"而不是"泄漏"★)★
asyncio.create_task(background()) # ★返回值被丢弃★
→ 事件循环只持有弱引用 → ★可能在执行中途被 GC★
→ 现象:任务"好像没执行",没有任何错误
✓ 修法:统一的 spawn() 工具(set + discard 回调)
★ 形态五:异常被吞后上层以为它还活着★
async def consumer():
while True:
item = await q.get()
await handle(item) # ★抛异常 → 整个 while 退出 → 任务结束★
spawn(consumer()) # ★没有 done_callback → 没人知道它死了★
→ 现象:队列越堆越多、消息不再被消费,★但没有任何报错★
✓ 修法:① 循环体内 try/except 兜住
② done_callback 记录异常
③ ★supervisor 模式:任务挂了自动重启★
★ 五种形态的共同点:
★都是"没有明确的所有者在关心这个任务的死活"★
→ 所以解法也是共同的:★统一的任务管理 + 超时 + 可观测★
五种泄漏形态各有典型场景。① while True 的 worker 没有退出机制——关闭时挂在 q.get() 上永远不结束,修法是显式 cancel、用停止事件配超时、或用哨兵值。② 等待一个永远不会完成的 await(Future 没人 set_result、锁的持有者崩了、生产者已死、对端不响应又没超时)——现象是任务数只增不减且 CPU 为 0,修法是所有可能永久等待的 await 都要有超时。③ 每请求 fire-and-forget 创建任务——下游卡住时任务数随 QPS 线性增长,几小时后内存耗尽。④ 忘记保存引用——这是「消失」而不是「泄漏」,现象是「任务好像没执行」且没有任何错误。⑤ 异常被吞后上层以为它还活着——消费者 while 循环因异常退出,队列越堆越多但没有任何报错。五种形态的共同点是**「没有明确的所有者在关心这个任务的死活」**,所以解法也是共同的:统一的任务管理 + 超时 + 可观测。
三、检测与定位
★ 一级指标:asyncio.all_tasks() 的数量★
async def monitor():
while True:
n = len(asyncio.all_tasks())
metrics.gauge("asyncio.tasks", n)
await asyncio.sleep(15)
怎么判读:
★围绕基线波动★ 正常(基线 = 常驻后台任务 + 在途请求)
★单调增长不回落★ ★泄漏★
★突然阶跃后不降★ 某次操作创建了一批不结束的任务
★长期贴着高位★ 下游变慢导致积压(不一定是泄漏)
★ 二级:按名字分类统计(★定位是哪类任务★)
asyncio.create_task(coro, name=f"webhook-{req_id}") # ★命名★
...
from collections import Counter
def histogram():
return Counter(t.get_name().split("-")[0] for t in asyncio.all_tasks())
# → {'webhook': 8421, 'consumer': 4, 'monitor': 1} ★一眼看出是 webhook 在堆★
★ 三级:看它们挂在哪一行
for t in asyncio.all_tasks():
if t.get_name().startswith("webhook"):
t.print_stack(limit=3)
break
# Stack for <Task pending name='webhook-123' ...>:
# File "hooks.py", line 22, in send_webhook
# resp = await session.post(url, json=data) ← ★卡在这里,没超时★
★ 线上手段:
① 信号触发 dump(★不用改代码★):
loop.add_signal_handler(signal.SIGUSR2, dump_tasks)
# kill -USR2 <pid>
② 诊断端点(内网 + 鉴权):
@app.get("/_debug/tasks")
async def tasks():
return {"count": len(asyncio.all_tasks()),
"hist": dict(histogram())}
③ py-spy dump --pid(★零侵入★,但看到的是"当前正在执行的"而非全部挂起的)
★ 内存视角的佐证:
任务泄漏往往伴随内存增长:
- 每个 Task ≈ 几 KB(协程帧 + 局部变量 + 引用的对象)
- ★真正占内存的是任务持有的对象★(请求体、响应缓冲、连接)
✓ tracemalloc 对比快照,看增长最快的分配点
★ 一个完整的排查流程:
① 监控看到 asyncio.tasks 单调增长
② 打直方图 → 发现是 "webhook-*" 在堆积
③ print_stack → 卡在 session.post(★没有超时★)
④ 修:加 timeout + Semaphore 限并发 + 注册表管理
⑤ 验证:任务数回到基线并稳定
检测泄漏有三级手段。一级是 len(asyncio.all_tasks()) 打点——判读标准是:围绕基线波动正常、单调增长不回落就是泄漏、突然阶跃后不降说明某次操作创建了一批不结束的任务、长期贴着高位则可能是下游变慢导致的积压。二级是按名字分类统计——所以创建任务时一定要命名(create_task(coro, name=f"webhook-{req_id}")),用 Counter 做直方图后能一眼看出是哪类任务在堆积。三级是 task.print_stack() 看它们具体挂在哪一行(往往就是某个没设超时的 await)。线上可以用信号触发 dump(loop.add_signal_handler(SIGUSR2, dump_tasks),不用改代码)或受鉴权保护的诊断端点。完整的排查流程是:监控发现增长 → 直方图定位任务类型 → print_stack 定位代码行 → 加超时和限并发 → 验证任务数回到基线。
四、TaskGroup:用结构化并发根治
★ 「结构化并发」的核心承诺:
★任务的生命周期不能超出创建它的代码块★
→ 就像 with 保证文件一定被关闭一样,TaskGroup 保证任务一定被结束
async with asyncio.TaskGroup() as tg: # 3.11+
tg.create_task(a(), name="a")
tg.create_task(b(), name="b")
# ★执行到这一行时:a 和 b 一定都已经结束★
# - 都成功 → 正常往下走
# - 有失败 → ★取消其余任务 → 等它们收尾 → 抛 ExceptionGroup★
# - 外部取消 → 取消所有子任务 → 等收尾 → 传播 CancelledError
★ 它同时解决了四个问题:
① ★生命周期★:不会有任务逃逸到块外(本题重点)
② ★强引用★:TaskGroup 自己持有任务,不用你操心 GC
③ ★异常传播★:ExceptionGroup 汇总,不会静默丢失
④ ★取消传播★:一个失败取消全部,不会有"孤儿任务继续跑"
★ 对比 gather:
┌──────────────┬────────────────┬──────────────────────┐
│ │ gather │ ★TaskGroup★ │
├──────────────┼────────────────┼──────────────────────┤
│ 失败时其余任务│ ★继续跑★ │ ★被取消★ │
│ 退出后有遗留?│ ★可能有★ │ ★保证没有★ │
│ 引用管理 │ 自己管 │ ★自动★ │
│ 异常 │ 只抛第一个 │ ★ExceptionGroup★ │
│ 版本 │ 一直有 │ 3.11+ │
└──────────────┴────────────────┴──────────────────────┘
★ 长期运行的后台任务怎么用 TaskGroup:
async def service(stop: asyncio.Event):
async with asyncio.TaskGroup() as tg:
tg.create_task(heartbeat(), name="hb")
tg.create_task(consume(), name="consume")
await stop.wait() # ★主协程在这里等停止信号★
raise _Shutdown # ★用异常触发取消所有子任务★
# 或者让子任务自己响应 stop 事件后退出
★ 注意:TaskGroup 要求所有子任务都结束才退出,
所以长期运行的任务必须能响应取消或停止信号
★ TaskGroup 不适合的场景:
✗ ★互相独立、一个失败不该影响其他★的后台任务
(TaskGroup 会"一损俱损")
→ 这种情况在每个任务内部 try/except 兜住,别让异常抛出去
✗ 需要"提交后立刻返回、稍后再收集"的场景
→ 用注册表模式
✗ 3.10 及以下 → 用 gather 或第三方 anyio.create_task_group
★ 3.10 及以下的替代:
from contextlib import asynccontextmanager
@asynccontextmanager
async def task_group():
tasks = []
try:
yield tasks.append # 调用方用 add(coro) 提交
finally:
ts = [asyncio.create_task(c) for c in tasks]
try:
await asyncio.gather(*ts)
except BaseException:
for t in ts: t.cancel()
await asyncio.gather(*ts, return_exceptions=True)
raise
★ 或直接用 anyio(它在 3.7+ 就提供了 task group)
TaskGroup 的核心承诺是「任务的生命周期不能超出创建它的代码块」——就像 with 保证文件一定被关闭,它保证执行到 async with 块之后那一行时,块内创建的所有任务一定都已结束(都成功、或者有失败则取消其余并抛 ExceptionGroup、或者外部取消则取消所有子任务并等收尾)。它同时解决四个问题:生命周期(不会有任务逃逸)、强引用(自己持有,不用担心 GC)、异常传播(ExceptionGroup 汇总)、取消传播(不会有孤儿任务继续跑)。相比之下 gather 失败时其余任务继续跑、退出后可能有遗留。它不适合的场景要清楚:互相独立、一个失败不该影响其他的后台任务(TaskGroup 是「一损俱损」的,这种情况应该在每个任务内部 try/except 兜住)、以及需要「提交后立刻返回」的场景(用注册表模式)。3.10 及以下可以用 anyio.create_task_group 或自己封装。
五、优雅关闭的完整流程
★ 标准四步:
① ★停止接收新任务★(关监听、停止消费队列、摘除服务发现)
② ★取消所有剩余任务★
③ ★等待它们真正收尾(带超时)★
④ ★关闭资源、关闭 loop★
async def shutdown(timeout=15):
# ② 收集并取消
current = asyncio.current_task()
tasks = [t for t in asyncio.all_tasks() if t is not current]
for t in tasks:
t.cancel()
# ③ ★等它们收尾★——这一步最容易被跳过
done, pending = await asyncio.wait(tasks, timeout=timeout)
for t in pending:
logging.warning("任务 %s 未在 %ds 内结束", t.get_name(), timeout)
# ④ 关闭资源
await session.close()
await pool.close()
★ 为什么第三步不能省:
cancel() 只是★向协程抛 CancelledError★,协程还要:
- 从当前 await 点抛出异常
- 执行 finally / async with 的 __aexit__
- 可能还要 await 一些清理操作(关连接、提交 offset、释放锁)
→ ★如果不等,loop 就关了 → 清理逻辑执行不完★
→ 现象:Task was destroyed but it is pending! + ★数据不一致★
★ asyncio.run 做了什么(★3.8+ 的行为★):
asyncio.run(main())
→ main 返回后:
① 取消所有剩余任务
② gather 它们(等收尾)
③ ★关闭异步生成器(loop.shutdown_asyncgens)★
④ ★关闭默认 executor(3.9+,loop.shutdown_default_executor)★
⑤ 关闭 loop
★ 所以简单场景下 asyncio.run 已经帮你做了大部分
★ 但它★不知道你的业务清理逻辑★(关连接池、提交 offset)→ 复杂服务要自己管
★ 完整的服务骨架:
async def main():
stop = asyncio.Event()
loop = asyncio.get_running_loop()
for sig in (signal.SIGTERM, signal.SIGINT):
loop.add_signal_handler(sig, stop.set) # ★Unix★
async with AsyncExitStack() as stack: # ★资源统一管理★
session = await stack.enter_async_context(aiohttp.ClientSession())
async with asyncio.TaskGroup() as tg:
tg.create_task(consume(stop), name="consume")
tg.create_task(heartbeat(stop), name="hb")
await stop.wait() # ★等信号★
# ★TaskGroup 退出时保证子任务都结束了★
# ★ExitStack 退出时关闭所有资源★
★ 协程侧要配合:
async def consume(stop):
while not stop.is_set():
try:
item = await asyncio.wait_for(q.get(), timeout=1) # ★可中断★
except asyncio.TimeoutError:
continue
try:
await handle(item)
except asyncio.CancelledError:
await rollback(item) # ★取消时的清理★
raise # ★必须重新抛出★
except Exception:
logging.exception("处理失败")
★ 三个易错点:
① ★cancel 后不等收尾★ → 清理没做完
② ★捕获 CancelledError 却不重新 raise★ → 任务"拒绝取消",关闭会超时
③ ★清理逻辑里的 await 又被取消★ → 用 asyncio.shield 保护关键收尾
优雅关闭的标准流程是四步:停止接收新任务 → 取消所有剩余任务 → 等它们真正收尾(带超时)→ 关闭资源和 loop。第三步最容易被跳过——因为 cancel() 只是向协程抛 CancelledError,协程还需要从当前 await 点抛出、执行 finally 和 async with 的 __aexit__、甚至可能还要 await 一些清理操作(关连接、提交 offset、释放锁);不等就关 loop 的结果是清理逻辑执行不完,现象就是 Task was destroyed but it is pending! 加上数据不一致。要知道 asyncio.run 已经做了大部分(3.8+ 会取消剩余任务、gather 等收尾、关闭异步生成器、关闭默认 executor),但它不知道你的业务清理逻辑(关连接池、提交 offset),复杂服务仍要自己管。协程侧也要配合:用 wait_for 让阻塞可中断、捕获 CancelledError 做清理后必须重新 raise(否则任务「拒绝取消」、关闭会超时),清理逻辑里的关键 await 可以用 asyncio.shield 保护。
六、实践模式
★ 模式一:统一的任务注册表(后台任务的标配)★
class TaskManager:
def __init__(self):
self._tasks: set[asyncio.Task] = set()
def spawn(self, coro, *, name=None):
t = asyncio.create_task(coro, name=name)
self._tasks.add(t)
t.add_done_callback(self._tasks.discard)
t.add_done_callback(self._on_done)
return t
def _on_done(self, t):
if not t.cancelled() and (e := t.exception()):
logging.error("任务 %s 失败", t.get_name(), exc_info=e)
async def close(self, timeout=10):
for t in list(self._tasks):
t.cancel()
if self._tasks:
await asyncio.wait(list(self._tasks), timeout=timeout)
def stats(self):
return {"count": len(self._tasks),
"names": [t.get_name() for t in self._tasks]}
★ 模式二:supervisor(挂了自动重启)★
async def supervised(factory, name, stop, backoff=1, max_backoff=60):
delay = backoff
while not stop.is_set():
try:
await factory() # 跑业务协程
delay = backoff # ★成功跑完就重置退避★
except asyncio.CancelledError:
raise
except Exception:
logging.exception("%s 崩溃,%ss 后重启", name, delay)
await asyncio.sleep(delay)
delay = min(max_backoff, delay * 2) # ★指数退避★
★ 适合:消费者、心跳、长连接这类"必须一直活着"的任务
★ 注意:要有退避,否则崩溃循环会打满 CPU
★ 模式三:把任务生命周期绑定到资源★
class Connection:
async def __aenter__(self):
self._reader_task = asyncio.create_task(self._read_loop())
return self
async def __aexit__(self, *exc):
self._reader_task.cancel()
await asyncio.gather(self._reader_task, return_exceptions=True)
★ 用 async with 保证"资源关闭时任务一定被取消"
★ 模式四:请求级任务用 TaskGroup(不逃逸到请求之外)★
async def handler(req):
async with asyncio.TaskGroup() as tg:
a = tg.create_task(fetch_a())
b = tg.create_task(fetch_b())
return combine(a.result(), b.result())
★ 客户端断开 → handler 被取消 → ★TaskGroup 自动取消 a 和 b★
★ 检查清单:
□ ★所有 create_task 都通过统一工具★(引用 + 命名 + 异常回调)
□ ★所有可能永久等待的 await 都有超时★
□ ★while True 的任务都能响应停止事件或取消★
□ 3.11+ ★优先用 TaskGroup★
□ 关闭流程:★cancel → 等收尾(带超时)→ 关资源★
□ ★捕获 CancelledError 后重新 raise★
□ 监控:★asyncio.all_tasks() 计数 + 按名字分类★
□ 定期演练:发 SIGTERM 看是否能在宽限期内干净退出
★ 一句话总结:
★"每个任务都要有所有者——要么被 await、要么在 TaskGroup 里、
要么在注册表里且有超时和异常回调。没有所有者的任务,
要么泄漏(永不结束),要么消失(被 GC),两者都很难查。"★
四种实践模式各有适用场景。任务注册表是后台任务的标配(引用 + 命名 + 异常回调 + 统一关闭 + 状态查询)。supervisor 模式用于「必须一直活着」的任务(消费者、心跳、长连接)——关键是要有指数退避,否则崩溃循环会打满 CPU。把任务生命周期绑定到资源(在 __aenter__ 里创建、__aexit__ 里取消并等待),用 async with 保证「资源关闭时任务一定被取消」。请求级任务用 TaskGroup——这样客户端断开导致 handler 被取消时,子任务会自动被取消,不会逃逸到请求之外。检查清单里最关键的四条:所有 create_task 走统一工具、所有可能永久等待的 await 都有超时、3.11+ 优先用 TaskGroup、关闭流程要「cancel → 等收尾 → 关资源」。最后建议定期演练:发一个 SIGTERM 看服务能否在宽限期内干净退出。
记忆钩子:「★任务泄漏=Task 被创建后永远不结束、也没人负责它★,五种形态:①while True 的 worker 关闭时没 cancel ②★等一个永不完成的 await★(Future 没人 set_result、队列没生产者、对端不响应★又没设超时★)③每请求 fire-and-forget 且无结束条件(★任务数随 QPS 线性增长★)④★忘记保存引用被 GC★(这是『消失』不是泄漏,因为★事件循环只持有弱引用★)⑤异常被吞后 while 循环退出、上层以为它还活着(队列越堆越多★却没有任何报错★)。共同点是★没有明确的所有者★。★检测的一级指标是 len(asyncio.all_tasks())★——围绕基线波动正常、★单调增长不回落就是泄漏★;二级是★给任务命名(create_task(coro, name=))后按名字做直方图★,一眼看出哪类在堆;三级是 ★task.print_stack() 看它挂在哪一行★(往往就是某个没超时的 await)。★根治靠结构化并发:TaskGroup(3.11+)保证『退出 async with 块时块内所有任务一定都已结束』★——不可能逃逸到块外,而且它自己持有强引用、汇总 ExceptionGroup、一个失败会取消其余;★gather 做不到★(失败时其余任务继续跑)。但 TaskGroup 是『一损俱损』的,★互相独立的后台任务要在内部 try/except 兜住★。★优雅关闭四步:停止接新 → cancel 所有任务 → 等它们真正收尾(带超时)→ 关资源★;★第三步最容易被跳过★——cancel 只是抛 CancelledError,协程还要执行 finally 和清理 await,不等就关 loop 会看到 ★Task was destroyed but it is pending★ 并伴随数据不一致。协程侧要配合:★用 wait_for 让阻塞可中断、捕获 CancelledError 清理后必须重新 raise★(否则任务『拒绝取消』、关闭超时)。一句话:★每个任务都要有所有者——要么被 await、要么在 TaskGroup 里、要么在注册表里且有超时和异常回调★。」
七、常见误区与追问
- 误区:任务跑完了就自动被回收,asyncio 不存在「泄漏」问题。 泄漏的前提恰恰是任务永远跑不完。五种典型形态:
while True的 worker 在服务关闭时没被 cancel、等待一个永远不会完成的await(Future 没人set_result、队列没有生产者、对端不响应且没设超时)、每个请求 fire-and-forget 创建任务但下游卡住(任务数随 QPS 线性增长)、异常被吞后循环退出但上层以为它还活着。这些任务会一直占着内存(协程帧 + 它持有的请求体、响应缓冲、连接)和事件循环的调度结构。最直接的证据就是len(asyncio.all_tasks())单调增长不回落——正常服务里这个数字应该围绕基线波动(基线 = 常驻后台任务 + 在途请求)。 - 误区:
asyncio.create_task()之后任务就归事件循环管了,不用保存返回值。 官方文档明确警告:事件循环只持有 Task 的弱引用。不保存强引用时,Task 可能在执行完成之前被垃圾回收,任务直接消失——没有异常、没有日志,表现为「这段代码好像没执行」。这和「泄漏」是两个相反但同源的问题:泄漏是任务永远不结束、消失是任务被提前回收,根因都是「没有所有者」。正确做法是用统一的spawn()工具:把 Task 存进模块级set(强引用)并add_done_callback(set.discard)(完成后移除,避免集合无限增长),同时再挂一个记录异常的回调。顺带一提,TaskGroup自己持有强引用,用它就不必操心这件事。 - 误区:调用
task.cancel()之后任务就结束了,可以直接关闭事件循环。cancel()只是「请求取消」——它向协程当前挂起的await点抛出CancelledError,协程还需要时间完成后续动作:异常沿协程栈传播、执行finally块和async with的__aexit__、以及可能存在的清理await(关连接、提交 offset、释放分布式锁)。如果不等就关闭 loop,这些清理逻辑就执行不完,现象是满屏的Task was destroyed but it is pending!并伴随真实的数据不一致(事务没提交、锁没释放)。正确流程是cancel()→await asyncio.wait(tasks, timeout=N)→ 对超时未结束的记录告警。还有一种更隐蔽的情况:协程捕获了CancelledError却不重新raise,此时任务会「拒绝取消」并继续运行,导致关闭流程超时。 - 误区:
gather和TaskGroup只是写法不同,异常处理略有差异。 差异的核心是生命周期保证。gather在某个任务抛异常时立即把异常抛给调用方,但不会取消其余任务——那些任务会继续在后台运行到底,逃逸到gather的调用范围之外,继续占资源、继续写数据、它们自己的异常也无人认领。TaskGroup提供的是「结构化」保证:执行到async with块之后的第一行时,块内创建的所有任务必然都已结束——要么全部正常完成,要么某个失败后其余被取消并完成收尾,异常汇总成ExceptionGroup抛出。这就像with保证文件一定被关闭一样,任务不可能逃逸出块。附带好处还有:自动持有强引用、取消传播(外部取消 handler 时子任务会跟着被取消)。所以 3.11+ 应该默认用TaskGroup,只在「部分失败可接受、要让所有任务都跑完」时才用gather(return_exceptions=True)。 - 误区:给所有
await加超时太麻烦,正常情况下不会一直等下去。 「正常情况」恰恰是最不该假设的前提。没有超时的await是任务泄漏的头号来源:对端进程被 kill 但 TCP 连接没有正常关闭(半开连接可以挂几十分钟到几小时)、生产者协程崩溃后队列再也没有新数据、持有锁的任务被取消后没释放、DNS 解析卡住、下游服务假死(TCP 连着但不返回数据)。这些场景下你的任务会永久挂起,而且 CPU 为 0、没有任何错误日志,只有all_tasks()的数字在慢慢往上爬。工程上的做法是在边界处统一加超时:网络调用用客户端库自带的timeout参数(aiohttp.ClientTimeout),队列取数据用asyncio.wait_for(q.get(), timeout=1)(顺便让 worker 能定期检查停止标志),业务逻辑用async with asyncio.timeout(N)(3.11+)。超时值要分层递减,避免上游先超时而下游还在空耗。 - 追问:
asyncio.all_tasks()的数字应该是多少才算正常? 它没有绝对的「正常值」,但有明确的判读方法。首先建立基线:基线 ≈ 常驻后台任务数(心跳、消费者、监控) + 平均在途请求数——后者可以用利特尔法则估算(QPS × 平均处理时长)。然后看趋势和形态:围绕基线小幅波动是健康的;单调增长且不回落几乎可以断定泄漏;突然阶跃后维持高位说明某次操作(比如一次批量任务)创建了一批不结束的任务;长期贴着高位但会回落通常是下游变慢导致的积压而非泄漏(这时应该看在途请求数和下游延迟)。定位时用create_task(coro, name=...)命名 + 按名字前缀做Counter直方图,能立刻看出是哪一类任务在堆积,再用task.print_stack()看它们卡在哪一行。这个指标应该和「事件循环滞后」「在途请求数」一起看,三者组合才能区分「被阻塞」「任务泄漏」「下游变慢」。 - 追问:长期运行的后台任务(消费者、心跳)应该怎么管理? 三个要素缺一不可。① 可停止:不要写裸的
while True: await q.get(),而要写成while not stop.is_set()并把阻塞调用包上超时(await asyncio.wait_for(q.get(), timeout=1)),这样每秒都会回到循环顶部检查停止标志;或者用哨兵值(往队列放None)唤醒它。② 抗崩溃:循环体内部必须try/except Exception兜住,否则一次业务异常就会让整个消费者退出,而队列会越堆越多却没有任何报错;更进一步是 supervisor 模式——外层包一个「挂了就重启」的循环,并且必须有指数退避(否则持续崩溃会打满 CPU)。③ 有所有者:放进任务注册表(带命名和异常回调)或TaskGroup,关闭时统一 cancel 并等待收尾。另外要注意TaskGroup的「一损俱损」语义对互相独立的后台任务可能太激进——那种情况应该让每个任务自己吞掉异常,不要让它抛到TaskGroup层面去连累兄弟任务。 - 追问:客户端断开连接时,handler 里创建的任务会自动被取消吗? 取决于你用什么创建的。ASGI 框架(FastAPI/Starlette)在客户端断开时会取消处理该请求的那个 Task——于是 handler 协程在当前
await点收到CancelledError。此时:用await直接等待的协程会跟着一起被取消(因为它们在同一个 Task 的调用栈上);用TaskGroup创建的子任务会被自动取消(这正是结构化并发的取消传播);但用裸asyncio.create_task()创建的任务不会——它们是独立的 Task,会继续运行到底,成为「孤儿任务」继续消耗下游资源。这也是请求级并发应该优先用TaskGroup的重要理由。如果某些工作确实需要在响应返回后继续(发通知、写审计日志),那就应该明确地把它交给一个后台任务管理器(注册表 + 超时 + 限并发),而不是随手create_task——否则在高 QPS 下这些「说好了要在后台完成」的任务会不断堆积。
八、加强记忆
任务泄漏 = Task 被创建后永远不结束、也没人负责它,五种形态:① while True 的 worker 关闭时没 cancel;② 等一个永不完成的 await(Future 没人 set_result、队列没生产者、对端不响应又没设超时);③ 每请求 fire-and-forget 且无结束条件(任务数随 QPS 线性增长);④ 忘记保存引用被 GC(这是「消失」而非泄漏,因为事件循环只持有弱引用);⑤ 异常被吞后 while 循环退出、上层以为它还活着(队列越堆越多却没有任何报错)。共同点是没有明确的所有者。检测的一级指标是 len(asyncio.all_tasks())——围绕基线波动正常、单调增长不回落就是泄漏;二级是给任务命名后按名字做直方图,一眼看出哪类在堆;三级是 task.print_stack() 看它挂在哪一行(往往就是某个没超时的 await)。根治靠结构化并发:TaskGroup(3.11+)保证「退出 async with 块时块内所有任务一定都已结束」——不可能逃逸到块外,而且它自己持有强引用、汇总 ExceptionGroup、一个失败会取消其余;gather 做不到(失败时其余任务继续跑)。但 TaskGroup 是「一损俱损」的,互相独立的后台任务要在内部 try/except 兜住。优雅关闭四步:停止接新 → cancel 所有任务 → 等它们真正收尾(带超时)→ 关资源;第三步最容易被跳过——cancel() 只是抛 CancelledError,协程还要执行 finally 和清理 await,不等就关 loop 会看到 Task was destroyed but it is pending 并伴随数据不一致。协程侧要配合:用 wait_for 让阻塞可中断、捕获 CancelledError 清理后必须重新 raise(否则任务「拒绝取消」、关闭超时)。最后记住一个实际场景:客户端断开时框架会取消 handler 的 Task,TaskGroup 里的子任务会跟着被取消,而裸 create_task 创建的会变成孤儿任务继续跑。一句话:每个任务都要有所有者——要么被 await、要么在 TaskGroup 里、要么在注册表里且有超时和异常回调。