首页
/ CPython asyncio 高层 API 全景索引:从 asyncio.run 到 TaskGroup、Streams 与同步原语

CPython asyncio 高层 API 全景索引:从 asyncio.run 到 TaskGroup、Streams 与同步原语

2026-09-06 18:37:50作者:魏侃纯Zoe

asyncio 是 CPython 标准库中基于 async/await 语法的并发框架。官方文档中的 High-level API Index 以一张索引页的形式,把全部“高层”(high-level)异步 API 按 Tasks、Queues、Subprocesses、Streams、Synchronization、Exceptions 六大类收拢在一起。本文以该索引为骨架,逐一讲解每个 API 的职责、适用场景与底层实现出处,帮你把这页索引真正变成可落地的选型地图——读完即可知道“并发该用 gather 还是 wait_for”“跨线程调度该用 to_thread 还是 run_coroutine_threadsafe”“超时该用 wait_for 还是 timeout()”。

一、这份索引在 asyncio 文档体系中的位置

CPython 官方文档将 asyncio 拆分为入口页 asyncio.rst高层 API 索引页(即本文依据的 asyncio-api-index.rst)、底层 API 索引页 asyncio-llapi-index.rst,以及按主题展开的细节页(asyncio-task.rstasyncio-queue.rstasyncio-subprocess.rstasyncio-stream.rstasyncio-sync.rstasyncio-exceptions.rst 等)。

索引页开篇点明其定位:"This page lists all high-level async/await enabled asyncio APIs."——它不重复讲解细节,而是给出“功能 → API”的一览表。与之相对的底层(low-level)API(事件循环、传输层 Transport、协议层 Protocol、Future、策略 Policy 等)则由 asyncio-llapi-index.rst 单独索引。理解这一分层是使用 asyncio 的第一步:日常业务开发优先选择下文将要介绍的高层 API,只有在实现自定义协议、自研事件循环或直接操作传输层时才需要下沉到底层。

从实现层面看,这些高层 API 几乎全部落在 CPython 仓库的 Lib/asyncio 包中,与文档章节对应关系如下:

文档章节 核心源码文件 主要对象
Tasks tasks.pyrunners.pytaskgroups.pytimeouts.pythreads.py runTaskTaskGrouptimeout()
Queues queues.py QueuePriorityQueueLifoQueue
Subprocesses subprocess.py create_subprocess_exec/shell
Streams streams.py open_connectionstart_serverStreamReader/Writer
Synchronization locks.py LockEventSemaphoreBarrier
Exceptions exceptions.py CancelledErrorBrokenBarrierError

下文按原索引的章节顺序展开。

二、Tasks:任务的创建、并发等待、超时与跨线程调度

索引把 Tasks 定义为“运行 asyncio 程序、创建 Task、以及带超时地同时等待多个对象的工具集合”。这是 asyncio 使用频率最高、也最容易混淆的一类 API,可再细分为四个层次。

1. 程序入口:asyncio.run()asyncio.Runner

  • asyncio.run(main):创建事件循环、运行一个协程、结束后关闭循环。签名(见 runners.py)为 run(main, *, debug=None, loop_factory=None)。它会每次新建事件循环并在结束时关闭("always creates a new event loop and closes it at the end"),同时负责终结异步生成器、关闭默认执行器;若被从“正在运行的事件循环中”调用,会直接抛出 RuntimeError。因此它被定位为 asyncio 程序的主入口,理想情况下整个程序只调用一次。
import asyncio

async def main():
    await asyncio.sleep(1)
    print("hello")

asyncio.run(main())  # 程序主入口
  • asyncio.Runner:一个控制事件循环生命周期的上下文管理器(runners.py)。官方索引将其概括为"A context manager that simplifies multiple async function calls"。它的意义在于:run() 方法可在同一个 Runner 生命周期内被调用多次,适合交互式控制台(如 IPython)、unittest 运行器、命令行工具等“从同步框架多次调用异步代码”的场景。其 docstring 明确指出 asyncio.run(main(), debug=True) 等价于:
with asyncio.Runner(debug=True) as runner:
    runner.run(main())

在实现上,Runner.close()runners.py)会依次执行取消全部剩余任务、shutdown_asyncgens()、以及带超时(THREAD_JOIN_TIMEOUT)的 shutdown_default_executor(),随后关闭事件循环并把进程级 set_event_loop(None),清理非常完整。它还内置了 SIGINT 处理逻辑(首次 Ctrl-C 取消主任务、再次触发抛 KeyboardInterrupt,见 _on_sigint)。

