从抽象视角理解async/await工作机制及asyncio任务调度优先级
高层抽象视角理解async/await工作机制与asyncio调度逻辑
需求背景
我希望从高层半抽象视角理解async/await的工作机制。网上已有不少长篇复杂的解释,但需要更简洁的、基于栈或队列的机制说明,以满足科研人员与工程师快速了解asyncio执行async函数时对各类代码块的优先级调度逻辑的需求。
示例代码:批量异步URL请求
以使用aiohttp批量发起URL的GET请求为例,代码如下:
import time import asyncio import aiohttp async def fetch_url(session, url): async with session.get(url) as response: return await response.text() async def fetch_urls(session, url): print(url + ' fetch 1 started') text_response = await fetch_url(session, url) print(url + ' fetch 1 ended') await asyncio.sleep(4) print(url + ' fetch 2 started') text_response = await fetch_url(session, url) print(url + ' fetch 2 ended') return text_response async def main(): start_time = time.time() urls = [ 'https://httpbin.org/delay/1', 'https://httpbin.org/delay/2', 'https://httpbin.org/delay/3', ] async with aiohttp.ClientSession() as session: tasks = [fetch_urls(session, url) for url in urls] results = await asyncio.gather(*tasks) for i, result in enumerate(results): print(f'Response {i+1} received with {len(result)} chars') print(f'Time taken: {time.time() - start_time:.2f} seconds') if __name__ == '__main__': asyncio.run(main())
上述代码可扩展更多不同URL,在fetch_urls函数中增加更多asyncio.sleep与fetch_url调用,且URL响应延迟和asyncio.sleep时长均可设为随机值。
自定义术语与抽象合理性疑问
为表述精准,定义以下术语:
- 代码块(chunk):指需按顺序执行的命令序列,每个命令执行完成后才能执行下一个,非async函数即为单个代码块。
async关键字修饰def时,会告知Python对应的async函数(更准确说是“协程coroutine”)不应被视为单个代码块,而是由await命令分隔的多个代码块序列。例如上述fetch_urls包含4个代码块:print/fetch、print/sleep、print/fetch、print/return,相邻代码块间通常存在未知时长的延迟。
请问该抽象方式是否合理?
调度优先级核心问题
若上述抽象合理,每个执行中的fetch_urls实例被称为“任务(task)”,可将首个任务表示为有序序列:t0 = (c00, d00, c01, d01, c02, d02, c03),其中t代表任务,c代表代码块,d代表延迟。
asyncio运行main时,似乎会按main中任务的定义顺序执行每个任务的首个代码块。代码块执行完成后,进入到下一个代码块的延迟阶段;延迟结束后,该任务的下一个代码块进入待调度队列,且似乎遵循FIFO(先进先出)优先级规则。但具体的调度优先级机制是什么?例如,当c03正在执行时,c13和c22先后进入待调度队列,c03执行完成后会优先执行c13还是c22?
补充验证:调度FIFO特性的代码示例
以下代码示例似乎表明asyncio的部分队列调度优先级为FIFO,而非按激活时间排序:
from time import perf_counter as pc import asyncio pc0 = pc() async def sleep_then_sum(k, d, n): print(f'sleep_then_sum {k} started at {pc()-pc0} sec') print(f'sleep_then_sum {k} falls asleep at {pc()-pc0} sec') await asyncio.sleep(d) print(f'sleep_then_sum {k} starts computing sum at {pc()-pc0} sec') sn = sum(range(n)) print(f'sleep_then_sum {k} done at {pc()-pc0} sec') return sn async def main(): m = 3 task_list = m * [None] delay = [1, 3, 2] n_to_sum = [2*10**8, 10**8, 10**8] async with asyncio.TaskGroup() as tg: for k in range(m): task_list[k] = tg.create_task(sleep_then_sum(k, delay[k], n_to_sum[k])) await asyncio.sleep(0.2) sn_list = [task.result() for task in task_list] return sn_list sn_list = asyncio.run(main()) print(sn_list) print(f'Total time taken: {pc()-pc0} sec')
内容的提问来源于stack exchange,提问作者SapereAude
相关产品推荐
相关产品推荐

