JEPA4Japan · 教程

使用 asyncio 进行异步 I/O

2,233字 7分钟阅读 #Python

在不阻塞事件循环的情况下协调协程、任务、超时和取消操作。

课程进度 课程大纲 已发布 24/24 课

一个事件循环传递一根接力棒

调用 async def 函数会创建一个协程对象。它的函数体此时还没有运行。你可以直接等待它,也可以将它包装为一个 Task,让事件循环对它进行调度。

  1. 协程对象 可以开始执行的工作
  2. Task 已调度的协程
  3. 事件循环 选择一个已就绪的任务
  4. await 可能会交出接力棒
Asyncio 通常在一个事件循环线程上协调许多等待操作。
import asyncio


async def hello():
    print("body runs")
    return "done"


async def main():
    coroutine = hello()
    print(type(coroutine).__name__)
    print(await coroutine)


asyncio.run(main())
coroutine
body runs
done

asyncio.run() 会创建一个入口事件循环,运行 main(),然后关闭该循环。不要在已有的事件循环中嵌套调用它。

只有当工作暂停时,Await 才会交出控制权

await 的意思是“我需要这个结果”。只有当被等待的操作确实暂停时,另一个任务才有机会运行。等待一个立即返回的协程,可能根本不会交出控制权。

  1. 任务 A 运行 普通代码会继续握着接力棒
  2. 真正的等待 I/O、事件、队列、计时器
  3. 任务 B 运行 事件循环传递接力棒
协作发生在真正的暂停点,而不是每一个写出来的 await 处。

阻塞调用和长时间运行的 CPU 循环会冻结事件循环线程;仅仅添加 async,既不会改变计算成本,也不会改变调度方式。

任务需要有所有者

asyncio.create_task() 会立即调度工作。请保留一个强引用,并明确由谁等待它的结果、观察它的失败,以及在关闭期间取消它。

asyncio.gather() 会调度可等待对象,并按照参数顺序返回结果,而不是按照完成顺序:

import asyncio


async def label(name, delay):
    await asyncio.sleep(delay)
    return name


async def main():
    results = await asyncio.gather(
        label("slow", 0.02),
        label("fast", 0.01),
    )
    print(results)


asyncio.run(main())
['slow', 'fast']

当某个普通子任务失败时,默认的 gather() 会传播第一个观察到的异常;它不会自动取消其他已提交的子任务。return_exceptions=True 会把失败转换为结果值,只有在明确约定每一项都可以独立失败时,这才合适。

Python 3.11+ 提供了用于结构化并发的 TaskGroup:

async def collect_structured(sources):
    tasks = []
    async with asyncio.TaskGroup() as group:
        for source in sources:
            tasks.append(group.create_task(fetch(source)))
    return [task.result() for task in tasks]

每个子任务都属于这个 async with 作用域。如果某个子任务因 CancelledError 以外的异常而失败,组会取消其余子任务,等待它们完成清理,然后在退出该作用域时抛出一个 ExceptionGroup。这是一个共同成败的边界;不要意外发布部分成功的结果。

截止时间会取消其拥有的作用域

Python 3.11+ 的 asyncio.timeout() 可以为一次或多次等待设置同一个截止时间:

async def load_with_deadline():
    try:
        async with asyncio.timeout(0.5):
            return await fetch_everything()
    except TimeoutError:
        return "deadline exceeded"

超时会取消上下文内部的当前任务,并在离开上下文时把这次取消转换为 TimeoutError,因此应在上下文外部捕获 TimeoutError。它限制的是所有者的等待时间;它无法证明远程操作已经回滚。

取消是协作式的。task.cancel() 会安排在之后的某个暂停点抛出 CancelledError。请在 finally 中释放资源;或者仅为了执行取消专用的清理而捕获它,然后重新抛出:

async def worker(resource):
    try:
        await resource.use()
    except asyncio.CancelledError:
        await resource.close()
        raise
  1. 截止时间 所有者停止等待
  2. 取消请求 在暂停时送达
  3. 清理 释放所拥有的资源
  4. 传播 重新抛出取消异常
清理是取消过程的一部分;吞掉取消异常会破坏结构化的所有权关系。

CancelledError 直接继承自 BaseException。结构化工具会在内部使用它,因此只有当某个组件有意负责一套完整的替代策略时,才应该抑制它。

转移阻塞 I/O,而不是繁重的 CPU 工作