2. Task 对象:TaskTaskGroupcreate_taskcurrent_taskall_tasks

  • asyncio.Task:索引中的定义很精炼——"Task object"。实现上它是“包在 Future 里的协程”(见 tasks.py 的类 docstring)。值得注意的是,纯 Python 的 Task 类定义在同文件靠前位置,而模块末尾(tasks.py)会尝试 import _asyncio用 C 加速实现 _asyncio.Task 覆盖同名符号——也就是说,正常安装的 CPython 中 asyncio.Task 实际指向 C 扩展版本,Python 版仅作为无 _asyncio 模块时的后备。Task 还支持 name= 参数自定义任务名(默认自动命名为 Task-1Task-2…)以及 eager_start 立即启动选项。
  • asyncio.create_task(coro):把协程包装成 Task 并安排其尽快运行,然后返回该 Task 对象。其实现只是取当前运行循环并调用 loop.create_task(coro)tasks.py),因此必须在已有正在运行的事件循环内调用(典型位置是另一个协程内部)。与 ensure_future() 相比,create_task 只接受协程且语义更明确,是官方推荐的“并发启动子任务”手段。
  • asyncio.current_task(loop=None):返回当前正在执行的 Task;若当前不在任何 Task 内则返回 None。常用于实现任务局部状态或在取消逻辑里定位自身。
  • asyncio.all_tasks(loop=None):返回某事件循环上“尚未完成”的全部 Task 集合(tasks.py)。实现会同时纳入普通调度任务与 eager 任务,用于关闭前的统一清理(Runner.close() 内部正是调用它来 _cancel_all_tasks)。
  • asyncio.TaskGroup:任务组的上下文管理器(taskgroups.py),其 docstring 中的范例即官方推荐用法:
async def main():
    async with asyncio.TaskGroup() as group:
        task1 = group.create_task(coro1())
        task2 = group.create_task(coro2())
    # 退出 with 块时,组内所有任务都已执行完毕

从源码可见其语义保证:退出上下文时等待组内所有任务;若某任务抛出了除 CancelledError 外的异常,则取消组内其余任务并等待它们退出,最终把多个异常合并为一个 ExceptionGroup 向上抛出(taskgroups.py__aexit__ 逻辑)。相比手工 gatherTaskGroup 能可靠处理“一个失败是否连带取消其余任务”的编排问题,且会即时暴露失败而不是干等到底。

3. 并发等待、超时与取消:sleepgatherwait_forshieldwaittimeout()as_completed

这一组 API 负责“同时等一堆东西”,但语义差异极大,是踩坑高发区。索引页特意为它们标注了 await 前缀以强调均为可等待操作。

  • await asyncio.sleep(delay, result=None):协程级休眠(非阻塞),到点后可选返回 result。它是“让出控制权”的最简单手段,常用于测试与轮询。
  • await asyncio.gather(*aws, return_exceptions=False):同时调度并等待多个 awaitable 完成,返回按传入顺序排列的结果列表。这是“并行等待同一批任务并收集全部结果”的标准选择。
  • await asyncio.wait_for(aw, timeout):给单个 awaitable 加超时(tasks.py)。超时发生时它会取消被等待对象并抛出 TimeoutError;若希望在超时后仍让后台任务继续跑,需要先用 shield() 包裹。实现细节上,timeout<=0 走一条特殊路径(先取消再等待完成,确保协程不会在未启动状态下泄漏);正常超时路径则直接复用 timeouts.timeout 上下文管理器(async with timeouts.timeout(timeout): return await fut)。
  • await asyncio.shield(aw):为 awaitable 提供“取消屏蔽”——当外层等待被取消时,被 shield 包裹的任务继续在后台运行。典型用途:清理/保存操作不应因客户端断开而中断。
  • await asyncio.wait(aws, *, timeout=None, return_when=ALL_COMPLETED):等待一批 Future/Task 完成到满足某个条件,返回 (done, pending) 两个集合(tasks.py)。其 docstring 特别强调两点:不抛 TimeoutError——超时时未完成对象放进 pending 集合返回;fs 集合中不允许直接传裸协程,必须先 create_task 包装。return_when 可选 FIRST_COMPLETED / FIRST_EXCEPTION / ALL_COMPLETED(常量定义在 tasks.py)。
  • asyncio.timeout(delay)asyncio.timeout_at(when):索引页注明它"Run with a timeout. Useful in cases when wait_for is not suitable",即在 wait_for 不适用的场景(如需要在超时发生后继续做清理或部分使用结果)使用。它以异步上下文管理器形式工作:
async def main():
    try:
        async with asyncio.timeout(1.5):
            await long_operation()
    except TimeoutError:
        print("操作超时,可在此继续做降级处理")

