Python 基础体系 · 第 51/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python asyncio 完整基础:事件循环、协程、Future 与调度

asyncio 是 Python 标准库中基于 async/await 语法实现异步并发的基础设施,主要面向 I/O 密集型任务和高层网络程序。它的核心不是“自动创建更多线程”,而是在一个事件循环线程中,让多个可暂停的任务交替运行。(docs.python.org)

要真正理解 asyncio,需要把几个概念分开:

  • 协程函数:用 async def 定义的函数。
  • 协程对象:调用协程函数后得到的对象,代表一段尚未执行或可暂停恢复的逻辑。
  • await:等待一个可等待对象,并在必要时把控制权交回事件循环。
  • 事件循环:负责取出就绪工作、等待 I/O 或定时器,并恢复相应任务。
  • Task:被事件循环调度的协程。
  • Future:表示一个尚未完成的结果或状态,通常用于连接底层回调式 API 与高级异步代码。

一、先区分并发、并行与异步

1. 并发不等于并行

设有两个任务:

任务 A:等待网络响应 2 秒
任务 B:等待网络响应 3 秒

串行执行时,总耗时近似为:

Tserial=TA+TB=2+3=5 秒T_{\text{serial}} = T_A + T_B = 2 + 3 = 5\text{ 秒}

如果两个任务可以同时发起网络请求,那么等待期间不需要占用 Python 代码的执行时间,总耗时接近:

Tconcurrentmax(TA,TB)=3 秒T_{\text{concurrent}} \approx \max(T_A, T_B) = 3\text{ 秒}

这叫并发:任务的生命周期发生重叠。

并行则要求多个任务在同一时刻由不同的执行单元运行,例如多个线程、多个进程或多个解释器。asyncio 的默认模型主要是事件驱动的协作式并发,而不是依靠多个 CPU 核心并行执行 Python 字节码。Python 官方文档也将事件驱动的协作式多任务与线程、进程等抢占式并发工具区分开来。(docs.python.org)

因此:

asyncio

适合:

  • 网络请求;
  • 套接字读写;
  • 数据库异步客户端;
  • 大量等待外部资源的任务;
  • 需要同时管理许多连接的服务器。

它不适合直接承担长时间的 CPU 密集计算。一个正在事件循环线程中运行的 CPU 任务,会阻塞同一线程中的其他任务和 I/O。(docs.python.org)


二、协程函数与协程对象

1. async def 定义的是协程函数

async def greet():
    print("hello")

greet 本身是一个协程函数。调用它:

coro = greet()

得到的是协程对象,而不是立即执行结果。

print(coro)
# <coroutine object greet at 0x...>

调用普通函数时,函数体通常会立即执行:

def normal_greet():
    print("hello")

normal_greet()
# hello

调用协程函数时,函数体不会因此立即执行:

async def async_greet():
    print("hello")

async_greet()
# 只创建协程对象,不打印 hello

Python 官方文档明确指出,单纯调用协程函数不会把它加入调度;需要通过 awaitasyncio.run()asyncio.create_task() 等方式运行。(docs.python.org)

2. 未执行的协程对象必须被处理

下面的代码会产生警告:

import asyncio

async def work():
    print("work")

async def main():
    work()  # 创建协程对象,但没有 await,也没有创建 Task

asyncio.run(main())

典型警告:

RuntimeWarning: coroutine 'work' was never awaited

正确写法有两种。

如果希望当前流程等待它完成:

async def main():
    await work()

如果希望它作为独立任务被调度:

async def main():
    task = asyncio.create_task(work())
    await task

这里的 task 不是普通的协程对象,而是由事件循环管理的 Task。后文会详细解释两者的区别。


三、await 到底做了什么

1. 可等待对象

能出现在 await expression 中的对象称为可等待对象,英文是 awaitableasyncio 中最主要的三类可等待对象是:

  1. 协程对象;
  2. asyncio.Task
  3. asyncio.Future。(docs.python.org)

例如:

await some_coroutine()
await some_task
await some_future

但是,“能被 await”并不自动意味着“会把控制权交给事件循环”。

2. await coroutine() 不一定切换任务

考虑下面的代码:

import asyncio

async def child():
    print("child")

async def main():
    print("before")
    await child()
    print("after")

asyncio.run(main())

