Python——协程asyncio库的使用


整体说明

  • 协程的本质是函数 (具体来说是定义时被 async 修饰的函数)
  • asyncio 是 Python 内置的 异步 I/O 框架 ,核心用于编写高效的并发代码,尤其适合处理 I/O 密集型任务(如网络请求、文件读写、数据库操作等)
  • asyncio 基于 协程(coroutine) 实现,通过事件循环(Event Loop)调度任务,避免了多线程的上下文切换开销,效率更高
  • asyncio 在较高的 Python 版本中才可以使用,建议在 Python 3.7 以上使用
  • 常见应用场景包括
    • 网络请求:并发调用 API(如结合 aiohttp 库)
    • 文件操作:异步读写文件(如结合 aiofiles 库)
    • 数据库操作:异步操作数据库(如 asyncpg 用于 PostgreSQL,motor 用于 MongoDB)
    • WebSocket 服务:实现高并发的实时通信(如 websockets 库)
    • 定时任务:通过 asyncio.create_task() + 循环实现简单定时任务
  • 仅适用于 I/O 密集型任务:
    • asyncio 是单线程的,CPU 密集型任务会阻塞事件循环,需结合 loop.run_in_executor() 提交到线程池/进程池
      • 注:run_in_executor(executor, func, *args) 将同步函数 func 放到一个线程池(或进程池)中执行,返回一个 可等待的 Future
        • 如果第一个参数是 None,表示使用默认的线程池(ThreadPoolExecutor)
    • await 只能在协程中使用:await 关键字不能在普通函数中使用,必须在 async def 定义的协程中
    • 事件循环是单线程的:协程的并发是“协作式”的,需通过 await 主动交出执行权,否则会独占事件循环

asyncio 相关核心概念

协程(Coroutine)

  • 协程是可暂停、可恢复的函数 ,用 async def 定义(语法),是 asyncio 的核心执行单元
  • 协程函数调用后不会立即执行 ,而是返回一个协程对象(coroutine object),需通过事件循环调度才能运行

事件循环(Event Loop)

  • asyncio 的“大脑”,负责调度所有协程任务:
    • 管理任务的暂停/恢复、监听 I/O 事件、分发任务执行权
  • 通常通过 asyncio.run() 自动创建和管理事件循环(推荐用法)
    • 注:这行代码会直到整个时间循环运行完成才最终返回

等待对象(Awaitable)

  • 可被 await 关键字修饰的对象
    • 包括:协程对象、Task、Future
  • await 会暂停当前协程,等待目标对象完成后再恢复,期间事件循环可调度其他协程执行(实现并发)

Task(任务)

  • 对协程的封装,将协程注册到事件循环中,使其可被调度执行
  • 通过 asyncio.create_task() 创建,会自动加入事件循环并运行

Future

  • 表示异步操作的“未来结果”,是低层级的对象(通常无需手动创建,Task 继承自 Future)