底层 Timeout 类(timeouts.py)把超时建模为一个"到期时刻":可用 when() 查询当前截止时间、用 reschedule(when) 动态延后/提前截止,超时取消的本质是到期时调用 task.cancel()timeout_at 则直接接受“循环时间戳”作为绝对截止点,适合做固定截止(deadline)语义。

  • for ... in asyncio.as_completed(aws):以迭代器方式按完成顺序逐个消费结果(“完成一个处理一个”,而非等全部完成后一次性拿列表)。支持 timeout= 参数。其实现基于一个内部 _AsCompletedIteratortasks.py),它把每个待等待对象 ensure_future 注册回调,完成的 Future 会按序放入内部 queue 供迭代器取出。

这五者如何取舍,可以浓缩为一张对照表:

需求 推荐 API 关键差异
并行等待一批并收集结果 gather 返回按顺序排列的结果列表
一批任务“同生共死”、失败即取消 TaskGroup 失败抛出 ExceptionGroup
单个操作限时,超时即放弃 wait_for 超时取消任务并抛 TimeoutError
限时后还要做清理/降级 timeout() 上下文管理器 退出块后可捕获 TimeoutError 继续处理
只关心完成状态、不收集结果 wait 返回 (done, pending),超时不抛错
按完成顺序逐个处理 as_completed 迭代器逐项产出

4. 跨线程协作:to_threadrun_coroutine_threadsafe

  • await asyncio.to_thread(func, /, *args, **kwargs):在独立 OS 线程中异步运行一个(通常是阻塞的)同步函数(threads.py)。它是“把同步阻塞调用扔出事件循环”的首选,例如 await asyncio.to_thread(requests.get, url),底层走事件循环的默认线程池执行器。
import asyncio

def blocking_io():
    return "阻塞操作的结果"

async def main():
    result = await asyncio.to_thread(blocking_io)
    print(result)
  • asyncio.run_coroutine_threadsafe(coro, loop):从另一个 OS 线程把协程调度到指定事件循环上执行(tasks.py),返回一个 concurrent.futures.Future 供调用线程做阻塞式结果获取。这是“异步主循环 + 同步 worker 线程”混合架构(如 GUI 回调、消息队列消费者里需要回投到主循环)的官方桥梁:
# 线程 A:运行着事件循环 loop
# 线程 B:向线程 A 投递协程
future = asyncio.run_coroutine_threadsafe(async_job(), loop)
result = future.result(timeout=5)   # 阻塞等待线程 A 的结果

配合索引页给出的“另见 Tasks 主题页”指引,这四类 API 的完整参数与异常细节可进一步参考 asyncio-task.rst(协程与任务主文档)、asyncio-runner.rst(Runner 专项)以及 asyncio-threading.rst(线程安全与跨线程调度专题)。

三、Queues:在 Task 之间分发工作负载

索引页对队列的定位是:"Queues should be used to distribute work amongst multiple asyncio Tasks, implement connection pools, and pub/sub patterns."——即任务之间传递消息的标准手段,典型用例包括任务分发(producer/consumer)、连接池、发布订阅模式。

asyncio 队列与 queue 模块的区别在于:它们是单事件循环内、协程间协作的队列,操作以 await 为核心(await q.get()await q.put(item)),不依赖操作系统级线程锁。三类队列定义在 queues.py

API 语义
asyncio.Queue(maxsize=0) FIFO 队列;maxsize=0 表示不限容量
asyncio.PriorityQueue 优先级队列,put 元素需可比较(如 (priority, item) 元组)
asyncio.LifoQueue 后进先出栈式队列

与线程版 queue 同名类在行为上尽量对齐,因此迁移成本低。一个经典的生产者-消费者骨架:

import asyncio

async def worker(name, q):
    while True:
        item = await q.get()          # 队列空则挂起等待
        try:
            print(f"{name} 处理 {item}")
        finally:
            q.task_done()             # 通知该任务已完成

async def main():
    q = asyncio.Queue(maxsize=10)
    workers = [asyncio.create_task(worker(f"w{i}", q)) for i in range(3)]
    for i in range(20):
        await q.put(i)
    await q.join()                    # 等待全部任务被消费完成
    for w in workers:
        w.cancel()

asyncio.run(main())

代码中用到了队列配套的 task_done() / await join() 协同机制:join() 会在所有已入队项都被消费并标记 task_done 前一直阻塞,是“等队列清空”的标准写法。队列的完整 API(put_nowaitget_nowaitemptyfullqsize 等)见 asyncio-queue.rst,官方还提供了用队列在多 Task 间分发工作负载的完整示例。

