Python多进程Pool.map分块时剩余任务均衡分配方案咨询
首先得明确为什么会出现你遇到的情况:multiprocessing.Pool.map的分块逻辑是按固定大小顺序切割任务列表,当总任务数没法被chunk整除时,最后一块就是剩余的任务数。比如你17个任务、chunk=6的场景,会被切成3块:[0-5]、[6-11]、[12-16]。Pool会按顺序把块分配给空闲进程,前两块各给一个进程,第三块则会被分配给第一个完成前序任务的进程(也就是你看到的Worker-1),直接导致这个进程要处理11个任务,另一个只处理6个,负载严重不均。
下面给你几个可行的解决方案:
方法1:自定义均匀分块(最可靠)
手动把任务拆分成大小尽可能均匀的子列表,确保每个进程处理的任务数差异不超过1,从根源上实现负载均衡。
示例代码:
import multiprocessing as mp import time def f(x): time.sleep(0.1) print(mp.current_process()) def split_tasks(tasks, num_processes): base_count = len(tasks) // num_processes extra_tasks = len(tasks) % num_processes task_chunks = [] start_idx = 0 for i in range(num_processes): # 前extra_tasks个进程多处理1个任务 chunk_size = base_count + 1 if i < extra_tasks else base_count task_chunks.append(tasks[start_idx:start_idx+chunk_size]) start_idx += chunk_size return task_chunks if __name__ == "__main__": process_num = 2 all_tasks = list(range(17)) # 拆分出均匀的任务块 balanced_chunks = split_tasks(all_tasks, process_num) with mp.Pool(process_num) as p: # 给每个进程分配一块任务 async_results = [p.map_async(f, chunk) for chunk in balanced_chunks] # 等待所有任务完成 for res in async_results: res.wait()
这个例子里,17个任务会被拆成9个和8个两块,两个进程各处理一块,完美均衡。
方法2:用imap_unordered动态分配
如果你的任务不需要保持输出顺序,可以用Pool.imap_unordered替代map。它会在进程完成一个chunk后立即分配下一个,而不是预先把所有块都分配好。即使最后有剩余的小块,也会被分配给先空闲的进程,避免单个进程独揽剩余任务。
示例代码:
import multiprocessing as mp import time def f(x): time.sleep(0.1) print(mp.current_process()) if __name__ == "__main__": with mp.Pool(2) as p: # 迭代获取结果,触发动态分配 for _ in p.imap_unordered(f, range(17), 6): pass
注意:如果任务执行耗时完全一致,可能还是会出现和之前一样的情况,但只要任务耗时有细微差异,动态分配就会让剩余任务流向空闲进程。
方法3:不指定chunk大小(简单快捷)
如果你的任务耗时较短,不需要通过大chunk减少进程间通信开销,可以直接省略chunk参数。默认情况下,Pool.map会自动把任务分成与进程数匹配的块,每个进程处理的任务数会尽量均匀。
示例代码:
import multiprocessing as mp import time def f(x): time.sleep(0.1) print(mp.current_process()) if __name__ == "__main__": with mp.Pool(2) as p: p.map(f, range(17))
此时17个任务会被自动拆成9个和8个,两个进程各处理一块,实现均衡。
内容的提问来源于stack exchange,提问作者user15964

