Python 中的 asyncio(异步 I/O 库)
用协程写好并发代码
Python 的 asyncio 库让你能用 async/await 写并发程序。它的核心是协程和事件循环,适合在一个线程里处理大量 I/O 密集型任务(比如网络请求、文件读写)。
这篇文章会带你理解异步I/O的原理,并用手敲代码的方式,学会定义、运行协程,以及什么时候该用它。
1. 同步 vs 异步:用等待时间证明问题
要说明的理论:同步代码在等待时会阻塞整个线程;异步代码在等待时会交出控制权,让其他任务运行。
为了证明这一点,我们做一个最简单的对比:执行三次“打印A → 等1秒 → 打印B”。
同步版本(阻塞等待):
import time
def sync_task():
print("A")
time.sleep(1) # 这里整个线程被卡住,什么也做不了
print("B")
def sync_main():
for _ in range(3):
sync_task()
start = time.perf_counter()
sync_main()
print(f"同步耗时: {time.perf_counter() - start:.2f}秒")
运行结果:
A
B
A
B
A
B
同步耗时: 6.02秒
异步版本(非阻塞等待):
import asyncio
async def async_task():
print("A")
await asyncio.sleep(1) # 这里交出控制权,事件循环可以干别的
print("B")
async def async_main():
# 同时启动三个任务
await asyncio.gather(async_task(), async_task(), async_task())
start = time.perf_counter()
asyncio.run(async_main())
print(f"异步耗时: {time.perf_counter() - start:.2f}秒")
运行结果:
A
A
A
B
B
B
异步耗时: 2.01秒
这个例子直接证明了:
- 三个
A几乎同时打印 → 三个任务被同时启动 - 然后集体等待1秒 → 只等了1秒,而不是3秒
await asyncio.sleep(1)没有阻塞,它只是“登记”了一个1秒后的唤醒,然后立刻让出控制权
2. async 和 await 的规则:用错误示例说明边界
要说明的理论:await 只能在 async def 内部使用;async def 函数调用后不会立即执行。
import asyncio
async def correct_coro():
print("协程开始")
await asyncio.sleep(0.5) # await 必须出现在 async def 内部
print("协程结束")
return 42
# 演示1:调用协程函数不会立即执行
result = correct_coro()
print(f"调用后得到: {result}") # <coroutine object ...>,不是42
print(f"它还没有执行,因为没有被 await 或 asyncio.run()")
# 演示2:await 不能在普通函数里
def normal_function():
# await asyncio.sleep(0) # 去掉注释会报 SyntaxError
pass
# 正确的执行方式
async def main():
value = await correct_coro() # 这里才真正执行并拿到返回值
print(f"await 后拿到: {value}")
asyncio.run(main())
输出:
调用后得到: <coroutine object correct_coro at 0x...>
它还没有执行,因为没有被 await 或 asyncio.run()
协程开始
协程结束
await 后拿到: 42
这个例子直接证明了:
- 协程函数被调用时只是创建了一个协程对象,不会运行任何代码
- 必须用
await或asyncio.run()才能真正执行 await不能出现在普通函数中,否则语法错误
3. 事件循环的工作机制:用一个“你来我往”的例子说明
要说明的理论:事件循环会在一个协程遇到 await 时,暂停它并切换到另一个就绪的协程。
下面用两个协程交替执行来证明这一点:
import asyncio
async def talker(name, wait_time, times):
for i in range(times):
print(f"{name}: 第{i+1}次说话")
await asyncio.sleep(wait_time) # 每次说完就交出控制权
print(f"{name}: 说完了")
async def main():
# 创建两个协程,一个等0.3秒,一个等0.7秒
await asyncio.gather(
talker("快速", 0.3, 4),
talker("慢速", 0.7, 2),
)
asyncio.run(main())
输出:
快速: 第1次说话
慢速: 第1次说话
快速: 第2次说话
快速: 第3次说话
慢速: 第2次说话
快速: 第4次说话
快速: 说完了
慢速: 说完了
这个例子直接证明了:
- 两个协程并不是一个接一个执行,而是交替进行
- 每当一个协程执行
await,事件循环就切换到另一个 - “快速”协程因为等待时间短,被执行得更频繁
- 这就是事件循环的核心工作方式:不是同时运行,而是轮流运行
4. 协程链:用“先拿用户信息再拿文章”证明依赖关系
要说明的理论:当一个协程的执行必须等待另一个协程的结果时,用 await 串联起来,形成链式依赖。
import asyncio
import random
async def step1_get_user(user_id):
"""第一步:获取用户基本信息(必须完成才能进入第二步)"""
delay = random.uniform(0.5, 1.0)
print(f" [步骤1] 获取用户{user_id}信息,需{delay:.1f}秒")
await asyncio.sleep(delay)
print(f" [步骤1] 用户{user_id}信息获取完成")
return {"id": user_id, "name": f"用户{user_id}"}
async def step2_get_orders(user):
"""第二步:获取用户的订单列表(依赖第一步返回的user对象)"""
delay = random.uniform(0.5, 1.0)
print(f" [步骤2] 获取{user['name']}的订单,需{delay:.1f}秒")
await asyncio.sleep(delay)
orders = [f"订单{i}" for i in range(1, 3)]
print(f" [步骤2] {user['name']}的订单获取完成")
return orders
async def process_user(user_id):
"""串联两个步骤:必须等第一步完成才能做第二步"""
user = await step1_get_user(user_id) # 等待第一步
orders = await step2_get_orders(user) # 等待第二步(依赖第一步的结果)
return user["name"], orders
async def main():
# 并发处理三个用户,但每个用户内部是串行的
results = await asyncio.gather(
process_user(1),
process_user(2),
process_user(3),
)
for name, orders in results:
print(f"{name}: {orders}")
asyncio.run(main())
输出(关键看顺序):
[步骤1] 获取用户1信息,需0.7秒
[步骤1] 获取用户2信息,需0.5秒
[步骤1] 获取用户3信息,需0.9秒
[步骤1] 用户2信息获取完成
[步骤2] 获取用户2的订单,需0.6秒
[步骤1] 用户1信息获取完成
[步骤2] 获取用户1的订单,需0.8秒
[步骤1] 用户3信息获取完成
[步骤2] 获取用户3的订单,需0.7秒
[步骤2] 用户2的订单获取完成
[步骤2] 用户1的订单获取完成
[步骤2] 用户3的订单获取完成
用户2: ['订单1', '订单2']
用户1: ['订单1', '订单2']
用户3: ['订单1', '订单2']
这个例子直接证明了:
- 对于同一个用户,步骤2 的打印永远出现在步骤1 的打印之后
- 但不同用户的步骤1 是并发执行的
await保证了依赖顺序:user必须先拿到,才能传给step2_get_orders
5. 队列 + 生产者消费者:用“任务分发”证明解耦能力
要说明的理论:当生产者和消费者的速度不匹配,且没有直接依赖时,用队列解耦,让多消费者自动负载均衡。
import asyncio
import random
async def producer(queue, total):
"""生产者:只管往队列里放任务,不管谁处理"""
for i in range(total):
await asyncio.sleep(random.uniform(0.1, 0.3)) # 生产间隔
task = f"任务_{i}"
await queue.put(task)
print(f"[生产] 放入 {task},队列大小={queue.qsize()}")
# 放结束标记
await queue.put(None)
print("[生产] 全部生产完毕")
async def consumer(queue, name):
"""消费者:谁空闲谁从队列取任务"""
while True:
task = await queue.get()
if task is None:
await queue.put(None) # 把结束标记传给下一个消费者
print(f"[消费 {name}] 收到结束信号,退出")
break
print(f"[消费 {name}] 开始处理 {task}")
await asyncio.sleep(random.uniform(0.4, 0.8)) # 处理耗时
print(f"[消费 {name}] 完成 {task}")
async def main():
q = asyncio.Queue(maxsize=2) # 队列最多缓存2个任务
await asyncio.gather(
producer(q, 5),
consumer(q, "甲"),
consumer(q, "乙"),
)
asyncio.run(main())
输出(注意谁空闲谁拿任务):
[生产] 放入 任务_0,队列大小=1
[消费 甲] 开始处理 任务_0
[生产] 放入 任务_1,队列大小=1
[消费 乙] 开始处理 任务_1
[生产] 放入 任务_2,队列大小=1
[消费 甲] 完成 任务_0
[消费 甲] 开始处理 任务_2
[生产] 放入 任务_3,队列大小=1
[消费 乙] 完成 任务_1
[消费 乙] 开始处理 任务_3
[生产] 放入 任务_4,队列大小=1
[消费 甲] 完成 任务_2
[消费 甲] 开始处理 任务_4
[消费 乙] 完成 任务_3
[生产] 全部生产完毕
[消费 甲] 完成 任务_4
[消费 甲] 收到结束信号,退出
[消费 乙] 收到结束信号,退出
这个例子直接证明了:
- 生产者只管放,消费者只管取,两者不直接调用对方
- 队列满了会自动阻塞生产者,空了会阻塞消费者
- 两个消费者自动负载均衡:谁先完成当前任务,谁就去拿下一个
None作为“毒丸”信号,优雅通知消费者结束
6. as_completed:用“先完先报”证明它和 gather 的区别
要说明的理论:gather 等待所有任务完成后统一返回;as_completed 谁先完成谁先返回,适合做进度反馈。
import asyncio
async def task_with_id(task_id, duration):
await asyncio.sleep(duration)
return f"任务{task_id}完成(耗时{duration}秒)"
async def demo_gather():
print("=== gather 演示 ===")
results = await asyncio.gather(
task_with_id(1, 3),
task_with_id(2, 1),
task_with_id(3, 2),
)
print(f"gather 一次性返回所有结果: {results}")
# 注意:虽然任务2只用了1秒,但我们要等3秒才能看到任何结果
async def demo_as_completed():
print("\n=== as_completed 演示 ===")
tasks = [
task_with_id(1, 3),
task_with_id(2, 1),
task_with_id(3, 2),
]
for coro in asyncio.as_completed(tasks):
result = await coro
print(f"as_completed 立即返回: {result}")
# 任务2完成的第一时间就被打印出来了
async def main():
await demo_gather()
await demo_as_completed()
asyncio.run(main())
输出:
=== gather 演示 ===
gather 一次性返回所有结果: ['任务1完成(耗时3秒)', '任务2完成(耗时1秒)', '任务3完成(耗时2秒)']
=== as_completed 演示 ===
as_completed 立即返回: 任务2完成(耗时1秒)
as_completed 立即返回: 任务3完成(耗时2秒)
as_completed 立即返回: 任务1完成(耗时3秒)
这个例子直接证明了:
gather:必须等最慢的那个(3秒),才能拿到任何结果as_completed:最快的任务(1秒)一完成,马上就能处理它的结果- 适用场景区别:需要全部结果用
gather;需要实时反馈进度用as_completed
7. 异常处理:用多个不同异常证明 return_exceptions 的作用
要说明的理论:gather 默认遇到第一个异常就会抛出,导致其他任务的结果丢失。设置 return_exceptions=True 可以把异常当结果返回,逐个处理。
import asyncio
async def safe_task(name, fail):
await asyncio.sleep(0.2)
if fail:
raise ValueError(f"{name} 失败了")
return f"{name} 成功"
async def demo_default_behavior():
print("=== 默认行为:一个异常导致全部丢失 ===")
try:
results = await asyncio.gather(
safe_task("A", False),
safe_task("B", True), # 这个会抛异常
safe_task("C", False),
)
print(results)
except ValueError as e:
print(f"捕获到异常: {e}")
print("注意:任务A和C的结果也丢失了,因为gather提前中断了")
async def demo_return_exceptions():
print("\n=== return_exceptions=True:异常当结果返回 ===")
results = await asyncio.gather(
safe_task("A", False),
safe_task("B", True),
safe_task("C", False),
return_exceptions=True # 关键参数
)
for r in results:
if isinstance(r, Exception):
print(f"这是一个异常: {r}")
else:
print(f"正常结果: {r}")
async def main():
await demo_default_behavior()
await demo_return_exceptions()
asyncio.run(main())
输出:
=== 默认行为:一个异常导致全部丢失 ===
捕获到异常: B 失败了
注意:任务A和C的结果也丢失了,因为gather提前中断了
=== return_exceptions=True:异常当结果返回 ===
正常结果: A 成功
这是一个异常: B 失败了
正常结果: C 成功
这个例子直接证明了:
- 默认情况下,一个任务抛异常,
gather会立即抛出,其他已完成任务的结果无法获取 - 设置
return_exceptions=True后,所有任务都会执行完毕,异常被包装成对象放在结果列表里 - 你可以遍历结果,分别处理正常值和异常值
总结:什么时候用 asyncio?
| 应该用 asyncio | 不应该用 asyncio |
|---|---|
| 程序大部分时间在等待网络、文件、数据库 | 程序大部分时间在做数学计算、视频处理 |
| 需要同时处理成千上万个连接 | 只有少量几个并发任务 |
| 想避免多线程的复杂性和竞态条件 | 任务需要真正利用多核CPU |
| 调用的库已经支持 async/await | 所有依赖库都是同步阻塞的 |
一句话原则:如果你的程序经常在等(网络、磁盘、数据库),asyncio 能帮你把等待的时间利用起来;如果你的程序一直在算(循环、矩阵、加解密),asyncio 帮不了你,该用 multiprocessing。
更多推荐



所有评论(0)