Python多进程Pool任务分配不符合预期原因及优化方案咨询
multiprocessing.Pool任务分配不符合预期的原因与解决方法
问题重现
代码:
import multiprocessing from time import sleep import os def multi_test(time): print(f'{os.getpid()}: {time}s') sleep(time) pool = multiprocessing.Pool(2) pool.map(multi_test, [1,1,1,1,4,1,1,1,1]) pool.close() pool.join()
实际输出:
68858: 1s 68857: 1s 68858: 1s 68857: 1s 68857: 4s 68858: 1s 68858: 1s 68858: 1s # 该行后有2秒延迟,因为另一个核心在执行下一个任务(为什么?) 68857: 1s
实际总耗时7秒,预期为6秒——原本认为4秒任务可与最后四个1秒任务并行,但最后一个1秒任务被分配给了执行4秒任务的进程,时间线:
<1><2><3><4><5><6><7> 68858: 1s 1s 1s 1s 1s 68857: 1s 1s <----4s---> 1s
当任务列表改为[1,1,4,1,1,1,1]时,结果符合预期,总耗时5秒:
<1><2><3><4><5> 68858: 1s 1s 1s 1s 1s 68857: 1s <----4s--->
原因分析
multiprocessing.Pool.map() 默认采用批量分块分配任务的策略:
- 它会根据任务总数和进程数量,将任务分成大致相等的块,每个块分配给一个工作进程。
- 分块大小计算公式为:
chunksize = 任务数 // 进程数 + (1 if 任务数 % 进程数 != 0 else 0)
在第一个测试案例中:
- 任务总数9,进程数2,分块大小为5,导致包含4秒任务的块末尾还有一个1秒任务。
- 块内任务必须按顺序执行,执行4秒任务的进程要等长任务完成后才能处理末尾的1秒任务;而另一个进程虽提前完成所有分配的任务,但无法跨块接管剩余任务,最终整体耗时被拖到7秒。
第二个测试案例中,任务列表的长度和顺序刚好让4秒任务所在的块没有后续任务,所有短任务都被分配到另一个进程,实现了并行执行,因此总耗时符合预期。
解决方法
要实现动态调度(空闲进程立即接管剩余任务),有两种方式:
1. 指定chunksize=1
强制map()逐个分配任务,而非批量分块:
import multiprocessing from time import sleep import os def multi_test(time): print(f'{os.getpid()}: {time}s') sleep(time) pool = multiprocessing.Pool(2) # 指定chunksize=1,逐个分配任务 pool.map(multi_test, [1,1,1,1,4,1,1,1,1], chunksize=1) pool.close() pool.join()
修改后,空闲进程会立即接手剩余任务,时间线变为:
<1><2><3><4><5><6> 68858: 1s 1s 1s 1s 1s 1s 68857: 1s 1s <----4s--->
总耗时6秒,符合预期。
2. 使用imap()或imap_unordered()
这两个方法默认逐个迭代任务并分配给空闲进程,效果和chunksize=1的map()一致:
import multiprocessing from time import sleep import os def multi_test(time): print(f'{os.getpid()}: {time}s') sleep(time) pool = multiprocessing.Pool(2) # 使用imap逐个分配任务 for _ in pool.imap(multi_test, [1,1,1,1,4,1,1,1,1]): pass pool.close() pool.join()
内容的提问来源于stack exchange,提问作者Matthias Schmitz
相关产品推荐
相关产品推荐

