发布于: -/最后更新: -/11 分钟/

协程库 asyncio 介绍

摘要

协程是一种用户态内的上下文切换技术,允许在单线程中实现并发操作。Python中的asyncio库提供了对协程的支持,包括定义协程函数、创建协程对象、使用await表达式等。事件循环是asyncio的核心,负责管理和调度协程的执行。Task对象用于封装协程,Future对象用于表示异步操作的结果。asyncio还提供了run函数、wait函数、run_in_executor函数等API,用于执行协程、等待任务完成和在执行器中运行函数。异步迭代器和异步上下文管理器也得到了支持。uvloop是asyncio事件循环的替代方案,提供了更高的性能。

协程 Coroutine

也称之为微线程,是一种用户态内的上下文切换技术。
简而言之,就是一种技术可以使得在一个线程的情况下实现代码快互相切换执行

以下代码中,假定#switch函数是用于实现这样的一种功能(用户态内的上下文切换)

Python
def f1():
	print(1)
	switch(f2)
	print(2)
	switch(f2)

def f2():
	print(3)
	switch(f1)
	print(4)

switch(f1)
plain
1
3
2
4

asyncio的优势

其优势在于它可以做到遇到IO阻塞情况下可自动进行切换

定义

  • 协程函数:

Python
# 1
async def function():
	pass

# 2
function
  • 协程对象 (Coroutine Object) awaitable:

Python
function()
  • “async” defines a coroutine

  • “await” suspends or yields execution to another coroutine from within a coroutine.

  • Awaitables

    • We say that an object is an awaitable object if it can be used in an await expression. Many asyncio APIs are designed to accept awaitables.

    • There are three main types of awaitable objects: coroutines, Tasks, and Futures.

事件循环

pseudo_code
协程列表 = [coro1, coro2, coro3, ...]

while 协程列表不为空:
    可执行的协程列表,已完成的协程列表 = 从协程列表中检查所有的协程并将'可执行'和'已完成'的协程返回
    
    for 就绪协程 in 可执行的协程列表:
        尝试执行就绪协程
        
    for 已完成的协程 in 已完成的协程列表:
        从协程列表中删除已完成的协程

可执行的协程列表:这个列表包含了可以立即执行的协程,也就是那些没有阻塞等待外部资源(通常是 I/O 操作)的协程。这是事件循环的关键部分,因为异步编程的主要目的之一是在等待 I/O 操作时允许其他协程执行。(ChatGPT)

协程函数的特性:

一个协程函数中可以有多个await协程对象/Task对象/Future对象

Python
async f1():
	pass

async f2():
	pass

async def main():
	await f1()
	await f2()
	...
  1. 协程标志: 使用 async 关键字定义的函数被标记为协程函数。这意味着这些函数可以在执行过程中被暂停和恢
    复,以允许其他协程执行。

  2. await 表达式: 在 async 函数中,可以使用 await 关键字来等待一个异步操作完成。当遇到 await 时
    ,协程函数会被暂停,等待异步操作完成后再继续执行。

  3. 非阻塞式: async 函数的特性使其在等待I/O操作等异步任务时不会阻塞整个程序的执行,而是可以切换到其他
    协程,从而提高并发性能。

  4. 生成器语法: async 函数的定义方式类似于生成器函数,使用 async def 关键字,而函数中的 yield 语句则变成了 await。

  5. 异步编程库支持: async 函数通常与异步编程库(如 asyncio、Trio 等)一起使用,这些库提供了管理、
    调度和执行协程的机制。

  6. 返回协程对象: 调用 async 函数会返回一个协程对象,而不是立即执行函数体。要执行协程,需要将其作为任
    务或事件注册到事件循环中。

  7. 异常处理: 协程函数中的异常会传播回调用者。可以使用 try 和 except 来捕获和处理协程函数中的异常。

  8. 并发性: async 函数允许在同一个线程中处理多个协程,从而在单线程中实现并发操作。这对于I/O密集型操作
    非常有用。

Task对象

Tasks are used to schedule coroutines concurrently.

When a coroutine is wrapped into a Task with functions like asyncio.create_task() the coroutine is automatically scheduled to run soon

在事件循环中添加多个任务的

