如何使用asyncio实现依赖任务与非依赖任务的并行执行及循环内任务的优化调度
看起来你在asyncio的任务并行调度上遇到了两个典型问题,我来帮你逐个解决~
一、解决第一部分:依赖前置任务后的并行执行
你提到第一部分的代码看起来是顺序执行,但其实asyncio.gather()本身就是用来并行运行多个协程的核心工具,可能是你测试时的误解?不过先确认核心逻辑:first_execution是同步函数,会立即执行并返回结果,之后我们要让check_first_execution(耗时10秒)和second_execution(耗时5秒)同时启动,这样后者会先完成。
你的原代码逻辑其实是对的,但为了更直观,我们可以用asyncio.create_task()显式创建后台任务,再等待,这样更清晰:
async def main(): choice = random.choice([0, 1]) result_from_first = first_execution(choice) # 第一部分:修正后的并行逻辑 # 先把check任务放到后台运行 check_task = asyncio.create_task(check_first_execution(result_from_first)) # 根据结果选择执行second或third if result_from_first == 1: await second_execution(result_from_first) else: await third_execution(result_from_first) # 等待后台的check任务完成(如果还没结束的话) await check_task # 如果你需要收集两个任务的结果,也可以用gather,效果完全一致: # if result_from_first ==1: # results = await asyncio.gather(check_first_execution(result_from_first), second_execution(result_from_first)) # else: # results = await asyncio.gather(check_first_execution(result_from_first), third_execution(result_from_first))
这样运行后,输出会完全符合你的预期:
First result 1 Second result 2 First check of value 1 complete
因为check_task和second_execution是同时启动的,后者耗时更短,会先输出结果,10秒后check任务完成并打印信息。
二、解决第二部分:循环内的任务调度优化
你现在的循环是每次先等second_execution完成,再等check_second_execution完成才进入下一次循环,这导致了完全的顺序执行。要实现“先连续执行多个second,同时check在后台运行”的效果,我们可以把check_second_execution作为后台任务提交,不用立即等待,等所有second任务都执行完后,再统一等待所有check任务结束:
优化方案:后台运行check任务,不阻塞循环
async def main(): # ... 前面的第一部分代码 ... # 第二部分:优化后的调度逻辑 list_results_from_first = [result_from_first+i for i in range(5)] # 用列表存储所有check任务,最后统一等待 check_tasks = [] for first_result in list_results_from_first: # 先执行second_execution,拿到结果 second_result = await second_execution(first_result) # 创建check任务放到后台运行,不等待直接进入下一次循环 check_task = asyncio.create_task(check_second_execution(second_result)) check_tasks.append(check_task) # 所有second任务都执行完后,等待所有后台的check任务完成 await asyncio.gather(*check_tasks)
这样运行后的输出会类似:
Second result 3 Second result 4 Second result 5 Second result 6 Second result 7 Second check of value 3 complete Second check of value 4 complete Second check of value 5 complete
check任务会在后台并行运行,各自完成后就会打印信息,而second任务会依次执行(符合你需求里的“先跑多个second”的逻辑)。
如果你的需求是second任务也想并行执行,还可以进一步优化,把所有second任务也放到后台:
async def main(): # ... 前面的第一部分代码 ... # 第二部分:second也并行的版本 list_results_from_first = [result_from_first+i for i in range(5)] # 一次性创建所有second任务并并行执行 second_tasks = [asyncio.create_task(second_execution(res)) for res in list_results_from_first] second_results = await asyncio.gather(*second_tasks) # 再启动所有check任务并行执行 check_tasks = [asyncio.create_task(check_second_execution(res)) for res in second_results] await asyncio.gather(*check_tasks)
这种情况下,所有second任务会同时启动,5秒后一起输出结果,10秒后所有check任务完成并打印信息,适合任务之间没有依赖的场景。
完整修正后的代码
整合两部分优化后的完整代码如下:
import asyncio import random def first_execution(choice): if choice == 1: print(f"First result {choice}") return choice else: print(f"First result {0}") return 0 async def check_first_execution(result_from_first): await asyncio.sleep(10) print(f"First check of value {result_from_first} complete") async def second_execution(result_from_first): await asyncio.sleep(5) print(f"Second result {result_from_first+1}") return result_from_first+1 async def check_second_execution(result_from_second): await asyncio.sleep(10) print(f"Second check of value {result_from_second} complete") async def third_execution(result_from_first): await asyncio.sleep(5) print(f"Third result {result_from_first+2}") return result_from_first+2 async def check_third_execution(result_from_third): await asyncio.sleep(10) print(f"Third check of value {result_from_third} complete") async def main(): choice = random.choice([0, 1]) result_from_first = first_execution(choice) # 第一部分:并行执行check_first和second/third check_first_task = asyncio.create_task(check_first_execution(result_from_first)) if result_from_first == 1: await second_execution(result_from_first) else: await third_execution(result_from_first) await check_first_task # 第二部分:循环内优化,后台运行check任务 list_results_from_first = [result_from_first+i for i in range(5)] check_tasks = [] for first_result in list_results_from_first: second_result = await second_execution(first_result) check_task = asyncio.create_task(check_second_execution(second_result)) check_tasks.append(check_task) await asyncio.gather(*check_tasks) asyncio.run(main())
备注:内容来源于stack exchange,提问作者HienPham

