You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何正确实现start_analysis()函数的100次并行执行?

问题分析

你的代码现在是串行执行的,原因有两个:

  1. 每次循环调用asyncio.run(),会新建并销毁一个事件循环,每次只执行一个start_analysis任务,完全没有并行效果。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 01:20:32