执行 await child() 时,child() 只是被当前协程直接调用并运行。由于 child() 内部没有等待一个会暂停它的对象,当前任务不会在这里切换给其他任务。

可以把它近似理解为:

result = child()  # 不是普通函数调用的语法表现,但控制流效果类似直接进入 child

更准确地说,协程对象会通过其协程协议运行;如果内部没有真正挂起的 await,控制权就会沿着当前协程调用链继续向下执行。

这意味着下面代码中的 worker_b 可能要等 worker_a 的三个循环全部结束后才开始:

import asyncio

async def worker_a():
    for _ in range(3):
        print("A")

async def worker_b():
    print("B")

async def main():
    task_b = asyncio.create_task(worker_b())

    for _ in range(3):
        await worker_a()

    await task_b

asyncio.run(main())

输出通常是:

A
A
A
B

如果改成:

async def main():
    task_b = asyncio.create_task(worker_b())

    for _ in range(3):
        await asyncio.create_task(worker_a())

    await task_b

那么每次 worker_a() 都被包装成一个 Task。当前任务等待 Task 时会挂起,事件循环可以调度其他就绪任务。

但这并不表示应该把所有协程都无条件包装成 Task。Task 会改变调度语义,也会增加任务管理和生命周期责任。当前协程只是需要顺序调用另一个异步函数时,直接使用 await child() 通常才是正确表达。

3. await 的协议层含义

从协议角度看:

await obj

会尝试使用对象的 __await__() 方法。该方法返回一个迭代器,协程在执行过程中可能通过 yield 暂停,把控制权交给上层驱动者;事件循环再根据底层对象完成情况恢复它。

可以构造一个极简的可等待对象:

class YieldOnce:
    def __await__(self):
        yield
        return "done"

使用:

import asyncio

async def main():
    result = await YieldOnce()
    print(result)

asyncio.run(main())

这里的 yield 不是业务结果,而是一个“暂停点”。真正的 asyncio I/O 对象会把这个暂停点与套接字、定时器或底层 Future 联系起来。Python 官方的概念性说明也把 await 描述为调用 __await__() 并传播其中的暂停信号。(docs.python.org)


四、事件循环:asyncio 的调度中心

1. 事件循环是什么

事件循环可以抽象成一个不断重复的驱动器:

取出就绪工作
    ↓
运行一个或多个回调/任务
    ↓
任务完成、异常或执行到 await
    ↓
等待 I/O、定时器或跨线程通知
    ↓
把新产生的工作放入就绪队列
    ↓
重复

事件循环中的任务通常在同一个线程内逐个运行。一个任务执行时,同一线程中的其他 Task 不会同时执行;当当前 Task 执行到能够暂停的 await 时,事件循环才有机会调度其他工作。(docs.python.org)

2. 就绪队列、定时器和 I/O 等待

以 CPython 的常见实现为例,事件循环至少需要管理几类工作:

  • 就绪队列:已经可以立即执行的回调;
  • 定时器队列:未来某个时间点才能执行的回调;
  • I/O 监听器:等待套接字或其他文件描述符变为可读、可写;
  • 任务恢复回调:某个 Future 完成后,恢复等待它的 Task。

可以用下面的伪代码描述一次循环迭代:

def event_loop_iteration():
    due_callbacks = move_due_timers_to_ready()

    timeout = calculate_selector_timeout(
        ready_queue=ready,
        scheduled_timers=timers,
    )

    io_events = selector.select(timeout)

    for event in io_events:
        schedule_callbacks_for(event)

    current_batch = take_current_ready_callbacks()

    for callback in current_batch:
        callback()

这段伪代码表达的是核心数据流,但不是 asyncio 对外承诺的精确 API 行为。就绪队列的内部数据结构、一次迭代处理多少回调以及不同平台的 I/O 细节属于实现层面;应用程序不应依赖这些内部细节。

在 Unix 上,SelectorEventLoop 基于 selectors 模块选择平台可用的高效 I/O 选择器;Windows 上还存在基于 IOCP 的 ProactorEventLoop。Python 3.13 起,asyncio.EventLoop 是当前平台上高效事件循环实现的别名。(docs.python.org)

3. I/O 多路复用解决了什么问题

假设程序管理 10,000 个 TCP 连接。传统阻塞模型可能为每个连接分配一个线程,并让线程阻塞在:

data = sock.recv(...)

事件循环模型则把多个套接字交给操作系统的 I/O 多路复用机制:

