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

Python Multiprocessing:如何配置进程池避免等待其他任务完成?

问题:多进程池完成部分任务后不自动处理新任务

我使用以下多进程代码,NCPUS_FOLDER设为4。发现当2个任务完成后,工作进程并未立即处理RUN_DIRS中的下2个任务,而是等待另外2个工作进程完成。查阅资料时,大多是询问如何让工作进程等待的方案(如使用异步版本),但我的场景是默认等待,需要改为不等待。请问是否是Pool的使用方式导致?如何让完成任务的工作进程继续处理新任务?

主程序代码:

from multiprocess import Pool
import time
import os

if __name__ == '__main__':
    t0 = time.time()
    with Pool(NCPUS_FOLDER) as pool:
        # 分配任务文件夹
        pool.map(main, RUN_DIRS)
    print(f"总运行时长:{time.time()-t0:.1f}秒")

main函数说明:每个并行工作进程会调用优化求解器,当打印语句输出时任务完成。在我的场景中,4个并行任务里有2个已输出完成语句,但未启动新任务。约束和成本函数评估过程中会创建文件夹、写入文件等。

main函数代码:

def main(run_dir):
    os.chdir(run_dir)
    
    # 创建嵌套求解器
    solver = BuckshotSolver(dim=len(TARGET), npts=4)
    # 指定搜索内部使用的优化算法
    solver.SetNestedSolver(PowellDirectionalSolver)

    # 设置约束条件
    solver.SetConstraints(constraint_func)

    # 寻找最小值,传入终止条件
    solver.Solve(cost_func, termination=VTR(tolerance=TOLERANCE), disp=True)

    print(f"在{run_dir}中找到的最优解为{solver.bestSolution}")

解答

你的问题确实是pool.map()的特性导致的。map()方法会将任务按进程数分成固定批次,必须等待当前批次的所有任务全部完成后,才会将下一批任务分配给工作进程。比如你有4个进程,若RUN_DIRS中有6个任务,map()会先把前4个任务分给4个进程,只有等这4个任务全完成,才会把剩下的2个分配给空闲进程——这就是你看到前两个任务完成后,进程不处理新任务的原因。

要让完成任务的进程立刻处理新任务,你需要改用异步分配任务的方法,以下是两种可行方案:

方案1:使用imap_unordered()

这个方法会在任务完成后立刻返回结果,并且一旦有进程空闲就会分配新任务,无需等待整批任务完成。修改后的主程序代码如下:

from multiprocess import Pool
import time
import os

if __name__ == '__main__':
    t0 = time.time()
    with Pool(NCPUS_FOLDER) as pool:
        # 用imap_unordered替代map,实现实时任务分配
        for _ in pool.imap_unordered(main, RUN_DIRS):
            # 若不需要处理任务结果,可省略循环体
            pass
    print(f"总运行时长:{time.time()-t0:.1f}秒")

注意:imap_unordered()返回结果的顺序与任务提交顺序无关,哪个任务先完成就先返回哪个。如果需要保持结果顺序与提交顺序一致,可以用imap(),它同样支持实时分配任务。

方案2:使用apply_async()

这种方式更灵活,可手动提交每个任务,并通过回调函数处理任务完成后的逻辑:

from multiprocess import Pool
import time
import os

def task_callback(_):
    # 任务完成后的回调操作,例如打印日志
    pass

if __name__ == '__main__':
    t0 = time.time()
    with Pool(NCPUS_FOLDER) as pool:
        # 逐个提交任务
        task_list = [
            pool.apply_async(main, args=(run_dir,), callback=task_callback)
            for run_dir in RUN_DIRS
        ]
        # 等待所有任务完成
        for task in task_list:
            task.get()
    print(f"总运行时长:{time.time()-t0:.1f}秒")

apply_async()会在进程空闲时立即分配新任务,不会等待整批任务结束。

额外注意事项

  • 确保main函数中的文件操作、文件夹创建等逻辑是进程安全的,避免多个进程同时操作同一资源引发问题。
  • 当RUN_DIRS中的任务数量远大于进程数时,异步方法的效率会远高于map(),因为不会出现进程空闲等待的情况。

内容的提问来源于stack exchange,提问作者Jesse Feng

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 11:27:44