除了使用asyncio.create_task()函数(在Python 3.7中被加入)以外,还可以用低层级的loop.create_task()或ensure_future()函数。不建议手动实例化Task对象。

Python
import asyncio as ai  
import time  
  
  
async def io1():  
    await ai.sleep(3)  
  
  
async def io2():  
    await ai.sleep(3)  
  
  
async def f1():  
    print("func1 begin")  
    # 忙等待开始, 由于f2协程任务为可执行任务所以切换上下文到f2  
    await io1()  
    print("func1 end")  
  
  
async def f2():  
    print("func2 begin")  
    # 执行到此,忙等待开始,由于此时没有可执行任务所以开始等待  
    res = await io2()  
    print("fun2 end", res)  
  
  
async def main():  
    print("MAIN begin")  
    # 创建Task对象, 并且将f1协程对象任务添加到事件循环  
    t1 = ai.create_task(f1())  
    # 创建Task对象, 并且将f2协程对象任务添加到事件循环  
    t2 = ai.create_task(f2())  
    print("MAIN end\n")  
  
    t0 = time.time()  
    # main()协程对象开始阻塞,事件循环...
    ret1 = await t1
    ret2 = await t2
    print(time.time() - t0)  
    print(ret1, ret2)  
  
  
if __name__ == '__main__':  
    ai.run(main())

常用写法

  1. 由run函数直接获取事件循环

Python
...

async def main():  
    # 创建Task对象, 并且将fn协程对象任务添加到事件循环  
    task_list = [
        ai.create_task(f1(), name="f1"),  
        ai.create_task(f2(), name="f2")  
    ]
    # 或者 (< 3.11)
    # DeprecationWarning: The explicit passing of coroutine objects to asyncio.wait() is deprecated since Python 3.8, and scheduled for removal in Python 3.11.

	task_list = [f1(), f2()]

    # done: 完成的任务  
    # pending: 超过timeout后未完成的任务  
    done, pending = await ai.wait(task_list, timeout=1)  
    print(done, pending)

ai.run(main())
  1. 手动获取事件循环

Python
...
task_list = [f1(), f2()]
# loop = ai.get_event_loop()  # 已弃用(测试环境Py 3.10)
loop = asyncio.get_event_loop_policy().get_event_loop()
done, pending = loop.run_until_complete(ai.wait(task_list))
print(done, pending)

asyncio.Future对象

A Future is a special low-level awaitable object that represents an eventual result of an asynchronous operation.
一个对象用来忙等待

Future内部具有_state_属性,当其状态为已完成的时候,则调用其的await表达式将停止阻塞

Task继承Future,Task对象内部await结果的处理基于Future对象

Python
async def main():
	loop = asyncio.get_running_loop()
	# 创建一个任务(Future对象),这个任务什么也不干
	fut = loop.create_future()
	# 等待任务的最终结果(Future对象),没有则会一直等待下去
	await fut

asyncio.run(main())
Python
async def set_after(fut):
	await asyncio.sleep(2)
	fut.set_result("114514")

async def main():
	loop = asyncio.get_running_loop()
	# 创建一个任务(Future对象),这个任务什么也不干
	fut = loop.create_future()

	# 创建一个任务(Task对象),绑定了set_after函数,该函数执行2s后,会给fut赋值
	# 当fut被设置结果后,它的状态就变为已完成了,于是调用者便不会被继续阻塞
	loop.create_task(set_after(fut))
	print(await fut)

asyncio.run(main())

深入理解事件循环

情况1:

Python
import asyncio  
import time  
  
  
async def set_after(fut):  
    await asyncio.sleep(2)  
    fut.set_result("114514")  
  
  
async def main():  
    loop = asyncio.get_running_loop()  
    # 创建一个任务(Future对象),这个任务什么也不干  
    fut = loop.create_future()  
  
    # 创建一个任务(Task对象),绑定了set_after函数,该函数执行2s后,会给fut赋值  
    # 当fut被设置结果后,它的状态就变为已完成了,于是调用者便不会被继续阻塞  
    t = loop.create_task(set_after(fut))  
    await t
    t0 = time.time()  
    time.sleep(3)  
    print(await fut)  
    total = time.time() - t0
  
  