连接 1:等待可读
连接 2:等待可读
连接 3:等待可写
...
连接 N:等待可读

程序不必为每个连接都占用一个正在阻塞的 Python 执行栈。事件循环向操作系统询问:

“哪些描述符现在已经准备好?”

操作系统返回就绪事件后,事件循环只调度相关回调或任务。

因此,异步 I/O 的关键不是“让一次网络读取变快”,而是避免在等待网络期间浪费执行线程,使同一事件循环能够管理更多同时等待的操作。


五、协作式调度:任务必须主动让出控制权

1. 协作式调度的定义

asyncio 默认采用协作式调度

当前任务运行到完成,或者主动执行一个可暂停的 await,事件循环才有机会运行其他任务。

这与线程的抢占式调度不同。线程可能在操作系统调度器决定的时间片边界被切换;而普通的 asyncio Task 不会在任意 Python 字节码指令之间被事件循环强制切换。

下面的函数会阻塞整个事件循环:

import asyncio
import time

async def blocking_worker():
    print("blocking start")
    time.sleep(2)  # 阻塞当前线程
    print("blocking end")

async def ticker():
    for i in range(3):
        print("tick", i)
        await asyncio.sleep(0.5)

async def main():
    await asyncio.gather(
        blocking_worker(),
        ticker(),
    )

asyncio.run(main())

预期表现类似:

blocking start
# 约 2 秒内没有 tick
blocking end
tick 0
tick 1
tick 2

asyncio.sleep() 是异步等待;time.sleep() 是线程阻塞。两者名字相似,但调度语义完全不同。

2. 事件循环延迟的因果关系

设事件循环线程正在执行一个阻塞调用,阻塞时间为 BB。在这段时间内:

  1. 已经就绪的其他 Task 不能运行;
  2. 已到期的定时器不能及时执行;
  3. 已经准备好的 I/O 回调不能及时处理;
  4. 等待这些 Task 的上层 Task 也无法恢复。

因此,一个看似局部的阻塞调用会把延迟传播到整个事件循环。

官方文档明确说明:如果一个 CPU 密集型函数直接运行 1 秒,同一事件循环中的其他 asyncio Task 和 I/O 操作都会至少延迟 1 秒。(docs.python.org)

3. 用线程执行阻塞函数

对于阻塞但不适合改写为异步 API 的函数,可以使用:

import asyncio
import time

def blocking_io():
    time.sleep(2)
    return "finished"

async def main():
    result = await asyncio.to_thread(blocking_io)
    print(result)

asyncio.run(main())

asyncio.to_thread() 的核心效果是:把同步函数放到线程中执行,而不是直接占用事件循环线程。

对于 CPU 密集型 Python 代码,线程通常不能自动带来真正的 Python 字节码并行;此时应根据任务性质考虑进程、解释器池、扩展模块或其他并行方案。底层事件循环也提供 run_in_executor(),可配合线程池、解释器池或进程池执行阻塞代码。(docs.python.org)


六、Task:被事件循环管理的协程

1. Task 的定义

Task 可以理解为:

协程对象 + 事件循环调度状态 + 完成回调 + 结果/异常

创建 Task 会自动把协程安排到事件循环中:

task = asyncio.create_task(coro())

但要注意:

  • coro() 只是创建协程对象;
  • create_task(coro()) 才是创建并安排 Task;
  • Task 被安排后,也不会抢占当前正在运行的 Task;
  • 它要等当前 Task 让出控制权,才有机会开始运行。

2. 一个完整的并发例子

import asyncio
import time

async def say_after(delay: float, message: str) -> str:
    await asyncio.sleep(delay)
    print(message)
    return message

async def main():
    started = time.perf_counter()

    task1 = asyncio.create_task(say_after(1, "hello"))
    task2 = asyncio.create_task(say_after(2, "world"))

    result1 = await task1
    result2 = await task2

    elapsed = time.perf_counter() - started

    print(result1, result2)
    print(f"elapsed: {elapsed:.1f}s")

asyncio.run(main())

预期输出:

hello
world
hello world
elapsed: 2.0s

如果改成顺序等待:

async def main():
    await say_after(1, "hello")
    await say_after(2, "world")

总耗时约为 3 秒。并发版本约为 2 秒,因为两个等待区间重叠。

3. Task 的生命周期

