线程是实现并发的经典手段,但内核调度、栈内存、锁竞争、上下文切换都不是免费的。当瓶颈是「等」而不是「算」(等网络、磁盘、数据库),开几万个线程去各自阻塞并不划算。asyncio 给了另一条路:单线程 + 事件循环 + 协作式调度。一个线程里跑成千上万个协程,谁 await 了就把执行权让出去,事件循环负责在「数据就绪」时把它接回来。

难点不在语法(async/await 很简单),而在心智模型:你必须始终清楚「现在谁在跑」「这一步会让出控制权吗」「阻塞了会发生什么」。本文按「概念 -> 调度 -> 并发原语 -> 桥接同步代码 -> 陷阱」的顺序拆解。示例以 Python 3.12 为基准,并标注 3.11+ 的新写法。

协程、事件循环与 await

理解 asyncio,先抓住三件事:

  • 协程(coroutine):用 async def 定义的函数,调用它不立即执行,而是返回一个「协程对象」,交给事件循环调度才真正运行。
  • 事件循环(event loop):不断轮询的任务调度器,维护一组就绪的协程和回调,哪个等的东西就绪了,就把执行权交还给它。
  • await:协程内部用它暂停自己、把控制权还给事件循环,直到被 await 的可等待对象完成。
1
2
3
4
5
6
7
8
9
import asyncio

async def greet(name: str) -> str:
print(f"开始问候 {name}")
await asyncio.sleep(0.1) # 模拟 IO:暂停 0.1s,期间让出控制权
return f"hello {name}"

# greet("ada") 直接调用只返回协程对象,不执行
result = asyncio.run(greet("ada")) # 'hello ada'

asyncio.run() 是程序入口:建循环、跑协程、关循环。每个程序只应有一个顶层 asyncio.run,不要在已运行的事件循环里再调它。

能被 await 的东西统称 awaitable,日常只需记住三类:协程对象、Task(被调度的协程)、Future(低层「将来有结果」对象,多数被 Task 包装,日常代码很少直接用)。

Task:把协程真正调度起来

直接 await coro串行的–当前协程停下来等它完成。要并发,必须先把协程包装成 Task 交给事件循环:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
import asyncio, time

async def work(name: str, seconds: float) -> str:
await asyncio.sleep(seconds)
return f"{name} done in {seconds}s"

async def main():
start = time.perf_counter()
t1 = asyncio.create_task(work("a", 1.0)) # 立即排入调度,返回 Task
t2 = asyncio.create_task(work("b", 1.0)) # 两个 task 已在并发跑
r1, r2 = await t1, await t2 # 等它们都完成
print(r1, r2, f"总耗时 {time.perf_counter()-start:.2f}s") # ~1.0s,不是 2.0s

asyncio.run(main())

关键区别:await work(...) 立即执行并等待(串行);create_task(work(...)) 立即排入调度并返回句柄,当前协程继续往下走(并发)。一个常见错误是「写了 async 但全程 await coro」,结果异步代码跑得和同步一样慢。

注意 asyncio.run 要求当前线程没有正在运行的事件循环,所以不能在协程内部再调它–在协程内部跑另一个协程直接 awaitcreate_task。在 Jupyter、某些测试框架里事件循环已在跑,这时改用 await main()

并发执行:gatherTaskGroupas_completed

asyncio.gather:等全部完成

1
2
3
4
5
6
7
8
9
10
async def fetch(url: str, seconds: float) -> str:
await asyncio.sleep(seconds)
return f"data from {url}"

async def main():
results = await asyncio.gather(
fetch("/a", 0.3), fetch("/b", 0.1), fetch("/c", 0.2),
)
# ['data from /a', 'data from /b', 'data from /c'],顺序与传入一致
# 总耗时 ~0.3s(取最慢的那个),而非 0.6s

gather 返回结果顺序与传入顺序一致,与完成先后无关。默认任一任务抛异常会立即传播,其余任务不会被取消。要改变行为:

