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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 19:59:54