一个 Task 的逻辑状态可以表示为:

PENDING
  ├── 协程正在运行
  ├── 协程在 await 处暂停
  ├── 等待 I/O、定时器或 Future
  │
  ├── DONE:协程正常返回
  ├── DONE:协程抛出异常
  └── CANCELLED:收到取消请求并结束

Task 完成后,其结果由协程的 return 值提供:

task = asyncio.create_task(say_after(1, "hello"))
value = await task

如果协程抛出异常:

async def fail():
    raise ValueError("bad input")

async def main():
    task = asyncio.create_task(fail())

    try:
        await task
    except ValueError as exc:
        print(f"caught: {exc}")

asyncio.run(main())

异常会在等待 Task 的位置重新抛出。若 Task 的异常从未被消费,程序可能产生:

Task exception was never retrieved

调试模式可以提供更完整的 Task 创建位置和异步调用信息。(docs.python.org)

4. 不要丢失后台 Task 的引用

下面的写法有生命周期风险:

async def main():
    asyncio.create_task(background_work())
    await asyncio.sleep(1)

如果这个后台任务的生命周期需要被明确管理,就应保留引用:

async def main():
    task = asyncio.create_task(background_work())

    try:
        await asyncio.sleep(1)
    finally:
        task.cancel()
        try:
            await task
        except asyncio.CancelledError:
            pass

更现代的方式是使用 TaskGroup,让一组子任务拥有明确的父级生命周期。TaskGroup 属于结构化并发主题,将在关联文章中展开;这里需要知道的是,create_task() 负责创建独立 Task,而 TaskGroup 负责把多个 Task 组织到一个受控作用域中。Python 文档将 TaskGroup 作为 create_task() 的更现代替代方案之一。(docs.python.org)


七、Future:一个尚未完成的结果容器

1. Future 不代表具体计算

Future 表示一个未来可能产生结果、异常或取消状态的对象。它本身通常不负责执行具体业务逻辑。

可以把它看成一个状态容器:

PENDING
  ├── set_result(value)     → DONE
  ├── set_exception(error)   → DONE
  └── cancel()              → CANCELLED

与协程不同:

  • 协程表示“如何执行一段逻辑”;
  • Future 表示“某个结果现在是否已经准备好”。

TaskFuture 的子类,因此 Task 既有协程执行能力,也有 Future 的结果和完成回调能力。(docs.python.org)

2. 手动完成 Future

下面的例子中,一个协程创建 Future,另一个协程在一秒后为它设置结果:

import asyncio

async def set_result_later(
    future: asyncio.Future[str],
    delay: float,
    value: str,
) -> None:
    await asyncio.sleep(delay)

    if not future.cancelled():
        future.set_result(value)

async def main():
    loop = asyncio.get_running_loop()
    future = loop.create_future()

    asyncio.create_task(
        set_result_later(future, 1, "done")
    )

    print("waiting...")
    result = await future
    print(result)

asyncio.run(main())

执行过程是:

  1. future 初始为 PENDING
  2. set_result_later() 被创建为 Task;
  3. main() 等待 future,暂时挂起;
  4. 一秒后,辅助 Task 调用 future.set_result("done")
  5. Future 状态变成 DONE
  6. 等待 Future 的 main() 被重新安排到事件循环;
  7. await future 得到字符串 "done"

如果没有任何代码调用 set_result()set_exception()cancel(),这个 Future 将永远处于等待状态。

3. Future 的完成回调

def on_done(future: asyncio.Future[str]) -> None:
    print("future completed")

async def main():
    loop = asyncio.get_running_loop()
    future = loop.create_future()
    future.add_done_callback(on_done)

    future.set_result("ok")
    print(await future)

asyncio.run(main())

完成回调不会作为同步函数直接执行,而是通过事件循环安排运行。Future 文档明确规定,Future 完成回调会使用 loop.call_soon() 调度。(docs.python.org)

4. 为什么应用代码通常少直接创建 Future

Future 是底层桥接工具,适合以下场景:

外部回调 API
    ↓
future.set_result() / set_exception()
    ↓
asyncio Task await future

例如,某个旧式库只有回调接口:

library.start_operation(
    callback=on_success,
    error_callback=on_error,
)

可以把它包装成 Future,让上层使用:

result = await async_wrapper()

但业务层如果只是要并发运行协程,通常应使用 create_task()gather()TaskGroup,而不是手动创建 Future。


