如何让Python multiprocessing始终固定占用8核运行并行计算任务
解决方案
方案1:使用multiprocessing.Pool(推荐,最简实现)
这是标准库内置的进程池实现,初始化时指定固定进程数后,会自动维持进程负载,空闲进程自动领取新任务,完全避免批处理等待的空窗问题,代码改动量最小。
import random from multiprocessing import Pool # 你的原有计算逻辑无需修改 def do_heavy_computation(num, idx): # 此处保留你原有的高负载计算逻辑 result = num * idx return result if __name__ == '__main__': random_numbers_list = [random.random()] * 10000000 # 初始化8进程的固定进程池 with Pool(processes=8) as pool: # 构造所有任务的参数列表 tasks = [(random_numbers_list[j], j) for j in range(len(random_numbers_list))] # starmap自动将任务分配给空闲进程,全程维持8个进程满负载 results = pool.starmap(do_heavy_computation, tasks)
如果不需要一次性获取所有返回结果,也可以用imap/imap_unordered方法迭代获取结果,内存占用更低,适合超大量任务场景。
方案2:使用concurrent.futures.ProcessPoolExecutor(更现代的API)
Python 3.2+引入的高层并发API,用法和线程池完全一致,切换并发模型更方便:
import random from concurrent.futures import ProcessPoolExecutor def do_heavy_computation(num, idx): result = num * idx return result if __name__ == '__main__': random_numbers_list = [random.random()] * 10000000 with ProcessPoolExecutor(max_workers=8) as executor: # 提交所有任务,executor自动调度给空闲进程 futures = [executor.submit(do_heavy_computation, random_numbers_list[j], j) for j in range(len(random_numbers_list))] # 迭代获取处理结果 for future in futures: result = future.result() # 此处写结果处理逻辑
方案3:手动用Queue实现任务调度(适合自定义调度场景)
如果有任务优先级、动态增减任务等特殊需求,可以自己实现生产者消费者模型,灵活控制调度逻辑:
import random from multiprocessing import Process, Queue def worker(task_queue, result_queue): # 工作进程循环领取任务,直到收到结束标记 while True: task = task_queue.get() if task is None: break num, idx = task result = do_heavy_computation(num, idx) result_queue.put(result) def do_heavy_computation(num, idx): return num * idx if __name__ == '__main__': random_numbers_list = [random.random()] * 10000000 # 限制队列长度避免内存溢出 task_queue = Queue(maxsize=16) result_queue = Queue() # 启动固定8个工作进程 workers = [Process(target=worker, args=(task_queue, result_queue)) for _ in range(8)] for p in workers: p.start() # 生产者投放任务 for j in range(len(random_numbers_list)): task_queue.put((random_numbers_list[j], j)) # 投放结束标记,每个进程对应一个 for _ in range(8): task_queue.put(None) # 等待所有进程执行完毕 for p in workers: p.join() # 处理所有返回结果 while not result_queue.empty(): res = result_queue.get() # 此处写结果处理逻辑
以上三种方案都可以完全避免原批处理逻辑的等待空窗问题,始终保持8个核心满负载运行,普通场景优先选择方案1即可。
内容的提问来源于stack exchange,提问作者Fernando Swenson
相关产品推荐
相关产品推荐