asyncio 基础用法

  • 定义和运行协程

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

    # 定义协程函数(async def 关键字)
    async def hello(name):
    print(f"Hello, {name}! (开始)")
    # 模拟 I/O 等待(必须用 await 修饰可等待对象)
    # 注:await 后面必须跟一个“可等待对象”(比如还可以是另一个协程函数)
    await asyncio.sleep(1) # 这里的含义是暂停 1 秒,期间事件循环可执行其他任务,即让当前的协程暂停执行 1 秒钟,同时把控制权交还给事件循环
    print(f"Hello, {name}! (结束)")

    # 运行协程(Python 3.7+ 推荐用 asyncio.run())
    asyncio.run(hello("asyncio")) # asyncio.run 启动事件循环并以 hello 协程为入口执行协程(等待执行完成才最终返回)
    # Hello, asyncio! (开始)
    ## 【等待】... 1s
    # Hello, asyncio! (结束)
    • 补充理解:await asyncio.sleep(1)
      • 当协程执行到 await asyncio.sleep(1) 时,它会向事件循环注册一个延迟任务,然后自身进入挂起(suspended)状态
      • 此时,事件循环会立刻接管控制权,去执行其他准备好的协程
      • 1 秒钟后,事件循环会重新唤醒这个协程继续往下执行
  • await 对接协程函数

    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
    26
    27
    28
    29
    30
    import asyncio

    # 内部协程:返回一个数字
    async def get_number():
    await asyncio.sleep(1)
    return 42

    # 外部协程:使用 return await
    async def wrapper_coroutine():
    # # 下面 await get_number():显式地等待子协程完成,如果子协程内部抛出异常,异常会在这里被捕获或处理
    return await get_number() # await 后接协程函数,等待 get_number() 执行完成,return 负责返回对应执行结果

    # # 等价写法1:
    # x = await get_number() # 获取协程执行结果
    # return x

    # # 等价写法2:
    # x = get_number() # 直接调用协程时,只会返回协程对象,不会执行内部代码
    # x = await x # 等待协程执行完成
    # return x

    # # 错误写法:
    # return get_number() # 直接调用协程时,只会返回协程对象,不会执行内部代码

    async def main():
    # result 的值就是 42
    result = await wrapper_coroutine() # await 后接协程函数
    print(result)

    asyncio.run(main())
  • 特别说明:

    • await 只能在 协程中被使用,普通函数中不能使用 await
    • 协程调用的三种方式总结:
      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
      26
      27
      import asyncio

      async def my_coroutine():
      print("Hello from coroutine!")
      return 42

      # 调用方式一:使用 asyncio.run() 作为程序入口调用
      # 下面会 打印 "Hello from coroutine!"
      asyncio.run(my_coroutine()) # 启动事件循环并以 my_coroutine 协程为入口 执行协程

      # 调用方式二:在另一个协程中使用 await 调用,然后外面的协程被其他调用方式调用
      async def main():
      # 下面这行会打印 打印 "Hello from coroutine!"
      result = await my_coroutine() # 这里会真正执行并等待结果
      # # 等价于 下面两行:
      # x = my_coroutine() # 调用后返回一个 <coroutine object>,此时函数体内的代码完全没有执行,必须手动使用 await 关键字去等待它,它才会开始运行
      # result = await x # 这里会真正执行并等待结果
      print(result) # 输出: 42
      asyncio.run(main()) # 启动事件循环并以 my_coroutine 协程为入口 执行协程

      # 调用方式三:使用 asyncio.create_task() 并发调度
      async def main():
      task = asyncio.create_task(my_coroutine()) # 提交给事件循环并发执行,注意是立即后台巡行,调用后返回一个 <Task pending> 对象
      # print(task) # <Task pending name='Task-8' coro=<my_coroutine() running at /Users/sanye/Workspace/IdeaProjects/Python/llm_demo/temp/temp.py:3>>
      result = await task
      print(result)
      asyncio.run(main()) # 启动事件循环并以 my_coroutine 协程为入口 执行协程

并发执行多个协程

  • 通过 asyncio.gather()asyncio.create_task() 实现并发(多个任务同时执行,总耗时接近最长任务的耗时)

  • 方式 1:asyncio.gather()(批量等待多个协程)

    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

    async def task1():
    await asyncio.sleep(2)
    return "Task 1 完成"

    async def task2():
    await asyncio.sleep(1)
    return "Task 2 完成"

    async def main():
    # 并发执行 task1 和 task2,等待所有完成后返回结果(按传入顺序)
    result1, result2 = await asyncio.gather(task1(), task2())
    # # 等价与下面的两行
    # results = asyncio.gather(task1(), task2())
    # result1, result2 = results

    print(result1)
    print(result2)

    # 总耗时大约 2 秒(而非 2+1=3 秒)
    asyncio.run(main()) # 启动事件循环并以 main 协程为入口 执行协程

    # Task 1 完成
    # Task 2 完成
  • 方式 2:asyncio.create_task()(手动创建任务)

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    import asyncio

    async def task(name, delay):
    await asyncio.sleep(delay)
    print(f"Task {name} 完成(延迟 {delay} 秒)")

    async def main():
    # 创建任务并自动加入事件循环
    task1 = asyncio.create_task(task("A", 2)) # 立即开始执行
    task2 = asyncio.create_task(task("B", 1)) # 立即开始执行

    # 等待任务完成(可单独等待,也可一起等待)
    await task1
    await task2

    asyncio.run(main()) # 启动事件循环并以 main 协程为入口 执行协程

    # Task B 完成(延迟 1 秒)
    # Task A 完成(延迟 2 秒)