八、协程、Task 与 Future 的关系

可以用下面的关系图概括:

flowchart TD
    A["async def 函数"] --> B["调用后得到协程对象"]
    B --> C{"如何运行?"}
    C -->|"await"| D["当前协程直接等待"]
    C -->|"create_task"| E["Task"]
    E --> F["事件循环调度"]
    E --> G["Task 继承 Future 能力"]
    H["底层回调或 I/O"] --> I["Future"]
    I --> J["set_result / set_exception / cancel"]
    J --> F
    D --> K["得到返回值或抛出异常"]
    F --> K

最容易混淆的是下面三种写法:

coro = work()

只创建协程对象。

result = await work()

当前协程直接等待 work() 完成;是否产生任务切换,要看 work() 内部是否真正挂起。

task = asyncio.create_task(work())
result = await task

创建一个由事件循环调度的 Task,再等待 Task 完成。


九、从 await 到事件循环的完整调度路径

考虑:

async def parent():
    child_task = asyncio.create_task(child())
    result = await child_task
    return result

其调度过程可以拆成以下步骤。

第一步:创建 Task

create_task() 创建 Task,并把“启动 Task 协程”的回调安排到事件循环的就绪队列。

此时:

parent:正在运行
child_task:PENDING,等待事件循环调度

第二步:parent 等待 Task

执行:

await child_task

当前 Task 暂停,并在 child_task 上登记一个完成回调:

child_task 完成后:
    把 parent 的恢复逻辑加入事件循环就绪队列

此时:

parent:等待 child_task
child_task:PENDING

第三步:事件循环调度 child

事件循环从就绪队列中取出 child 的启动回调,运行 child()

如果 child() 执行:

await asyncio.sleep(1)

那么 child 会暂停,定时器负责在约一秒后把恢复 child 的回调重新放入就绪队列。

第四步:child 完成

如果 child 返回:

return 42

那么:

child_task:PENDING → DONE
child_task.result():42

它的完成回调被安排到事件循环中。

第五步:恢复 parent

事件循环再次调度 parent 的恢复回调:

result = await child_task

得到 42,然后 parent 继续执行。

这个过程说明,await 并不是“阻塞线程直到结果出现”,而是建立“结果完成后恢复当前协程”的关系。线程在等待期间可以回到事件循环,处理其他就绪工作。


十、阻塞调用、异步 I/O 与同步文件操作

1. 不是所有 I/O 都自动异步

下面这些调用虽然属于 I/O,但如果 API 本身是同步的,仍然会阻塞事件循环:

data = open("large-file.txt").read()
response = requests.get(url)
time.sleep(1)

“这是 I/O”与“这是非阻塞异步 I/O”是两件不同的事。

异步代码需要使用:

  • 原生 asyncio 网络 API;
  • 支持 async/await 的第三方客户端;
  • asyncio.to_thread()
  • run_in_executor()
  • 或其他明确的异步适配层。

2. 计算密集型循环同样会阻塞

async def cpu_bound():
    total = 0
    for i in range(100_000_000):
        total += i
    return total

即使函数定义为 async def,只要循环内部没有可暂停的 await,它仍然会连续占用事件循环线程。

async def 不会把同步计算自动变成后台线程,也不会自动分配 CPU 核心。

3. 让出控制权不能替代拆分计算

有时可以在长循环中加入:

await asyncio.sleep(0)

这会给事件循环一次调度机会:

async def cooperative_loop():
    for i in range(1_000_000):
        do_small_piece(i)

        if i % 1000 == 0:
            await asyncio.sleep(0)

但这只改善响应性,不会让总计算量消失。频繁切换还会引入调度开销。对于真正的 CPU 密集任务,应把工作放到合适的并行执行单元,而不是仅依赖 sleep(0)


十一、事件循环的生命周期

1. 推荐使用 asyncio.run()

最常见的程序入口是:

import asyncio

async def main():
    await asyncio.sleep(0.1)
    print("done")

if __name__ == "__main__":
    asyncio.run(main())

asyncio.run() 会:

  1. 创建事件循环;
  2. 执行传入的可等待对象;
  3. 等待顶层协程完成;
  4. 结束异步生成器;
  5. 关闭默认执行器;
  6. 关闭事件循环。