1
2
3
await asyncio.gather(*tasks, return_exceptions=False)  # 默认,异常立即抛
results = await asyncio.gather(*tasks, return_exceptions=True)
# 成功的位是值,失败的位是 Exception 实例,适合「尽力而为」的批量请求

TaskGroup(3.11+):更安全的并发

TaskGroupgather 的结构化替代品,最大的区别是异常时自动取消组内所有任务

1
2
3
4
5
6
7
async def main():
async with asyncio.TaskGroup() as tg:
t1 = tg.create_task(fetch("/a", 0.3))
t2 = tg.create_task(fetch("/b", 0.1))
t3 = tg.create_task(fetch("/c", 0.2))
# 离开 async with 时所有任务都已完成
print(t1.result(), t2.result(), t3.result())

语义是「要么全成功,要么一起失败」:任一任务抛异常,组内其余任务被取消,所有异常打包成 ExceptionGroup 抛出。新代码优先用 TaskGroup,错误处理边界比 gather 清晰得多。

as_completed:谁先完成谁先处理

gather 要等全部完成才返回。想在每个任务完成时立刻处理(如流式输出),用 as_completed

1
2
3
4
tasks = [fetch(f"/{c}", 0.3 - i * 0.08) for i, c in enumerate("abc")]
for coro in asyncio.as_completed(tasks):
result = await coro # 谁先完成就先 await 到谁
print(result) # 完成顺序:b, c, a

日常用 gather/TaskGroup 就够了。还有 asyncio.wait,返回 (done, pending) 两个集合,配合 return_when=FIRST_COMPLETED 用于「第一个完成就走」的竞速逻辑(如超时竞速),它不返回结果本身,需要自己调 .result()

超时与取消:让协程可以被中断

asyncio.timeout(3.11+)

1
2
3
4
5
6
7
8
9
10
async def slow() -> str:
await asyncio.sleep(10)
return "done"

async def main():
try:
async with asyncio.timeout(0.5):
result = await slow()
except TimeoutError:
print("超时了") # 0.5s 后触发

它会在超时后取消作用域内被 await 的任务,抛出 TimeoutError。上下文管理器形式能覆盖一段含多个 await 的代码块,比旧的 wait_for(coro, timeout=0.5) 更灵活。wait_for 仍可用,功能等价。

取消的传播

取消通过 CancelledError 实现。事件循环取消一个 Task 时,会在该协程当前 await 的位置抛入 asyncio.CancelledError。协程可以捕获它做清理,但不应该吞掉

1
2
3
4
5
6
7
async def work():
try:
await asyncio.sleep(10)
except asyncio.CancelledError:
await cleanup() # 必要的清理:关连接、回滚事务
raise # 重新抛出,让取消正常传播
return "done"

CancelledError 在 3.8+ 继承自 BaseException 而非 Exception,所以 except Exception: 不会误捕它–这是有意为之,避免普通异常处理代码意外吞掉取消信号。task.cancel() 主动取消一个任务,它会在任务下次被调度时抛 CancelledError。取消是协作式的,被取消的协程得跑到下一个 await 才能感知到。

异步上下文管理器、迭代器与生成器

异步世界里 withfor、生成器都有异步版本,对应 __aenter__/__aexit____aiter__/__anext__async yield

资源(数据库连接、HTTP 会话)的获取和释放如果本身是异步操作,用 async with

1
2
3
4
5
6
7
8
9
10
11
class AsyncDB:
async def __aenter__(self):
self.conn = await connect()
return self
async def __aexit__(self, exc_type, exc, tb):
await self.conn.close()
return False

async def main():
async with AsyncDB() as db:
await db.query("SELECT 1")

contextlib.asynccontextmanager 和同步版 contextmanager 用法一致,yield 前后可以是 await

1
2
3
4
5
6
7
8
9
from contextlib import asynccontextmanager

@asynccontextmanager
async def db_session():
conn = await connect()
try:
yield conn
finally:
await conn.close()

数据分批异步到达时(游标分页、消息流)用异步迭代器:

1
2
3
4
5
6
7
8
9
10
11
class AsyncLines:
def __aiter__(self): return self
async def __anext__(self):
line = await self.reader.readline()
if not line:
raise StopAsyncIteration
return line.decode().rstrip()

async def main():
async for line in AsyncLines(reader):
print(line)

async def + yield 就是异步生成器,用 async for 消费,适合「边产生边消费」的流式场景,内存占用恒定:

1
2
3
4
5
6
7
async def poll_status(url: str):
while True:
status = await fetch_status(url)
yield status
if status == "done":
return
await asyncio.sleep(1)

同步原语:异步版的锁、信号量与队列

单线程协程模型下,协程之间「同时」访问共享状态时仍可能交错(一个协程 await 让出后,另一个改了同一个变量),因此也需要同步原语。它们的接口与 threading 模块对应,但只用于协程之间,不能跨线程当锁用。

1
2
3
4
5
6
7
8
9
lock = asyncio.Lock()
async def critical_section():
async with lock: # 同一时刻只有一个协程能进入
await update_shared_state()

sem = asyncio.Semaphore(10) # 最多 10 个并发
async def fetch_limited(url: str):
async with sem:
return await fetch(url)

Semaphore 是爬虫和 API 客户端限流的惯用法,配合 TaskGroup 用,能避免一次 gather 几千个任务打爆对方。注意:在没有 await 的纯计算里协程根本不会切换,这种情况下加锁是多余的–锁保护的是「跨 await 的不变量」。

asyncio.Queueput/get 都是协程方法,队列满或空时会挂起等待:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
async def producer(q: asyncio.Queue):
for i in range(5):
await q.put(i)
await q.put(None) # 哨兵

async def consumer(q: asyncio.Queue):
while True:
item = await q.get()
if item is None:
break
print(f"处理 {item}")
q.task_done()

async def main():
q = asyncio.Queue(maxsize=3)
await asyncio.gather(producer(q), consumer(q))

task_done() 配合 q.join() 可以等到所有 put 进去的任务都被标记完成。此外 asyncio.Event(事件通知,wait()set())、asyncio.Condition(带通知的锁)按需取用。这些原语的存在说明一个事实:单线程不等于没有并发问题

桥接同步代码:to_threadrun_in_executor

异步程序的尴尬在于:很多库是同步的。在协程里直接 time.sleep(1)requests.get(...) 会卡住整个事件循环–所有其他协程都跟着停。解决办法是把阻塞调用扔到线程池里。

1
2
3
4
5
6
7
8
9
import asyncio, time

def blocking_io():
time.sleep(1) # 同步阻塞
return "done"

async def main():
# 不阻塞事件循环,阻塞发生在独立线程里
result = await asyncio.to_thread(blocking_io)

asyncio.to_thread(3.9+)内部用默认 ThreadPoolExecutor。它适合 CPU 极轻但会阻塞的 IO 调用(老旧的同步 HTTP 库、文件 IO、调用 C 扩展)。如果阻塞函数本身有 GIL 竞争(CPU 密集),线程池帮不了你,该用进程。

需要自定义线程池大小或用进程池时,用更底层的 loop.run_in_executor

1
2
3
4
5
6
from concurrent.futures import ThreadPoolExecutor

async def main():
loop = asyncio.get_running_loop()
with ThreadPoolExecutor(max_workers=4) as pool:
result = await loop.run_in_executor(pool, blocking_io)

to_thread 就是 run_in_executor(None, func) 的语法糖。把一个同步阻塞函数包装成协程函数是连接新旧代码的常见模式:

1
2
3
def sync_func(x): ...   # 某个阻塞的第三方库
async def async_func(x):
return await asyncio.to_thread(sync_func, x) # 不阻塞事件循环,但没让它变快

实战:并发抓取 + 限流 + 超时

把前面几节拼成一个真实骨架:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
import asyncio

sem = asyncio.Semaphore(5) # 全局并发上限

async def fetch_one(client, url: str):
async with sem: # 限流
try:
async with asyncio.timeout(3.0): # 单个超时
resp = await client.get(url)
return await resp.json()
except TimeoutError:
print(f"超时: {url}")
return None