asyncio.run(main())

total ≈ 3秒

情况2:

Python
...
async def main():  
    loop = asyncio.get_running_loop()  
    # 创建一个任务(Future对象),这个任务什么也不干  
    fut = loop.create_future()  
  
    # 创建一个任务(Task对象),绑定了set_after函数,该函数执行2s后,会给fut赋值  
    # 当fut被设置结果后,它的状态就变为已完成了,于是调用者便不会被继续阻塞  
    t = loop.create_task(set_after(fut))  
    t0 = time.time()  
    time.sleep(3)  
    print(await fut)  
    total = time.time() - t0
...

total ≈ 5秒

原因:

情况1:

  1. 当执行到await t时,main()协程对象开始阻塞。

  2. 事件循环执行按照任务列表中,切换上下文开始执行set_after(fut)协程对象,执行到await asyncio.sleep(2)时,该协程对象开始阻塞

  3. 当2秒过去后,set_after(fut)执行完fut.set_result(114514)后结束,协程对象完成,t任务完成

  4. await t停止阻塞,...

  5. await fut由于已经有结果,所以取操作耗时约为0

  6. 最后为3秒

情况2:

  1. ...

  2. 执行到await fut阻塞

  3. 事件循环,切换上下文开始执行set_after(fut)协程对象

  4. ...

  5. 最后为5秒

常用API

  • asyncio.run(coro, *, debug=None)

    • Execute the coroutine coro and return the result.

    • This function runs the passed coroutine, taking care of managing the asyncio event loop, finalizing asynchronous generators, and closing the threadpool.

    • This function cannot be called when another asyncio event loop is running in the same thread.

    • If debug is True, the event loop will be run in debug mode. False disables debug mode explicitly. None is used to respect the global Debug Mode settings.

    • This function always creates a new event loop and closes it at the end. It should be used as a main entry point for asyncio programs, and should ideally only be called once.

  • coroutine asyncio.wait(aws, *, timeout=None, return_when=ALL_COMPLETED)

    • Run Future and Task instances in the aws iterable concurrently and block until the condition specified by return_when.

    • aws列表中可以为 Future对象 或 Task对象 或 协程对象

    • The aws iterable must not be empty and generators yielding tasks are not accepted.

    • Returns two sets of Tasks/Futures: (done, pending).

    • 可以是协程对象的原因:

Python
async def wait(fs, *, timeout=None, return_when=ALL_COMPLETED):
	...
	"""
	Wrap a coroutine or an awaitable in a future.  
 
	If the argument is a Future, it is returned directly.  
	"""
	fs = {ensure_future(f, loop=loop) for f in fs}
	return await _wait(fs, timeout, return_when, loop)
  • loop.rununtil_complete(_future)

    • Run until the future (an instance of Future) has completed.

    • If the argument is a coroutine object it is implicitly scheduled to run as a asyncio.Task.

    • Return the Future’s result or raise its exception.

  • awaitable loop.runin_executor(_executor, func, *args)

    • Arrange for func to be called in the specified executor.

    • The executor argument should be an concurrent.futures.Executor instance. The default executor is used if executor is None.

asyncio + 不支持异步的模块

结合多线程库中的线程池等效使用异步

不同点:

  • 底层使用的是线程池,原理不同,因此为多线程操作,更耗费资源

    • 而原始asyncio为单线程操作

Python
import time
import asyncio
import concurrent.futures

def f():
	"""
	不支持异步的函数
	"""
	time.sleep(2)
	return "114514"

def method1():
	# 1. Run in the default loop's executor (Default is ThreadPoolExecutor)
	# Step 1. The internal will use ThreadPoolExecutor's submit method to apply for a thread to execute #f function
	# Step 2. Then it will call asyncio.wrap_future to wrap concurrent.futures.Future as asyncio.Future
	# Due to concurrent.futures.Future doesn't support await
	fut = loop.run_in_executor(None, f)
	result = await fut
	print('default thread pool', result)

def method2():
	# 2. Run in a custom thread pool
	# Suited to IO-bound Tasks, not CPU-bound tasks.
	with concurrent.futures.ThreadPoolExecutor() as pool:
		result = await loop.run_in_executor(
			pool, f)
		print('custom thread pool', result)

