Python multiprocessing Manager.Queue多进程未提速问题及优化咨询
多进程处理队列未提效的原因及优化方案
我有一个包含100000个元素的队列,想用多进程逐个处理并统计耗时,实现代码如下:
import multiprocessing from datetime import datetime def worker(queue): while not queue.empty(): try: item = queue.get_nowait() # Process the item print(f"Process {multiprocessing.current_process().name} processing item: {item}") except: break def main(): # Create a queue with 100,000 elements using Manager manager = multiprocessing.Manager() queue = manager.Queue() for i in range(100000): queue.put(i) # Starttime start=datetime.now() # Create 8 worker processes processes = [] for _ in range(NR_WORKERS): p = multiprocessing.Process(target=worker, args=(queue,)) processes.append(p) p.start() # Wait for all processes to finish for p in processes: p.join() # Endtime print(queue.empty(), "all processes finished in:", datetime.now()-start) if __name__ == "__main__": global NR_WORKERS NR_WORKERS = 8 main()
但执行后发现,使用8个进程与1个进程的耗时几乎相同(开启打印时为15-18秒);关闭打印语句后,两种方式耗时均约4秒。为何多进程未带来性能提升?该如何优化?
为什么多进程没带来性能提升?
Manager.Queue的IPC开销与锁竞争:multiprocessing.Manager().Queue()由管理进程托管,每次get/put都依赖RPC通信,开销远大于原生队列。多个进程同时访问时会触发严重锁竞争,大部分时间都在等待队列访问权,完全发挥不出并行优势。- 任务过于轻量:你的处理逻辑几乎没有实际计算(要么空操作要么打印),进程的时间都耗在IPC交互上,多进程的并行计算优势被IPC开销彻底抵消,甚至单进程因为没有竞争反而更快。
- 打印操作的全局阻塞:
print会触发GIL(全局解释器锁)和控制台输出锁,多进程同时打印时会排队等待输出,所有进程都卡在这个环节,此时多进程和单进程的耗时自然无差异。
优化方案
1. 替换为原生multiprocessing.Queue
原生队列基于管道和锁实现,IPC开销远低于Manager.Queue,修改队列创建代码即可:
# 移除Manager相关代码,改用原生队列 queue = multiprocessing.Queue()
2. 分块分配任务,彻底减少IPC次数
把100000个元素分成若干块,每个进程处理一块数据,避免队列竞争:
def worker(task_chunk): for item in task_chunk: # 这里写实际的处理逻辑 pass def main(): items = list(range(100000)) nr_workers = 8 chunk_size = len(items) // nr_workers # 分割任务块,最后一块补全剩余元素 chunks = [items[i:i+chunk_size] for i in range(0, len(items), chunk_size)] if len(chunks) > nr_workers: chunks[-2].extend(chunks[-1]) chunks.pop() start = datetime.now() processes = [] for chunk in chunks: p = multiprocessing.Process(target=worker, args=(chunk,)) processes.append(p) p.start() for p in processes: p.join() print("all processes finished in:", datetime.now()-start) if __name__ == "__main__": main()
3. 优化打印或改用进程安全日志
如果需要输出,不要用print,改用logging模块配置进程安全的Handler:
import logging from multiprocessing import get_logger def worker(task_chunk): logger = get_logger() logger.setLevel(logging.INFO) # 添加进程安全的FileHandler或其他Handler handler = logging.FileHandler("process.log") logger.addHandler(handler) for item in task_chunk: logger.info(f"Process {multiprocessing.current_process().name} processing item: {item}")
4. 确保任务有足够计算量
如果实际处理逻辑还是轻量,多进程本身优势有限,可考虑:
- 合并小任务,减少进程调度和IPC开销
- 改用多线程(适合IO密集型任务,注意GIL对CPU密集型任务的限制)
内容的提问来源于stack exchange,提问作者Kilian Kramer
相关产品推荐
相关产品推荐

