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 秒
串行执行时,总耗时近似为:
如果两个任务可以同时发起网络请求,那么等待期间不需要占用 Python 代码的执行时间,总耗时接近:
这叫并发:任务的生命周期发生重叠。
并行则要求多个任务在同一时刻由不同的执行单元运行,例如多个线程、多个进程或多个解释器。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 官方文档明确指出,单纯调用协程函数不会把它加入调度;需要通过 await、asyncio.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 中的对象称为可等待对象,英文是 awaitable。asyncio 中最主要的三类可等待对象是:
- 协程对象;
asyncio.Task;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. 事件循环延迟的因果关系
设事件循环线程正在执行一个阻塞调用,阻塞时间为 。在这段时间内:
- 已经就绪的其他 Task 不能运行;
- 已到期的定时器不能及时执行;
- 已经准备好的 I/O 回调不能及时处理;
- 等待这些 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 表示“某个结果现在是否已经准备好”。
Task 是 Future 的子类,因此 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())
执行过程是:
future初始为PENDING;set_result_later()被创建为 Task;main()等待future,暂时挂起;- 一秒后,辅助 Task 调用
future.set_result("done"); - Future 状态变成
DONE; - 等待 Future 的
main()被重新安排到事件循环; 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() 会:
- 创建事件循环;
- 执行传入的可等待对象;
- 等待顶层协程完成;
- 结束异步生成器;
- 关闭默认执行器;
- 关闭事件循环。
它不能在同一线程中已有事件循环运行时再次调用,通常应作为应用程序的顶层入口,并且一个程序通常只调用一次。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
这里需要注意:
- 三个 Task 会先被安排;
B最先完成,但以异常结束;gather()将异常传播给调用者;- 主流程捕获异常后,主动取消尚未完成的任务;
- 使用
return_exceptions=True等待取消收尾,避免取消异常再次打断清理过程; 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:
- 调用协程函数只得到协程对象,不会自动执行。
- Task 是被事件循环调度的协程;Future 是表示未来结果的状态对象。
await只有在等待对象能够挂起当前协程时,才会把控制权交还给事件循环。- 事件循环线程中的任何长时间同步代码,都会延迟同一线程中的所有 Task 和 I/O 回调。
从数据流看,完整路径是:
async def
↓ 调用
协程对象
↓ await 或 create_task
Task / 当前协程执行
↓ await I/O、定时器或 Future
事件循环接管
↓ selector / 定时器 / 回调
Future 完成
↓
恢复等待中的 Task
↓
返回值、异常或取消
asyncio 的性能和正确性都建立在同一个前提上:任务必须在等待外部资源时主动挂起,并且在退出时明确处理结果、异常、取消和资源关闭。理解了这条控制流,后续的 create_task、TaskGroup、异常组、结构化并发,以及同步代码与异步代码之间的边界,才会从 API 记忆变成可以推导的运行模型。
系列导航与关联阅读
- 系列入口:Python 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python CLI 工程:argparse、子命令、退出码、管道和可测试性
- 下一篇:Python 异步任务:create_task、TaskGroup、异常组和结构化并发
- 延伸:Python GIL 与自由线程构建:互斥边界、扩展兼容和并行选择
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论