用协程写好并发代码

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. asyncawait 的规则:用错误示例说明边界

要说明的理论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

这个例子直接证明了

  • 协程函数被调用时只是创建了一个协程对象,不会运行任何代码
  • 必须用 awaitasyncio.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

Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