def method3():
	# 3. Run in a custom process pool
	# Suited to CPU-bound Tasks, probably not IO-bound tasks.
	with concurrent.futures.ProcessPoolExecutor() as pool:
		result = await loop.run_in_executor(
			pool, f)
		print('custom thread pool', result)


async def main():
	method1()
	loop = asyncio.get_running_loop()

Implemented method:

Python
def run_in_executor(self, executor, func, *args):  
    self._check_closed()  
    if self._debug:  
        self._check_callback(func, 'run_in_executor')  
    if executor is None:  
        executor = self._default_executor  
        # Only check when the default executor is being used  
        self._check_default_executor()  
        if executor is None:  
            executor = concurrent.futures.ThreadPoolExecutor(  
                thread_name_prefix='asyncio'  
            )  
            self._default_executor = executor  
    return futures.wrap_future(  
        executor.submit(func, *args), loop=self)
Python
def wrap_future(future, *, loop=None):  
    """Wrap concurrent.futures.Future object."""  
    if isfuture(future):  
        return future  
    assert isinstance(future, concurrent.futures.Future), \  
        f'concurrent.futures.Future is expected, got {future!r}'  
    if loop is None:  
        loop = events._get_event_loop()  
    new_future = loop.create_future()  
	"""
	Chain two futures so that when one completes, so does the other.  
	"""
    _chain_future(future, new_future)
    return new_future

案例 Case

Python
import asyncio  
  
import requests  
  
  
async def download(url, i):  
    loop = asyncio.get_event_loop()  
    fur = loop.run_in_executor(None, requests.get, url)  
    print(f"Start to download No.{i} image")  
    response = await fur  
    print(f"Done No.{i} image")  
    return response  
  
  
if __name__ == '__main__':  
    URL_LIST = [  
        "https://fastly.picsum.photos/id/685/1920/1080.jpg?hmac=GjjlhGiZFP-hXkJ4S2r2UwMqVqeBH6ky7FAe3DTgrmg",  
        "https://fastly.picsum.photos/id/527/1920/1080.jpg?hmac=FuLyw1LQ-LThCpTUUvrpF-OJQwsp18wwp-d6hmlI9E0",  
        "https://fastly.picsum.photos/id/330/1920/1080.jpg?hmac=1U_9bH5VS8l-_PJBf1Yrfp-d9NuhgOr4Zjditft-7sw",  
        "https://fastly.picsum.photos/id/809/200/300.jpg?hmac=jC-cQrqqx-NPPfMItPjmHx8XKCKi5WRG46ds3qYReKI",  
    ]  
    tasks = [download(URL_LIST[i], i) for i in range(len(URL_LIST))]  
    loop = asyncio.get_event_loop_policy().get_event_loop()  
    done, pending = loop.run_until_complete(asyncio.wait(tasks))  
    print(len(done))

异步迭代器

Asynchronous Iterators and “async for”
An asynchronous iterable is able to call asynchronous code in its iter implementation, and asynchronous iterator can call asynchronous code in its next method. To support asynchronous iteration:

  1. An object must implement an __aiter__ method (or, if defined with CPython C API, tp_as_async.am_aiter slot) returning an asynchronous iterator object.

  2. An asynchronous iterator object must implement an __anext__ method (or, if defined with CPython C API, tp_as_async.am_anext slot) returning an awaitable.

  3. To stop iteration __anext__ must raise a StopAsyncIteration exception.

异步上下文管理器

Asynchronous Context Managers and “async with”

An asynchronous context manager is a context manager that is able to suspend execution in its enter and exit methods.

To make this possible, a new protocol for asynchronous context managers is proposed. Two new magic methods are added: __aenter__ and __aexit__. Both must return an awaitable.

Python
import asyncio

class AsyncContextManager:
    async def __aenter__(self):
        await log('entering context')
        return await asyncio.sleep(3, self)

    async def __aexit__(self, exc_type, exc, tb):
        await log('exiting context')

async def func():
	async with AsyncContextManager() as f:
		...

asyncio.run(func())

uvloop

是asyncio事件循环的替代方案

pip3 install uvloop

Python
import asyncio
import uvloop
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
...

正文结束