四、Subprocesses:启动子进程与执行 Shell 命令

该章节只有两个高层入口,覆盖“在 asyncio 程序里跑子进程”的两种形式,定义于 subprocess.py

API 语义与要点
await asyncio.create_subprocess_exec(program, *args, ...) 直接执行指定程序(不经 Shell),参数以列表/可变参数形式传递
await asyncio.create_subprocess_shell(cmd, ...) 通过系统 Shell 执行命令字符串

两者都返回 asyncio.subprocess.Process 对象,可 await proc.communicate() 与子进程交换 stdin/stdout/stderr,或 await proc.wait() 等待其结束并读取 proc.returncode。需要警惕的是:create_subprocess_shell 存在 shell 注入风险,对包含不可信输入的字符串命令应尽量避免;能用 exec 表达就优先用 exec

import asyncio

async def main():
    proc = await asyncio.create_subprocess_exec(
        "ls", "-l",
        stdout=asyncio.subprocess.PIPE,
        stderr=asyncio.subprocess.PIPE,
    )
    stdout, stderr = await proc.communicate()
    if proc.returncode == 0:
        print(stdout.decode())

asyncio.run(main())

子进程 API 在索引中仅出现两个顶层函数,但实际是一个完整子系统:Windows 上依赖 Proactor 事件循环,Unix 上使用子进程监视器(child watcher)等底层机制,均由 asyncio.rstasyncio-subprocess.rst 展开说明。

五、Streams:基于网络 IO 的高层读写抽象

索引把 Streams 定位为 "High-level APIs to work with network IO"。它们是“无需编写 Protocol 回调即可进行 TCP/Unix 网络通信”的高层封装,核心实现在 streams.py

API 职责 源码位置
await asyncio.open_connection(host, port) 建立 TCP 连接,返回 (reader, writer) streams.py
await asyncio.open_unix_connection(path) 建立 Unix 域套接字连接 streams.py
await asyncio.start_server(client_cb, host, port) 启动 TCP 服务端,为每个连接调用回调 streams.py
await asyncio.start_unix_server(client_cb, path) 启动 Unix 域套接字服务端 streams.py
asyncio.StreamReader 高层接收对象:await read() / readline() streams.py
asyncio.StreamWriter 高层发送对象:write() + await drain() streams.py

StreamReader/StreamWriter 构成的读写抽象屏蔽了底层流量控制细节——尤其注意 write() 只是写进缓冲区,真正“把数据冲刷出去并等待对端可写”要用 await writer.drain()。一个最小可运行的 echo 服务端/客户端对例如下:

import asyncio

async def handle(reader, writer):
    data = await reader.read(100)
    writer.write(data)            # 原样回显
    await writer.drain()          # 等待数据真正下发
    writer.close()
    await writer.wait_closed()

async def client():
    reader, writer = await asyncio.open_connection("127.0.0.1", 8888)
    writer.write(b"ping")
    await writer.drain()
    echo = await reader.read(100)
    print("收到回显:", echo)
    writer.close()
    await writer.wait_closed()

async def main():
    server = await asyncio.start_server(handle, "127.0.0.1", 8888)
    await asyncio.gather(client(), server.serve_forever())

asyncio.run(main())

Streams 是在其上构建 HTTP/WebSocket 等服务的基础层,完整 API 与更多协议细节(如 readexactlyreaduntillimit 缓冲上限及其可能引发的 LimitOverrunError/IncompleteReadError)见 asyncio-stream.rst,索引页给出的 TCP 客户端完整示例也可在该主题页找到。

六、Synchronization:线程风格同步原语在协程中的等价物

索引定位为 "Threading-like synchronization primitives that can be used in Tasks"——它们把 threading 的同步概念移植到协程世界,接口风格几乎一一对应,全部定义于 locks.py

API 对应 threading 概念 说明
asyncio.Lock threading.Lock 互斥锁,保护共享资源的临界区
asyncio.Event threading.Event 事件标志:await event.wait() 阻塞至 event.set()
asyncio.Condition threading.Condition 条件变量,等待特定条件满足
asyncio.Semaphore threading.Semaphore 信号量,限制并发进入的数量
asyncio.BoundedSemaphore threading.BoundedSemaphore 有界信号量,防止 release 次数超过 acquire 次数
asyncio.Barrier threading.Barrier 栅栏,等待 n 个参与者到齐后一起放行

与线程版最大的语法差异是等待全部以 await 表达,且大多可直接用作异步上下文管理器:

import asyncio

