前言
asyncio是 Python 标准库里的异步框架,但它不是「一个函数」,而是一整套协作式并发的基础设施。新手常犯的错,是把asyncio的 API 当成threading的等价物来用——随手await两下,却发现根本没有并发。
这篇是API 速查 + 语义说明:先讲清「协作式调度」到底意味着什么,再把最常用的几个入口——create_task、gather、wait_for、as_completed、Queue、Semaphore、Lock、to_thread——逐个说清用途和版本。
先记一条贯穿全文的语义:协程只有在await时才会把控制权交还给事件循环。你写的任何一段不含await的同步代码,都会独占循环,其他所有任务一起等它。理解这一点,下面每个 API 的取舍就都顺了。
一、协作式调度意味着什么
「协作式」是相对「抢占式」说的。线程由操作系统抢占式调度——时间片到了就强制切换,代码管不着。协程反过来:必须自己在await处让位。这带来两个实践结果:
- 好处:切换点完全由你掌控,不用担心「改到一半被切走」,所以共享数据不需要锁——只要两个协程之间没有
await,它们就不会交错。 - 坏处:一个协程如果长时间不
await,事件循环就卡住。阻塞调用(time.sleep、同步 I/O、大循环)是异步程序的头号杀手。
# 适用于 Python 3.8+
import asyncio
async def ticker(name, n):
for i in range(n):
print(name, i)
await asyncio.sleep(0) # await sleep(0) 主动让位一次,让别的任务有机会跑
async def main():
await asyncio.gather(ticker("A", 3), ticker("B", 3))
asyncio.run(main())asyncio.sleep(0)是「主动让出一次控制权」的常用技巧,它不会真的睡,只是把执行权还给循环。
二、create_task 与 gather:把任务跑起来
asyncio.create_task(coro)(3.7 新增)把协程包装成一个Task并立刻交给事件循环排队。关键区别在于:
# 适用于 Python 3.8+
import asyncio
async def work(n):
await asyncio.sleep(n)
return n * 2
async def serial():
a = await work(1) # 串行:等完 1 秒再走下一行
b = await work(1) # 再等 1 秒,总共约 2 秒
return a, b
async def concurrent():
t1 = asyncio.create_task(work(1)) # 立刻排队,不等它
t2 = asyncio.create_task(work(1)) # 立刻排队
return await asyncio.gather(t1, t2) # 一起等,总耗时约 1 秒asyncio.gather(*aws, return_exceptions=False)用来等一组任务全部完成,并按传入顺序返回结果列表。return_exceptions=True时,某个任务抛出的异常会当成结果放进列表,而不是让gather整体失败。
需要说明一个易混点:gather里如果某个任务抛异常,默认情况下其他任务不会被自动取消,它们会继续跑;而 3.11 新增的TaskGroup行为不同——组里任一任务失败,TaskGroup会把其余任务取消。两者不是简单替代关系。
三、wait_for 与 as_completed:控制等待方式
await asyncio.wait_for(aw, timeout)给单个可等待对象加超时。超时后它会取消该任务并抛asyncio.TimeoutError。
# 适用于 Python 3.8+
import asyncio
async def slow():
await asyncio.sleep(10)
async def main():
try:
await asyncio.wait_for(slow(), timeout=1.0)
except asyncio.TimeoutError:
print("超时了")asyncio.as_completed(aws, *, timeout=None)则是「谁先完成就先处理谁」,返回一个可迭代对象。注意它不会取消未完成的任务,只是改变你获取结果的顺序,适合「先拿到先处理」的场景,比如分批写文件。
# 适用于 Python 3.8+(3.13 起也可用 async for 迭代)
import asyncio
async def job(i):
await asyncio.sleep(0.1 * i)
return i
async def main():
tasks = [asyncio.create_task(job(i)) for i in range(5)]
for fut in asyncio.as_completed(tasks):
print(await fut) # 按完成先后打印,不是按 0,1,2,3 的顺序在 3.11 及以上,如果想在协程里用async with asyncio.timeout(seconds):这种上下文管理器写法,也是官方推荐的方式,它比wait_for更贴近结构化并发。
四、Queue、Semaphore、Lock:协作式下的同步原语
标准库asyncio提供了自己的一套同步原语,不能用queue.Queue或threading.Lock替代——那些是给线程用的阻塞原语,放进协程里会卡死循环。
asyncio.Queue(maxsize=0):maxsize <= 0表示不限;await q.put(x)在队满时挂起,await q.get()在队空时挂起。生产者消费者模型标配。配套还有join()与task_done()。asyncio.Semaphore(value=1):信号量,用来限制同时进行的任务数量。这是控制并发上限最常用的工具。asyncio.Lock:互斥锁,用async with lock:使用。在纯协程代码里如果临界区没有await,其实不需要锁;但一旦临界区里有await,就可能被其他协程穿插,此时需要锁。
# 适用于 Python 3.8+
import asyncio
async def main():
sem = asyncio.Semaphore(3) # 最多 3 个任务同时进行
queue = asyncio.Queue()
async def worker(n):
async with sem: # 超过 3 个就排队等信号量
await asyncio.sleep(0.1)
await queue.put(n * n)
await asyncio.gather(*(worker(i) for i in range(10)))
results = []
while not queue.empty():
results.append(queue.get_nowait())
print(len(results))信号量的意义就在「把并发上限压到一个合规的值」——这一点在后文讲抓取时会再强调。
五、to_thread:把同步代码搬离事件循环
asyncio.to_thread(func, /, *args, **kwargs)是3.9 新增的。它把同步函数丢到一个线程里执行,不阻塞事件循环,返回一个可 await 的对象。
# 适用于 Python 3.9+(3.8 请用 loop.run_in_executor)
import asyncio, time
def blocking_io():
time.sleep(1)
return "done"
async def main():
result = await asyncio.to_thread(blocking_io)
print(result)
asyncio.run(main())官方文档特别提醒:由于 GIL 的存在,to_thread()通常只适合 I/O 密集的函数;只有在没有 GIL 的构建上(比如自由线程构建),它才能用于 CPU 密集函数。所以别指望用to_thread()加速纯计算——那还是要靠进程池。
3.8 及更早版本没有to_thread(),等价写法是:
# 适用于 Python 3.8(to_thread 出现之前)
import asyncio
async def main():
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(None, blocking_io)
print(result)run_in_executor的第一个参数传None表示使用默认的ThreadPoolExecutor。
常见坑点
- 连写两个 await 以为是并发
❌a = await f(); b = await g()(串行)
✅a, b = await asyncio.gather(f(), g())
- 用
time.sleep做异步等待
❌async def f(): time.sleep(1)(卡死循环)
✅await asyncio.sleep(1)
- 在协程里用线程版同步原语
❌import queue; queue.Queue().get()(阻塞事件循环)
✅await asyncio.Queue().get()
- 创建了 Task 却没保存引用
❌asyncio.create_task(f())后不持有返回值,任务可能被垃圾回收提前消失
✅ 保存到变量或放进集合:t = asyncio.create_task(f())
- 以为
to_thread能加速 CPU 计算
❌ 把密集计算丢给to_thread期待提速
✅ 认清 GIL,CPU 密集改用ProcessPoolExecutor
- 把
wait_for当成「给任务加超时」
❌ 以为超时后子任务自动结束了
✅ 超时会取消等待,必要时在except里自行清理
- 在 3.11+ 混用 gather 与 TaskGroup 的异常语义
❌ 以为TaskGroup和gather失败后行为一样
✅ 记住TaskGroup会连带取消兄弟任务,gather默认不会
总结
| API | 用途 | 版本 |
|---|
create_task | 把协程立即排入事件循环 | 3.7+ |
gather | 并发跑一组并收齐结果 | 3.5+ |
wait_for | 给单个等待加超时 | 3.5+ |
as_completed | 按完成先后取结果 | 3.5+ |
Queue/Semaphore/Lock | 协作式同步原语 | 3.4+ |
to_thread | 把同步函数挪到线程 | 3.9+ |
TaskGroup | 结构化并发任务组 | 3.11+ |
asyncio的入门要点只有两条:并发要先create_task再统一等,等待必须用await而不是阻塞调用。把这两条装进肌肉记忆,剩下的 API 都是在这个框架上做加减法。