它不能在同一线程中已有事件循环运行时再次调用,通常应作为应用程序的顶层入口,并且一个程序通常只调用一次。Python 3.14 中,asyncio.run() 的参数可以是任意可等待对象,不再限于协程对象。(docs.python.org)

下面的代码在已有事件循环的环境中可能失败:

async def handler():
    asyncio.run(main())  # 如果当前线程已有事件循环,会报错

在异步函数内部,应直接:

await main()

2. asyncio.Runner 适合多次运行顶层异步调用

如果需要在同一个事件循环和上下文中多次运行顶层异步函数,可以使用 Runner

import asyncio

async def first():
    return "first"

async def second():
    return "second"

with asyncio.Runner() as runner:
    print(runner.run(first()))
    print(runner.run(second()))

退出 with 代码块时,Runner 会关闭异步生成器、关闭默认执行器、关闭事件循环并释放相关上下文。(docs.python.org)

3. 手动管理事件循环

底层代码也可以手动创建和关闭事件循环:

import asyncio

async def main():
    print("hello")

loop = asyncio.new_event_loop()

try:
    loop.run_until_complete(main())
finally:
    loop.run_until_complete(loop.shutdown_asyncgens())
    loop.close()

事件循环关闭时:

  • 事件循环不能仍在运行;
  • 未处理的待执行回调会被丢弃;
  • 关闭操作不可逆;
  • 关闭后不能继续调用其他事件循环方法。

asyncio.run() 已经负责这些收尾工作,因此普通应用不应重复手动关闭由它创建的事件循环。(docs.python.org)

4. 关闭不是取消所有业务任务的替代品

事件循环关闭只代表调度基础设施停止,不等于业务资源已经正确释放。应用程序仍然需要处理:

  • 未完成的 Task;
  • 网络连接;
  • 异步生成器;
  • 执行器中的线程;
  • 数据库连接;
  • 服务器 socket;
  • 外部子进程。

一个合格的关闭流程通常是:

停止接收新请求
    ↓
通知后台任务退出
    ↓
等待任务完成或取消
    ↓
关闭连接和异步生成器
    ↓
关闭执行器
    ↓
关闭事件循环

如果任务没有退出点,或者阻塞调用仍在运行,关闭流程就可能延迟甚至无法按预期完成。


十二、取消是协作式的生命周期控制

取消一个 Task:

task.cancel()

并不是从另一个线程强行杀死它,而是请求在 Task 的协程中注入 asyncio.CancelledError。如果协程正在可取消的 await 上等待,它通常会在那里收到取消异常:

import asyncio

async def worker():
    try:
        while True:
            print("working")
            await asyncio.sleep(1)
    except asyncio.CancelledError:
        print("cleanup")
        raise

async def main():
    task = asyncio.create_task(worker())

    await asyncio.sleep(0.1)
    task.cancel()

    try:
        await task
    except asyncio.CancelledError:
        print("cancelled")

asyncio.run(main())

输出类似:

working
cleanup
cancelled

finally 适合释放资源:

async def worker():
    resource = await acquire()

    try:
        await use(resource)
    finally:
        await resource.close()

如果捕获 CancelledError 后不重新抛出,就等于吞掉取消请求,可能导致上层认为任务已经正常完成。是否吞掉取消必须有明确的控制流理由,不能作为普通异常随意处理。


十三、诊断常见失败表现

1. 协程从未被等待

失败代码:

async def main():
    work()

诊断信号:

RuntimeWarning: coroutine 'work' was never awaited

修复:

await work()

或者:

task = asyncio.create_task(work())
await task

2. Task 异常无人读取

失败代码:

async def main():
    asyncio.create_task(fail())
    await asyncio.sleep(1)

如果 fail() 抛出异常,而程序从未等待它或读取其结果,就可能出现:

Task exception was never retrieved

修复方式是明确管理 Task:

task = asyncio.create_task(fail())

try:
    await task
except Exception as exc:
    print(f"task failed: {exc}")

或者把任务放入 TaskGroup,让异常传播和兄弟任务取消具有结构化边界。

3. 事件循环响应变慢

典型原因包括:

time.sleep(...)
requests.get(...)
subprocess.run(...)
大量 CPU 循环
同步日志处理器
同步文件读写

诊断时可启用调试模式:

asyncio.run(main(), debug=True)

