如何正确实现start_analysis()函数的100次并行执行?
问题分析
你的代码现在是串行执行的,原因有两个:
- 每次循环调用
asyncio.run(),会新建并销毁一个事件循环,每次只执行一个start_analysis任务,完全没有并行效果。 start_analysis里的func_1/2/3都是同步调用,没有await关键字,即使放到同一个事件循环里,也会阻塞整个循环,无法切换到其他任务。
修正方案
根据你的需求,分两种场景处理:
场景1:func是IO密集型(比如文件读写、网络请求)
把同步函数包装成可异步等待的任务,同时一次性创建所有任务,让事件循环调度它们并行执行:
import asyncio async def start_analysis(trial, data): ranges = generate_ranges() # 用asyncio.to_thread把同步IO函数包装成异步任务,允许事件循环切换 res_1 = await asyncio.to_thread(func_1, ranges) res_2 = await asyncio.to_thread(func_2, ranges) res_3 = await asyncio.to_thread(func_3, ranges) map_results(trial, res_1, res_2, res_3) async def main(): data = retrieve_data() # 批量创建100个异步任务 tasks = [start_analysis(trial, data) for trial in range(1, 101)] # 等待所有任务完成 await asyncio.gather(*tasks) if __name__ == "__main__": asyncio.run(main())
场景2:func是CPU密集型(比如大量计算)
asyncio单线程异步无法利用多核CPU,这时候用多进程并行才是最高效的:
from multiprocessing import Pool def start_analysis(trial): # 数据量大的话建议在主进程读取后传入,避免子进程重复加载 data = retrieve_data() ranges = generate_ranges() res_1 = func_1(ranges) res_2 = func_2(ranges) res_3 = func_3(ranges) map_results(trial, res_1, res_2, res_3) def main(): # 进程池默认使用CPU核心数,可手动指定参数如Pool(8) with Pool() as pool: # 并行执行100次任务 pool.map(start_analysis, range(1, 101)) if __name__ == "__main__": main()
关键说明
- 用asyncio实现并行的核心是:在同一个事件循环中创建多个任务,且任务中有可等待的异步操作(await),这样事件循环才能在等待时切换到其他任务。
- 如果
func_1/2/3本身就是异步函数(定义为async def),直接调用await func_xxx(ranges)即可,无需to_thread包装。
内容的提问来源于stack exchange,提问作者HEIWO78
相关产品推荐
相关产品推荐

