Python 异步编程:基础、进阶与常见坑点
Python 异步编程:基础、进阶与常见坑点
异步编程(asyncio)是 Python 处理高并发 IO 场景的核心工具,心智模型和同步代码差异很大。本文按"基础用法 → 进阶用法 → 常见坑点"的顺序梳理,帮助系统性掌握。
一、基础使用
1. 核心概念:协程不是线程
asyncio 是单线程的协作式并发:所有协程跑在同一个线程里,通过事件循环(event loop)在"等待 IO"的间隙切换执行其他任务,而不是像多线程那样由操作系统抢占式调度。它适合 IO 密集型任务(网络请求、文件读写、数据库查询),不适合 CPU 密集型任务(没有真正并行计算能力)。
2. 协程函数与 await
import asyncio
async def fetch(name):
await asyncio.sleep(1) # 模拟IO等待,非阻塞
return f"{name} done"
async def main():
result = await fetch("task1")
print(result)
asyncio.run(main()) # 程序入口,负责创建并运行事件循环async def 定义协程函数,调用它不会立即执行,而是返回一个协程对象,必须用 await(或交给事件循环调度)才会真正运行。asyncio.run() 是最常见的同步代码入口,负责创建事件循环、运行协程、关闭循环。
3. 顺序 await 仍然是串行的
async def main():
await fetch("task1") # 等待1秒
await fetch("task2") # 再等待1秒,总耗时2秒,没有并发只用 await 顺序调用并不会带来并发收益,真正的并发需要进阶部分的 gather/create_task。
二、进阶使用
1. 并发执行多个协程:gather
async def main():
results = await asyncio.gather(
fetch("task1"), fetch("task2"), fetch("task3")
)
print(results) # 三个任务并发执行,总耗时约等于最慢的那个2. create_task:手动调度并保留任务引用
async def main():
t1 = asyncio.create_task(fetch("task1"))
t2 = asyncio.create_task(fetch("task2"))
# 任务创建后立即开始调度执行(不需要立刻await)
r1 = await t1
r2 = await t2相比 gather,create_task 适合需要更精细控制任务生命周期(比如提前取消、单独处理异常)的场景。
3. 超时控制与任务取消
try:
result = await asyncio.wait_for(fetch("slow"), timeout=3)
except asyncio.TimeoutError:
print("超时了")
task = asyncio.create_task(fetch("task1"))
task.cancel() # 主动取消任务4. 限制并发数:Semaphore
批量发起请求时,往往需要限制同时进行的并发数量,避免压垮目标服务:
sem = asyncio.Semaphore(5) # 最多5个并发
async def limited_fetch(name):
async with sem:
return await fetch(name)
await asyncio.gather(*(limited_fetch(f"t{i}") for i in range(100)))5. 生产者-消费者模式:asyncio.Queue
async def producer(q):
for i in range(5):
await q.put(i)
await q.put(None) # 哨兵值,通知消费者结束
async def consumer(q):
while True:
item = await q.get()
if item is None:
break
print("处理", item)
async def main():
q = asyncio.Queue()
await asyncio.gather(producer(q), consumer(q))6. 混合调用阻塞代码:run_in_executor
遇到没有异步版本的阻塞库(比如某些数据库驱动、CPU 密集计算),应该丢进线程池/进程池执行,避免卡住事件循环:
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(None, blocking_function, arg1)7. 异步上下文管理器与异步迭代器
常见于数据库连接、HTTP 会话等资源管理,使用 async with 和 async for:
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
data = await resp.json()8. 协程内的异常处理与 return_exceptions
results = await asyncio.gather(*tasks, return_exceptions=True)
# 失败的任务对应位置是异常对象本身,不会中断整体执行三、常见坑点
1. 忘记 await,协程从未执行
async def main():
fetch("task1") # 错误:只创建了协程对象,从未执行
await fetch("task2") # 正确不会报错,只会收到 RuntimeWarning: coroutine was never awaited,是最隐蔽的坑之一。
2. 在协程里调用阻塞函数,会卡住整个事件循环
time.sleep()、同步版 requests.get()、普通文件 IO 等都是阻塞调用,会让整个事件循环停摆,其他协程全部无法执行。必须用 asyncio.sleep() 或 run_in_executor 替代。
3. create_task 创建的任务没有保留引用,可能被提前回收
事件循环只会弱引用 task,如果没有变量保存它,垃圾回收可能在任务还没跑完时就把它销毁:
tasks = set()
t = asyncio.create_task(fetch("task1"))
tasks.add(t)
t.add_done_callback(tasks.discard)4. gather 默认不会自动取消其他任务
某个协程抛异常时,gather 会立刻抛出,但其他协程仍在后台运行,容易造成"看似已处理异常,实际还有任务在跑"的状态,需要 return_exceptions=True 拿到完整结果。
5. 处理 CancelledError 时忘记重新抛出
async def worker():
try:
await asyncio.sleep(10)
except asyncio.CancelledError:
print("清理资源")
raise # 必须重新抛出,否则取消机制被破坏6. 事件循环不能嵌套运行
已经处于运行中的事件循环里(如 Jupyter Notebook)再调用 asyncio.run(),会报 RuntimeError: cannot be called from a running event loop。此时应直接 await,不要再包一层 asyncio.run()。
7. 协程专用同步原语不能跨线程使用
asyncio.Lock、asyncio.Event、asyncio.Queue 是为单个事件循环内的协程设计的,不能在多线程间共享(这点和 threading.Lock 不同)。跨线程通信需要 asyncio.run_coroutine_threadsafe()。
8. 把异步当成"自动并行"
asyncio 不能让 CPU 密集型代码变快,单线程依然受 GIL 限制;它解决的是 IO 等待期间的"空闲浪费",不是计算性能问题。CPU 密集任务应该用 ProcessPoolExecutor。
小结清单
- 基础:
async def+await是基本单元,asyncio.run()是入口,顺序await不会带来并发。 - 进阶:并发用
gather/create_task,限流用Semaphore,生产消费用Queue,混合阻塞代码用run_in_executor,资源管理用async with/async for。 - 坑点:忘记
await、协程里调用阻塞函数、任务引用丢失被回收、gather异常处理细节、CancelledError要重新抛出、事件循环不能嵌套、同步原语不能跨线程、误把异步当并行加速 CPU 任务。