Python协程与异步IO深入理解——从生成器到asyncio
本文我们将从生成器的底层原理出发,一步步剖析协程的演进过程,最后深入 asyncio 的事件循环机制。文章包含大量可运行的代码示例和我在实际项目中遇到的坑,希望能帮你真正理解异步编程的精髓。
大家好,我是小陈工,一个有着 9 年 Python 后端开发经验的老程序员。今天我们来聊聊 Python 异步编程中那个让人又爱又恨的话题——协程。
你可能已经用过 asyncio,写过 async/await,但你是否真的理解协程是怎么来的?为什么一个看似简单的语法背后,藏着如此复杂的调度机制?今天这篇文章,我将从最底层的生成器开始,带你走过协程演进的每一个关键节点,最后深入 asyncio 的核心设计。
先问大家一个问题: 你在使用 asyncio 过程中,最常遇到的性能瓶颈是什么?是数据库查询太慢?还是网络请求太多?欢迎在评论区分享你的踩坑经历,我们一起探讨解决方案。
一、从生成器到协程:一场静默的革命
如果你用过 Python 的 yield 关键字,那你已经接触过协程的“祖先”了。在 Python 2.5 引入生成器表达式时,没人想到它会是未来异步编程的基石。
1.1 生成器的本质:可暂停的函数
先看一个最简单的生成器:
def simple_generator():
print("开始")
yield 1
print("继续")
yield 2
print("结束")
gen = simple_generator()
print(next(gen)) # 输出:开始\n1
print(next(gen)) # 输出:继续\n2
# print(next(gen)) # 这里会抛出 StopIteration
关键点: 生成器函数在 yield 处暂停,下次调用 next() 时从暂停处恢复。这种“暂停-恢复”的能力,正是协程的核心特征。
不过,早期的生成器只能往外发送数据(yield),不能接收外部传入的数据。直到 Python 2.5 引入了 .send() 方法,生成器才真正变成了双向通道:
def coroutine_v1():
print("启动")
while True:
x = yield # 接收外部传入的值
print(f"接收到: {x}")
coro = coroutine_v1()
next(coro) # 启动生成器,执行到第一个 yield
coro.send(10) # 输出:接收到: 10
coro.send("hello") # 输出:接收到: hello
这时候,生成器开始有了“协程”的雏形——它可以在暂停时接收外部数据,处理后再继续执行。
1.2 yield from:协程进化的关键一步
Python 3.3 引入的 yield from 语法,彻底改变了游戏规则。它允许一个生成器委托另一个生成器执行:
def generator_a():
yield from [1, 2, 3]
yield from generator_b()
def generator_b():
yield "a"
yield "b"
for item in generator_a():
print(item)
# 输出:1 2 3 a b
看起来平平无奇,对吧?但 yield from 的真正威力在于它建立了调用链和异常传播机制。这让“生成器嵌套生成器”变得可行,为后来的异步编程框架打下了基础。
我的个人经验: 我曾在一次系统重构中,用 yield from 实现了一个简单的任务调度器。当时我们有一个需要按顺序调用多个第三方 API 的需求,每个 API 都有复杂的依赖关系。用 yield from 构建的生成器链,让代码的可读性比回调地狱提升了不止一个档次。
二、async/await:Python 的“官方”协程
Python 3.5 引入的 async/await 语法,终于让协程成为了语言的一等公民。这两个关键字带来了全新的协程类型和更清晰的语义。
2.1 三种协程类型
在 Python 3.5+ 中,我们有了三种不同的协程概念:
- 生成器协程:基于
yield的传统方式,兼容旧代码 - 原生协程:用
async def定义的函数,返回coroutine对象 - 异步生成器:用
async def定义且包含yield的函数,Python 3.6+
import types
# 1. 生成器协程
def gen_coro():
yield from range(3)
# 2. 原生协程
async def native_coro():
await asyncio.sleep(1)
return "done"
# 3. 异步生成器
async def async_gen():
yield 1
await asyncio.sleep(0.1)
yield 2
print(isinstance(gen_coro(), types.GeneratorType)) # True
print(isinstance(native_coro(), types.CoroutineType)) # True
print(isinstance(async_gen(), types.AsyncGeneratorType)) # True
重要区别: 原生协程不能用 yield,也不能用 yield from;异步生成器不能用 return 返回值。这些限制看似严格,实际上让代码的意图更加清晰。
2.2 await 的工作原理
await 关键字可能是异步编程中最容易被误解的部分。很多人以为 await 就是“等待”,但实际上它是“让出控制权”。
async def task_a():
print("A: 开始")
await asyncio.sleep(2) # 让出控制权,事件循环可以去执行其他任务
print("A: 结束")
return "A结果"
async def task_b():
print("B: 开始")
await asyncio.sleep(1) # 让出控制权
print("B: 结束")
return "B结果"
async def main():
# 错误写法:串行执行
# result_a = await task_a() # 等2秒
# result_b = await task_b() # 再等1秒,总耗时3秒
# 正确写法:并发执行
task1 = asyncio.create_task(task_a())
task2 = asyncio.create_task(task_b())
result_a = await task1
result_b = await task2
# 总耗时约2秒(取最长任务的时间)
print(f"结果: {result_a}, {result_b}")
asyncio.run(main())
关键理解: await 并不是“阻塞等待”,而是告诉事件循环:“我这里需要等待某个结果,你先去执行其他准备好的任务,等我的结果准备好了再来通知我。”
三、实战踩坑:我在生产环境遇到的 5 个典型问题
理论知识说再多,不如实际踩几个坑来得深刻。下面是我在 9 年开发中遇到的 5 个典型异步编程问题,每一个都曾让我们的系统付出过代价。
3.1 坑一:同步库混入异步世界,导致整个事件循环冻结
场景还原: 我们有一个基于 FastAPI 的微服务,需要调用第三方 API。当时团队新来的同事为了图方便,直接在异步路由中使用了 requests.get()。
# ❌ 错误示范:在异步函数中使用同步阻塞调用
from fastapi import FastAPI
import requests # 同步库!
app = FastAPI()
@app.get("/api/data")
async def get_data():
# 这个调用会阻塞整个事件循环!
response = requests.get("https://api.example.com/data")
return response.json()
问题现象: 当并发请求量稍大时(比如 100 QPS),整个服务开始响应缓慢,甚至完全无响应。监控显示 CPU 使用率很低,但请求排队时间越来越长。
根本原因: requests.get() 是同步阻塞调用,它会独占事件循环所在的线程。在它执行期间,事件循环无法调度其他协程,导致所有并发请求都被“冻住”。
解决方案:
# ✅ 正确做法1:使用异步HTTP客户端
import httpx
@app.get("/api/data")
async def get_data():
async with httpx.AsyncClient() as client:
response = await client.get("https://api.example.com/data")
return response.json()
# ✅ 正确做法2:如果必须用同步库,放到线程池中执行
import asyncio
from concurrent.futures import ThreadPoolExecutor
thread_pool = ThreadPoolExecutor(max_workers=10)
@app.get("/api/data")
async def get_data():
loop = asyncio.get_running_loop()
# 将同步调用转移到线程池,不阻塞事件循环
response = await loop.run_in_executor(
thread_pool,
requests.get,
"https://api.example.com/data"
)
return response.json()
我的反思: 这件事让我意识到,团队的技术规范多么重要。我们后来建立了明确的编码规范——在异步上下文中,禁止直接使用任何同步阻塞库。所有新引入的依赖都必须经过异步兼容性审查。
3.2 坑二:Task 对象管理不当,导致内存泄漏
场景还原: 我们需要定期从消息队列拉取任务并处理。最初的实现中,我们为每个消息创建一个 Task,但没保留引用,也没处理异常。
# ❌ 危险写法:创建任务后丢失引用
async def process_message(message):
# 模拟耗时处理
await asyncio.sleep(5)
print(f"处理完成: {message}")
async def message_consumer():
while True:
message = await queue.get()
# 危险!创建任务后立即丢弃引用
asyncio.create_task(process_message(message))
问题现象: 服务运行几天后,内存使用量持续增长。即使消息处理完成,内存也不释放。最终导致 OOM(Out Of Memory)崩溃。
根本原因: asyncio.Task 对象持有协程的强引用。Python 3.8 之前,如果没有显式保留 Task 的引用,垃圾回收时可能会出现问题。即使任务执行完成,如果存在未捕获的异常,Task 对象可能不会被正确清理。
解决方案:
# ✅ 正确做法1:使用 TaskGroup(Python 3.11+)
async def message_consumer():
async with asyncio.TaskGroup() as tg:
while True:
message = await queue.get()
tg.create_task(process_message(message))
# ✅ 正确做法2:显式管理任务生命周期
async def message_consumer():
tasks = set()
while True:
message = await queue.get()
task = asyncio.create_task(process_message(message))
tasks.add(task)
task.add_done_callback(tasks.discard) # 任务完成后自动移除引用
# ✅ 正确做法3:控制并发数量,避免创建过多任务
async def message_consumer():
semaphore = asyncio.Semaphore(50) # 最多同时处理50个消息
while True:
message = await queue.get()
async with semaphore:
await process_message(message)
我的经验: 异步编程中的资源管理比同步编程更加复杂。我们后来在项目中强制使用 TaskGroup(需要 Python 3.11+),对于旧版本则使用封装好的任务管理器。同时,我们建立了内存监控告警机制,一旦发现内存异常增长就立即介入排查。
3.3 坑三:异步上下文管理器使用不当,引发死锁
场景还原: 在数据库事务管理中,我们需要确保连接的正确关闭。最初我们混用了同步锁和异步代码。
# ❌ 危险组合:同步锁 + async代码
import threading
lock = threading.Lock()
async def unsafe_transaction():
with lock: # 阻塞式获取锁!
# 假设这里有一些数据库操作
await asyncio.sleep(1) # ⚠️ 在持有同步锁期间切换协程
# 这可能导致死锁或其他协程无法获取锁
问题现象: 在高并发场景下,偶尔会出现某些请求永远卡住,不报错也不完成。重启服务后暂时恢复正常,但一段时间后问题复现。
根本原因: Python 的 GIL(全局解释器锁)在 await 时会被释放,但同步锁(threading.Lock)不会。当一个协程持有同步锁然后 await 时,其他协程可能被调度执行,如果它们也需要同一个锁,就会产生死锁。
解决方案:
# ✅ 正确做法:使用异步锁
from asyncio import Lock
lock = Lock()
async def safe_transaction():
async with lock: # 异步获取锁,会在 await 时正确释放控制权
await asyncio.sleep(1)
# 其他协程可以在此期间执行,但无法获取同一个锁
关键区别对比:
| 特性 | threading.Lock | asyncio.Lock |
|---|---|---|
| 获取方式 | 阻塞调用 lock.acquire() | 可等待 await lock.acquire() |
| GIL 行为 | 持有 GIL | 在 await 时释放 GIL |
| 协程安全 | 不安全,可能死锁 | 安全,专为协程设计 |
| 使用场景 | 线程同步 | 协程同步 |
3.4 坑四:CPU 密集型操作导致事件循环冻结
场景还原: 我们需要在异步服务中对大量数据进行计算处理。开发同学直接在主协程中写了计算逻辑。
# ❌ 错误:在协程中执行CPU密集型计算
async def process_batch(data_list):
results = []
for data in data_list:
# 假设这是复杂的计算,没有await点
result = expensive_computation(data) # 纯CPU计算
results.append(result)
return results
问题现象: 当处理大数据批次时,整个服务完全无响应,连简单的健康检查请求都超时。计算结束后才恢复正常。
根本原因: asyncio 本质是协作式多任务,每个 await 都是调度点。纯 CPU 运算不会自动让出控制权,长时间占用事件循环线程,导致其他协程饥饿。
解决方案:
# ✅ 正确做法1:分块处理,定期插入 await asyncio.sleep(0)
async def process_batch(data_list):
results = []
for i, data in enumerate(data_list):
result = expensive_computation(data)
results.append(result)
# 每处理100个数据让出控制权一次
if i % 100 == 0:
await asyncio.sleep(0) # 让事件循环有机会调度其他任务
return results
# ✅ 正确做法2:将CPU密集型任务转移到进程池
from concurrent.futures import ProcessPoolExecutor
process_pool = ProcessPoolExecutor(max_workers=4)
async def process_batch(data_list):
loop = asyncio.get_running_loop()
# 将计算转移到进程池
results = await loop.run_in_executor(
process_pool,
batch_compute, # 这个函数需要在模块顶层定义
data_list
)
return results
# ✅ 正确做法3:使用 asyncio.to_thread(Python 3.9+)
async def process_batch(data_list):
# 将CPU密集型任务转移到单独的线程
results = await asyncio.to_thread(batch_compute, data_list)
return results
我的建议: 对于 CPU 密集型任务,最好的选择是进程池。线程池虽然也可以,但 Python 的 GIL 会限制多核性能。如果计算量不大,使用 asyncio.to_thread() 或定期 await asyncio.sleep(0) 是不错的折中方案。
3.5 坑五:CancelledError 处理缺失,导致资源泄漏
场景还原: 我们在处理数据库连接时,没有正确处理取消异常。
# ❌ 危险:cleanup 操作可能永远不会执行
async def write_to_db():
conn = await acquire_connection()
try:
await conn.execute("INSERT ...")
finally:
await conn.close() # ⚠️如果任务被取消,这行代码可能永远不会执行!
问题现象: 长时间运行后,数据库连接池被耗尽,新请求无法获取连接。重启服务后连接释放,但问题会再次出现。
根本原因: asyncio.CancelledError 继承自 BaseException 而非 Exception,常规的 try/except Exception 无法捕获它。如果任务在 await conn.execute() 处被取消,会直接抛出 CancelledError,跳过 finally 块中的清理代码。
# ✅ 健壮写法:正确处理 CancelledError
async def robust_db_op():
conn = None
try:
conn = await acquire_connection()
await conn.execute("...")
except (Exception, asyncio.CancelledError):
# 捕获所有异常,包括取消异常
if conn:
try:
await conn.close()
except Exception:
log.error("清理失败")
raise # 重要:CancelledError 必须重新抛出
# ✅ 或者使用 asyncio.shield() 保护关键清理操作
async def shielded_db_op():
conn = await acquire_connection()
try:
await conn.execute("...")
finally:
# 使用 shield 防止取消中断清理
await asyncio.shield(conn.close())
关键点: CancelledError 必须重新抛出,否则会破坏整个取消机制。清理操作自身也需要异常处理,避免一个异常掩盖另一个异常。
四、最佳实践:来自 9 年经验的 5 条建议
基于以上踩坑经历,我总结了 5 条 Python 异步编程的最佳实践:
4.1 明确职责边界:异步 vs 同步
- 异步部分:负责 I/O 密集型操作(网络请求、数据库查询、文件读写)
- 同步部分:负责 CPU 密集型计算、简单逻辑处理
- 绝对不要在异步函数中直接调用同步阻塞库
4.2 严格控制并发数量
# 使用 Semaphore 限制并发数
async def limited_concurrent(tasks, max_concurrent=100):
semaphore = asyncio.Semaphore(max_concurrent)
async def sem_task(task):
async with semaphore:
return await task
return await asyncio.gather(*[sem_task(t) for t in tasks])
4.3 正确处理异常和取消
- 使用
try/except (Exception, asyncio.CancelledError)捕获所有异常 - 重要清理操作使用
asyncio.shield()保护 - 取消异常必须重新抛出
4.4 合理设置超时
async def safe_request():
try:
# 设置 10 秒超时
response = await asyncio.wait_for(
http_client.get(url),
timeout=10.0
)
return response
except asyncio.TimeoutError:
log.warning("请求超时")
return None
4.5 建立完善的监控体系
- 监控事件循环延迟
- 监控协程数量和状态
- 监控任务队列长度
- 设置告警阈值,及时发现问题
五、互动与思考
问大家几个问题,欢迎在评论区讨论:
-
你在项目中是如何平衡异步和同步代码的? 是全部异步化,还是只在关键路径使用异步?
-
当需要调用一个没有异步版本的第三方库时,你会怎么做? 用线程池包装,还是寻找替代方案?
-
你觉得异步编程最大的挑战是什么? 是心智模型的转变,还是调试困难,或是团队协作问题?
-
有没有遇到过一些特别的异步编程陷阱? 欢迎分享你的踩坑经历,让更多人受益。
六、总结
Python 协程从生成器一路走来,经历了 yield from 到 async/await 的演进,最终在 asyncio 中形成了完整的异步编程生态。理解这个演进过程,不仅能帮你更好地使用现有工具,还能让你在面对未来变化时更加从容。
异步编程适合 I/O 密集型场景,但在 CPU 密集型任务中可能带来额外复杂度。选择合适的工具,遵循最佳实践,异步编程才能真正成为提升系统性能的利器。
最后送大家一句话: 异步编程的核心不是技术复杂度,而是思维模式的转变。一旦你理解了“让出控制权”的精髓,你会发现异步世界其实并没有那么可怕。
更多推荐


所有评论(0)