asyncio.to_thread(blocking_function, ...) 可以适配旧式的阻塞 I/O 调用,使事件循环线程保持响应。请限制同时进行的调用数量,并使用阻塞 API 自带的超时机制。取消正在等待的任务,并不能强制已经在线程中运行的任意代码停止。

不要仅仅因为大量纯 Python CPU 工作是同步的,就把它交给 to_thread()。在标准 CPython 3.12 的 GIL 下,它通常仍然无法并行运行 Python 字节码,而且可能损害事件循环的响应速度。微小的计算可以直接执行;对于经过测量确认、相互独立的 CPU 任务,可以考虑使用有界进程池或单独的工作进程。

小项目:离线元数据收集器

创建 metadata_collector.py。这个仅使用 Python 3.9+ 标准库的项目会使用本地延迟,而不是访问网络。显式的任务引用、gather()、排序和 wait_for() 可以带来确定性的结果和可测试的截止时间。在 Python 3.11+ 中,如果多个等待操作共享同一份时间预算,可以用超时作用域替代 wait_for()。

import asyncio
from dataclasses import dataclass


@dataclass(frozen=True)
class Source:
    slug: str
    title: str
    words: int
    delay: float


async def fetch(source):
    await asyncio.sleep(source.delay)
    return (
        source.slug,
        source.title.strip(),
        max(1, (source.words + 199) // 200),
    )


async def collect(sources):
    sources = list(sources)
    slugs = [source.slug for source in sources]
    if len(slugs) != len(set(slugs)):
        raise ValueError("slugs must be unique")

    tasks = [asyncio.create_task(fetch(source)) for source in sources]
    results = await asyncio.gather(*tasks)
    return sorted(results)


async def main():
    sources = [
        Source("models", " Open models ", 620, 0.02),
        Source("agents", "Reliable agents", 340, 0.01),
        Source("evals", "Practical evals", 199, 0.00),
    ]
    items = await collect(reversed(sources))

    print("Article metadata")
    for slug, title, minutes in items:
        print(f"- {slug}: {title} ({minutes} min)")

    timed_out = False
    try:
        await asyncio.wait_for(
            fetch(Source("slow", "Slow", 10, 10)),
            timeout=0.01,
        )
    except asyncio.TimeoutError:
        timed_out = True

    print(f"Sorted: {[slug for slug, _, _ in items]}")
    print(f"Timeout observed: {timed_out}")


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

运行 python3 metadata_collector.py:

Article metadata
- agents: Reliable agents (2 min)
- evals: Practical evals (1 min)
- models: Open models (4 min)
Sorted: ['agents', 'evals', 'models']
Timeout observed: True

输入顺序被反转,而且延迟各不相同,但发布的顺序仍然经过排序。请记住这个边界:创建任务的成本很低,下游容量却不是无限的。真实服务需要信号量或其他明确的并发限制、自身的超时设置,以及有文档说明的重试策略。

三个小任务与第 24 章检查清单

  1. 观察接力棒。 比较逐个等待三个获取操作和使用 gather() 的区别,同时保持相同的结果顺序。
  2. 验证清理。 在一个 Event 确认任务已启动后取消该任务;断言它的 finally 块已经运行,并确认等待它时仍然会抛出 CancelledError。
  3. 共同成败。 在 Python 3.11+ 中,让一个 TaskGroup 子任务失败,并验证同组中正在等待的另一个子任务会在组抛出异常之前完成清理。

当你做到以下几点时,就可以开始学习第 24 章:

  • 你能区分协程对象、Task 和事件循环;
  • 你知道,只有当可等待对象确实暂停时,await 才会交出控制权;
  • 每个创建出来的任务都有一个所有者负责等待它并观察其结果;
  • 你知道 gather() 会保持参数顺序,但不会让任务共同成败;
  • 你知道 TaskGroup 拥有其子任务,并会在发生失败时取消同组任务;
  • 你会围绕预期的作用域设置截止时间;
  • 在 CancelledError 继续向上传播之前,清理工作会先运行;
  • to_thread() 用于适配有界的阻塞 I/O,而不是实现 CPU 并行;
  • CPU 密集型工作绝不会阻塞事件循环线程;
  • 你可以运行收集器,并复现其稳定的输出。

接下来,你将测量一个完整的程序,只改进已经证实存在的瓶颈,对程序进行打包,并使用发布检查清单将它正式发布。