async def worker(lock, i):
    async with lock:              # 等价于 await lock.acquire() + finally lock.release()
        print(f"worker {i} 进入临界区")
        await asyncio.sleep(0.1)

async def main():
    lock = asyncio.Lock()
    await asyncio.gather(*(worker(lock, i) for i in range(5)))

asyncio.run(main())

信号量的典型用途是“限制对某资源的并发度”(如控制同时进行的 HTTP 请求数):

sem = asyncio.Semaphore(3)        # 最多 3 个并发

async def limited_request(url):
    async with sem:
        return await fetch(url)

Event 常用于任务间的“开始/停止”信号通知,Barrier 用于多任务对齐(如并行阶段化计算),各自的完整示例见 asyncio-sync.rst,索引页也单独链接了 Event 与 Barrier 的专项示例。需要明确的是,这些原语只在同一事件循环内有效,不能替代多线程/多进程场景下真正的跨线程锁。

七、Exceptions:asyncio 专属异常

索引的 Exceptions 章节只列了两个最具代表性的异常,但完整清单在 asyncio-exceptions.rst 中,全部实现在 exceptions.py

  • asyncio.CancelledError:任务被取消时抛出(见 Task.cancel())。这是 asyncio 取消机制的核心,其使用需特别注意两点:其一,在 exceptions.py 中它继承自 BaseException 而非 Exception,因此不会被 except Exception 捕获——如果你用裸 except Exception 兜底,取消信号会被正确穿透;若用 except BaseException 则需要先 except asyncio.CancelledError 处理取消。其二,取消并不总是“立即生效”:任务可能在取消请求到达时处于不可中断点,真正的取消点是下一个 await
  • asyncio.BrokenBarrierError:栅栏被破坏(如有参与者被取消或显式调用 Barrier.abort())时由 Barrier.wait() 抛出(exceptions.py),用于告知其余等待者“栅栏不再可用”。
import asyncio

async def main():
    task = asyncio.create_task(long_running())
    task.cancel()                  # 请求取消
    try:
        await task
    except asyncio.CancelledError:
        print("任务已被取消")       # 必须显式捕获才能拿到取消信号

asyncio.run(main())

其余常见异常还包括 TimeoutError(已与内置 TimeoutError 统一为同一对象,见 exceptions.py)、InvalidStateErrorIncompleteReadErrorLimitOverrunErrorSendfileNotAvailableError 等,它们大多由 Streams/底层传输触发,使用到相关功能时应一并查阅。

八、如何基于本仓库继续深入

这份索引页本身是“查询目录”,要形成完整技能闭环,建议在 CPython 仓库内做三层递进式阅读:

  1. 读文档主题页:每个章节对应的详细文档都承接了索引页的交叉引用——Tasks 详见 asyncio-task.rst,队列见 asyncio-queue.rst,子进程见 asyncio-subprocess.rst,流见 asyncio-stream.rst,同步原语见 asyncio-sync.rst,异常与事件循环主题分别见 asyncio-exceptions.rstasyncio-eventloop.rst;若需要进入底层,可对照 asyncio-llapi-index.rst
  2. 读标准库源码:本文所有 API 的权威实现都在 Lib/asyncio 包内——tasks.py(Task/等待原语)、runners.py(run/Runner)、taskgroups.pytimeouts.pythreads.pyqueues.pylocks.pystreams.pysubprocess.pyexceptions.py。docstring 往往比文档更贴代码,如 RunnerTaskGroupwait_for 的行为约束都直接写在 docstring 里,且常量默认值(waitreturn_when=ALL_COMPLETED、队列的 maxsize=0)均以源码为准。
  3. 跑测试验证语义:CPython 仓库在 Lib/test 下维护着庞大的回归测试集。当你不确定某个 API 的边界行为(例如 wait_for 超时后任务是否被真正取消、TaskGroup 如何合并多异常)时,阅读对应测试用例是最快的实证方式。

小结:从“索引页”到“正确的 API 直觉”

回到 asyncio-api-index.rst 的六大分组,可以提炼出一条使用主线:asyncio.run() 启动程序 → TaskGroup/create_task 建立并发 → gather/as_completed/wait 选择等待策略 → wait_for/timeout() 控制时限 → Queue 与同步原语在任务间协调 → to_thread/run_coroutine_threadsafe 打通线程边界 → Streams/Subprocesses 负责 IO。把这 27 个高层 API 的分工与“彼此差异”记牢,就能绕开绝大多数“该用哪个函数”的困惑,而这张索引正是把全部答案压缩在一页之内的最佳速查表。

登录后查看全文
热门项目推荐
相关项目推荐