本文我们将从生成器的底层原理出发,一步步剖析协程的演进过程,最后深入 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+ 中,我们有了三种不同的协程概念:

  1. 生成器协程:基于 yield 的传统方式,兼容旧代码
  2. 原生协程:用 async def 定义的函数,返回 coroutine 对象
  3. 异步生成器:用 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.Lockasyncio.Lock
获取方式阻塞调用 lock.acquire()可等待 await lock.acquire()
GIL 行为持有 GILawait 时释放 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 建立完善的监控体系

  • 监控事件循环延迟
  • 监控协程数量和状态
  • 监控任务队列长度
  • 设置告警阈值,及时发现问题

五、互动与思考

问大家几个问题,欢迎在评论区讨论:

  1. 你在项目中是如何平衡异步和同步代码的? 是全部异步化,还是只在关键路径使用异步?

  2. 当需要调用一个没有异步版本的第三方库时,你会怎么做? 用线程池包装,还是寻找替代方案?

  3. 你觉得异步编程最大的挑战是什么? 是心智模型的转变,还是调试困难,或是团队协作问题?

  4. 有没有遇到过一些特别的异步编程陷阱? 欢迎分享你的踩坑经历,让更多人受益。

六、总结

Python 协程从生成器一路走来,经历了 yield fromasync/await 的演进,最终在 asyncio 中形成了完整的异步编程生态。理解这个演进过程,不仅能帮你更好地使用现有工具,还能让你在面对未来变化时更加从容。

异步编程适合 I/O 密集型场景,但在 CPU 密集型任务中可能带来额外复杂度。选择合适的工具,遵循最佳实践,异步编程才能真正成为提升系统性能的利器。

最后送大家一句话: 异步编程的核心不是技术复杂度,而是思维模式的转变。一旦你理解了“让出控制权”的精髓,你会发现异步世界其实并没有那么可怕。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