调试模式可以帮助发现未等待协程、未消费异常以及长时间占用事件循环的代码。官方概念说明提到,长时间垄断执行的协程会被调试模式记录。(docs.python.org)


十四、跨线程使用事件循环

事件循环通常运行在一个线程中,而大多数 asyncio 对象并不是线程安全的。其他线程不能随意操作 Task、Future 或事件循环对象。

从其他线程安排普通回调,应使用:

loop.call_soon_threadsafe(callback, *args)

从其他线程提交协程,应使用:

future = asyncio.run_coroutine_threadsafe(coro(), loop)
result = future.result()

这里返回的是 concurrent.futures.Future,不是 asyncio.Future。它适合在线程边界之外同步获取协程结果。Python 官方文档明确区分了这两种 Future,并提供了跨线程调度 API。(docs.python.org)


十五、一个端到端示例:并发 I/O、异常与关闭

下面的程序模拟三个异步请求,其中一个请求失败:

import asyncio
from collections.abc import Coroutine
from typing import Any

async def fetch(name: str, delay: float, fail: bool = False) -> str:
    print(f"{name}: start")
    await asyncio.sleep(delay)

    if fail:
        raise RuntimeError(f"{name}: failed")

    print(f"{name}: done")
    return f"{name}: result"

async def main() -> None:
    tasks: list[asyncio.Task[str]] = [
        asyncio.create_task(fetch("A", 0.3)),
        asyncio.create_task(fetch("B", 0.1, fail=True)),
        asyncio.create_task(fetch("C", 0.2)),
    ]

    try:
        results = await asyncio.gather(*tasks)
    except RuntimeError as exc:
        print(f"caught: {exc}")

        for task in tasks:
            if not task.done():
                task.cancel()

        await asyncio.gather(*tasks, return_exceptions=True)
    else:
        print(results)

if __name__ == "__main__":
    asyncio.run(main())

可能输出:

A: start
B: start
C: start
caught: B: failed

这里需要注意:

  1. 三个 Task 会先被安排;
  2. B 最先完成,但以异常结束;
  3. gather() 将异常传播给调用者;
  4. 主流程捕获异常后,主动取消尚未完成的任务;
  5. 使用 return_exceptions=True 等待取消收尾,避免取消异常再次打断清理过程;
  6. asyncio.run() 最后负责关闭事件循环和执行器。

在新代码中,TaskGroup 通常能更清晰地表达这类“子任务属于同一个作用域”的关系:

import asyncio

async def fetch(name: str, delay: float, fail: bool = False) -> str:
    await asyncio.sleep(delay)

    if fail:
        raise RuntimeError(f"{name} failed")

    return f"{name} ok"

async def main() -> None:
    try:
        async with asyncio.TaskGroup() as group:
            group.create_task(fetch("A", 0.3))
            group.create_task(fetch("B", 0.1, fail=True))
            group.create_task(fetch("C", 0.2))
    except* RuntimeError as errors:
        for error in errors.exceptions:
            print(error)

asyncio.run(main())

这段代码涉及异常组和结构化并发,关键基础是:任务不再是被创建后无人负责的孤立对象,而是绑定在 TaskGroup 的作用域中。其异常传播、兄弟任务取消和退出等待机制属于下一层主题。


十六、最终建立一个准确的 mental model

可以用下面四句话检查自己是否真正理解了 asyncio

  1. 调用协程函数只得到协程对象,不会自动执行。
  2. Task 是被事件循环调度的协程;Future 是表示未来结果的状态对象。
  3. await 只有在等待对象能够挂起当前协程时,才会把控制权交还给事件循环。
  4. 事件循环线程中的任何长时间同步代码,都会延迟同一线程中的所有 Task 和 I/O 回调。

从数据流看,完整路径是:

async def
  ↓ 调用
协程对象
  ↓ await 或 create_task
Task / 当前协程执行
  ↓ await I/O、定时器或 Future
事件循环接管
  ↓ selector / 定时器 / 回调
Future 完成
  ↓
恢复等待中的 Task
  ↓
返回值、异常或取消

asyncio 的性能和正确性都建立在同一个前提上:任务必须在等待外部资源时主动挂起,并且在退出时明确处理结果、异常、取消和资源关闭。理解了这条控制流,后续的 create_taskTaskGroup、异常组、结构化并发,以及同步代码与异步代码之间的边界,才会从 API 记忆变成可以推导的运行模型。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。