async def fetch_all(client, urls: list[str]):
async with asyncio.TaskGroup() as tg: # 异常时自动取消全部
tasks = [tg.create_task(fetch_one(client, u)) for u in urls]
return [t.result() for t in tasks if t.result() is not None]

async def main():
urls = [f"https://api.example.com/{i}" for i in range(50)]
results = await fetch_all(client, urls) # client 用 aiohttp/httpx 异步客户端
print(f"成功 {len(results)} 条")

asyncio.run(main())

几个关键点:Semaphore 限流、timeout 防慢请求拖累整体、TaskGroup 保证异常时及时取消、全程用异步 HTTP 客户端而不是 requests

陷阱

阻塞了事件循环–最常见的错误。协程里调 time.sleep(1)requests.get(url)json.load(open("big.json")) 都会卡住整个循环。规则:协程里只能用异步 IO 库,没有异步替代品时用 to_thread 包一层。开启 asyncio.get_event_loop().set_debug(True) 会打印耗时超过 100ms 的回调,便于定位。

协程没有被调度–漏写 await 时协程不执行,只留下一个没被消费的协程对象和 RuntimeWarning: coroutine was never awaited。看到这个警告就去找哪里漏了 await

在异步代码里用同步原语threading.Lockacquire() 是阻塞调用,在协程里用会卡住事件循环。协程间互斥要用 asyncio.Lock。反过来,asyncio.Lock 也不能跨线程用。

gather 不取消其他任务–默认 return_exceptions=False 时某个任务抛异常,其余任务仍在后台跑。不希望这样就用 TaskGroup(自动取消),或手动 for t in tasks: t.cancel()

取消被吞掉–在协程里 except BaseException: 又不 raise,取消信号被吃掉,任务变成「取消不了」的僵尸。清理后必须重新抛出 CancelledError

混用事件循环实现uvloop 在 Linux 上能显著提升性能,但不要在代码里硬编码切换,应集中在程序入口处 uvloop.install()(包在 try/except ImportError 里)。

何时该用 asyncio,何时不该

维度适合 asyncio不太适合
任务类型IO 密集(网络、数据库、文件)CPU 密集(数值计算、压缩、加密)
并发量大量连接(成百上千)少量任务(几个)
依赖库有成熟的异步客户端只能用同步阻塞库
团队熟悉协程心智模型异步经验有限,追求直观

决策要点:

  • 瓶颈是 IO 等待 -> asyncio 收益最大,一个进程撑几万连接是常态。
  • 瓶颈是 CPU 计算 -> 用 multiprocessing 或 C 扩展,asyncio 帮不了你(GIL 下单线程协程不会让计算变快)。
  • 需要混合 -> asyncio 做主调度,CPU 密集任务用 run_in_executor 扔到进程池。
  • 只是写个小脚本 -> 直接同步代码更简单,别为了异步而异步。

一个常见误判:以为「异步 = 更快」。对单个请求,异步代码因多了协程调度开销往往比同步还慢一点。异步的优势在并发吞吐–同时处理大量等待型任务时,单线程协程的内存和调度成本远低于多线程。

总结

asyncio 的核心是「单线程事件循环 + 协作式调度」。全部概念可以压缩成几句话:async def 定义协程,await 暂停并让出控制权,asyncio.run 启动循环;要并发就把协程变成 Task,用 TaskGroup/gather 等待;超时用 asyncio.timeout,取消通过 CancelledError 传播,清理后必须重新抛出;协程间互斥用 asyncio.Lock/Semaphore/Queue,它们是 threading 的协程版,不能混用;阻塞调用必须扔进线程池或换成异步库,否则卡死事件循环。

掌握 asyncio 的真正门槛,是建立「哪里会让出控制权」的直觉。每一处 await 都是一个潜在的切换点,所有共享状态的正确性都建立在对这些切换点的预判上。最后记住一条铁律:异步代码的价值在并发 IO,不在单个任务的加速。如果你的程序既不并发也不等待,写异步只会增加复杂度而没有回报。