附录:处理异常

  • 协程中的异常需通过 try/except 捕获,或在 gather() 中通过 return_exceptions=True 收集异常
    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
    26
    import asyncio

    async def task1():
    await asyncio.sleep(2)
    return "Task 1 完成"

    async def faulty_task():
    await asyncio.sleep(1)
    raise ValueError("任务执行失败!") # 模拟抛出异常

    async def main():
    # 方式1:捕获单个协程的异常
    try:
    await faulty_task()
    except ValueError as e:
    print(f"捕获异常:{e}") # 捕获异常:任务执行失败!

    # 方式2:批量捕获多个协程的异常(return_exceptions=True)
    results = await asyncio.gather(
    task1(), # 正常任务,输出 "Task 1 完成"
    faulty_task(), # 异常任务,抛出异常 ValueError("任务执行失败!")
    return_exceptions=True # 不终止,返回异常对象
    )
    print(results) # 输出:["Task 1 完成", ValueError("任务执行失败!")]

    asyncio.run(main())

协程的进阶特性

超时控制(asyncio.wait_for()

  • 限制协程的执行时间,超时则抛出 TimeoutError
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    import asyncio

    async def long_task():
    await asyncio.sleep(3) # 模拟耗时 3 秒的任务

    async def main():
    try:
    # 限制任务 2 秒内完成,超时则取消任务并抛出异常
    result = await asyncio.wait_for(long_task(), timeout=2)
    except asyncio.TimeoutError:
    print("任务超时被取消!")

    asyncio.run(main())

任务取消(Task.cancel()

  • 手动取消正在执行的任务,被取消的任务会抛出 CancelledError
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    import asyncio

    async def endless_task():
    try:
    while True:
    print("任务运行中...")
    await asyncio.sleep(1)
    except asyncio.CancelledError:
    print("任务被取消!")
    raise # 可选:重新抛出,让调用方知道任务被取消

    async def main():
    task = asyncio.create_task(endless_task())
    await asyncio.sleep(2) # 运行 2 秒后取消
    task.cancel()
    await task # 必须等待任务处理取消逻辑

    asyncio.run(main())

异步上下文管理器(async with

  • 用于异步资源的获取和释放(如异步数据库连接、异步文件),需实现 __aenter____aexit__ 方法
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    21
    import asyncio

    class AsyncResource:
    async def __aenter__(self):
    print("获取异步资源")
    await asyncio.sleep(0.5)
    return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
    print("释放异步资源")
    await asyncio.sleep(0.5)

    async def main():
    async with AsyncResource() as res:
    print("使用异步资源")

    asyncio.run(main())

    # 获取异步资源
    # 使用异步资源
    # 释放异步资源

异步迭代器(async for

  • 用于迭代异步生成的数据(如异步流、分页接口),需实现 __aiter____anext__ 方法
    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
    26
    import asyncio

    class AsyncIterator:
    def __init__(self, limit):
    self.limit = limit
    self.count = 0

    def __aiter__(self):
    return self

    async def __anext__(self):
    if self.count >= self.limit:
    raise StopAsyncIteration
    self.count += 1
    await asyncio.sleep(0.5) # 模拟异步获取数据
    return self.count

    async def main():
    async for num in AsyncIterator(3):
    print(f"迭代得到:{num}")

    asyncio.run(main())

    # 迭代得到:1
    # 迭代得到:2
    # 迭代得到:3

附录:协程函数可以被继承和重写

  • 在 Python 中,协程函数(使用 async def 定义的函数)本质上仍然是类的方法
    • 它们遵循 Python 所有的面向对象规则,因此可以被继承、重写(Override),甚至可以通过 super() 调用父类的协程方法
  • 相关示例:
    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
    26
    27
    28
    29
    30
    import asyncio

    class BaseAPI:
    async def fetch_data(self, url):
    print(f"[BaseAPI] 正在从 {url} 获取数据...")
    await asyncio.sleep(1) # 模拟网络请求
    return {"status": 200, "data": "基础数据"}

    class AdvancedAPI(BaseAPI):
    # 1. 重写父类的协程方法
    async def fetch_data(self, url):
    print(f"[AdvancedAPI] 拦截请求: {url}")

    # 2. 通过 super() 调用父类的协程方法(注意必须加 await)
    result = await super().fetch_data(url)

    # 增加额外的处理逻辑
    result["enhanced"] = True
    return result

    async def main():
    api = AdvancedAPI()
    data = await api.fetch_data("https://example.com")
    print("最终结果:", data)

    asyncio.run(main())

    # [AdvancedAPI] 拦截请求: https://example.com
    # [BaseAPI] 正在从 https://example.com 获取数据...
    # 最终结果: {'status': 200, 'data': '基础数据', 'enhanced': True}

借助协程能实现 async with 语句

  • async with 是 Python 异步编程中的异步上下文管理器语法 ,基于协程实现
    • 可以把它理解为同步代码中 with 语句的异步版本,专门用于在协程环境中优雅、安全地管理那些需要异步初始化和清理的资源(如异步网络连接、数据库会话、异步文件读写等)
  • async with 的存在完全依赖于协程机制,只能async def 定义的协程函数内部使用

with 对比 async with

  • 同步 with :底层依赖对象的 __enter____exit__ 方法
    • 这两个方法执行时是同步的,会阻塞当前线程
  • 异步 async with :底层依赖对象的 __aenter____aexit__ 方法
    • 这两个方法必须使用 async def 定义 (即它们本身就是协程函数)
    • 当执行 async with 时,解释器会自动 await 这两个方法,从而在等待资源获取或释放时挂起当前协程,将控制权交还给事件循环,实现非阻塞并发

async with 核心工作机制

  • 当写下 async with obj as x: 时,Python 解释器在底层会做以下事情:
    • 1)调用 obj.__aenter__()await 它,获取资源对象并赋值给 x
    • 2)执行 with 代码块内部的逻辑
    • 3)无论代码块是正常结束还是抛出异常,都会自动调用 obj.__aexit__(exc_type, exc_val, exc_tb)await 它,完成资源的清理

使用示例:实际项目中的常见应用

  • 在实际开发中,不需要每次都手写 __aenter__,而是直接使用第三方异步库提供的对象:
    1
    2
    3
    4
    5
    6
    7
    8
    import aiohttp

    async def fetch_data(url):
    # aiohttp.ClientSession 是一个异步上下文管理器
    async with aiohttp.ClientSession() as session: # 这个类非常常用
    # session.get() 返回的响应对象也是一个异步上下文管理器
    async with session.get(url) as response:
    return await response.text()

自定义异步上下文管理器示例

  • 自定义异步上下文管理器示例代码
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    import asyncio

    class AsyncDatabaseSession:
    async def __aenter__(self):
    print("正在异步建立数据库连接...")
    await asyncio.sleep(1) # 模拟异步网络IO,不阻塞事件循环
    return self

    async def __aexit__(self, exc_type, exc_val, exc_tb):
    print("正在异步关闭数据库连接...")
    await asyncio.sleep(0.5) # 模拟异步清理操作

    async def main():
    # 只能在协程中使用 async with
    async with AsyncDatabaseSession() as session:
    print("正在安全地使用数据库会话...")
    await asyncio.sleep(1)

    asyncio.run(main())

asyncio.Semaphore() 的介绍和使用

  • asyncio.Semaphore 是 Python 异步编程中用于控制并发任务数量的核心同步原语
    • 理解:就像一个”限流器”,能确保同时执行的协程数量不超过预设的上限(常用于 控制网络请求并发数等)

标准用法

  • 在异步编程中,asyncio.Semaphore 的标准用法是配合 async with 管理上下文:

    1
    2
    3
    4
    5
    6
    7
    8
    sem = asyncio.Semaphore(max_conccurency)
    async def worker():
    # 统计耗时位置1:统计排队耗时 + 执行耗时
    async with sem: # 获取令牌(若满则阻塞等待)
    # 执行受限的 I/O 密集型操作(如网络请求、数据库读写)
    # # 统计耗时位置2:仅统计到执行耗时
    await do_io()
    # 离开上下文自动释放令牌
    • 最佳实践:将需要消耗物理资源(网络连接、文件句柄)的代码 放在 async with 内部,将本地计算/日志统计 放在外部
    • 若要

信号量(Semaphore)与锁(Lock)的区别

  • Semaphore 用于 限流(控制并发任务数)
    • 防止同时发起过多 请求,打爆 API 或数据库连接池
  • Lock 用于 互斥(保护共享数据)
    • 防止多个协程同时修改指定数据

高阶通用模式示例

  • 通用模板:
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    21
    22
    # 1. 初始化限流器与共享锁
    sem = asyncio.Semaphore(concurrency)
    lock = asyncio.Lock()
    results = []

    async def worker(data):
    # 2. 外部记录开始(用于计算排队延迟)
    start = time.time()

    # 3. 获取信号量,执行受限操作
    async with sem:
    result = await expensive_io_task(data) # 仅这里受并发限制

    # 4. 退出信号量立即计算总耗时(含排队)
    result['total_latency'] = time.time() - start

    # 5. 使用 Lock 保护共享状态写入(不占用信号量)
    async with lock:
    results.append(result)
    update_summary(result)
    await write_to_file(result) # 异步写文件(或用同步写但持有锁)
    pbar.update(1) # 更新进度条

附录:使用 asyncio.create_subprocess_exec 创建子进程

  • 官方说明文档:3.14.6 Documentation » Python 标准库 » 网络和进程间通信 » asyncio — 异步 I/O » 子进程

  • asyncio.create_subprocess_exec 是 Python 用于异步创建和管理子进程的高层 API

    • 可以在不阻塞事件循环的前提下执行外部程序,并与其输入/输出流进行交互
  • asyncio.create_subprocess_exec 函数签名与参数

    1
    2
    async asyncio.create_subprocess_exec(program, *args, stdin=None, stdout=None,
    stderr=None, limit=None, **kwds)
    • program (必需): 要执行的程序路径(字符串)
    • *args :程序的命令行参数 ,作为可变参数传入,可以一次输出多个(可按照顺序随意列出)
    • stdin, stdout, stderr :标准输入、输出、错误流的重定向方式,常用取值有:
      • asyncio.subprocess.PIPE: 创建一个管道,此时通过 Process 对象的 stdin/stdout/stderr 进行读写
      • asyncio.subprocess.DEVNULL: 重定向到 os.devnull
      • 也可以是已有的文件描述符或文件对象
    • limit :为 stdoutstderrStreamReader 包装器设置缓冲区大小上限(字节数)
    • **kwds :其他关键字参数(如 env, cwd 等),具体可参考 loop.subprocess_exec() 文档
    • 该协程返回一个 asyncio.subprocess.Process 实例

返回值说明

  • asyncio.subprocess.Process 实例可以进行如下操作:
    • 等待进程结束 :使用 await proc.wait()
    • 获取输出 :通过 proc.stdout.read() 等方法读取
    • 发送输入 :使用 proc.stdin.write()
    • 终止进程 :调用 proc.terminate()proc.kill()
    • 获取进程的 returncodepid 等信息

使用注意

  • 这是一个协程,必须使用 await 调用
  • 如果 Process 对象在子进程仍在运行时被垃圾回收 ,子进程会被杀掉
  • create_subprocess_exec 直接执行程序 ,不经过 shell 解析,因此没有 shell 注入风险 ,安全性更高
  • 可以配合 asyncio.gather 等工具并行执行和管理多个子进程

asyncio.create_subprocess_exec 使用示例

  • 异步执行 ls -l -a 命令并读取其输出:
    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    15
    16
    17
    18
    19
    20
    21
    22
    import asyncio
    import asyncio.subprocess

    async def run_ls():
    # 创建子进程,捕获 stdout 和 stderr
    proc = await asyncio.create_subprocess_exec(
    'ls', '-l', '-a',
    stdout=asyncio.subprocess.PIPE,
    stderr=asyncio.subprocess.PIPE
    )

    # 等待进程结束并获取输出
    stdout, stderr = await proc.communicate()

    print(f"进程退出码: {proc.returncode}")
    if stdout:
    print(f"[标准输出]\n{stdout.decode()}")
    if stderr:
    print(f"[标准错误]\n{stderr.decode()}")

    # 运行异步函数
    asyncio.run(run_ls())