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

多进程结果收集:用Queue还是manager.list?兼询SimpleQueue性能

问题背景与疑问

注意:本问题与另一问题核心差异在于任务分发与结果收集的时机。

现有如下代码:

import multiprocessing as MP
from typing import List
import asyncio
import time

WORKER_COUNT = 4  # 示例值

def process_data_worker(job_queue, state_dict, faileds_list):
    # 模拟CPU密集型处理逻辑
    ident = MP.current_process()._identity[0]
    while True:
        item = job_queue.get()
        if item == "DIE":
            state_dict[ident] = "IDLE"
            break
        try:
            # 实际CPU密集型处理逻辑
            pass
        except Exception:
            faileds_list.append(item)
        state_dict[ident] = "IDLE"

async def fetch_data(job_queue):
    # 模拟异步获取数据并推入任务队列
    for i in range(10):
        await asyncio.sleep(0.1)
        job_queue.put(f"chunk_{i}")

if __name__ == "__main__":
    mp_jobqueue = MP.Queue()
    mp_mgr = MP.Manager()
    mp_state = mp_mgr.dict()
    mp_faileds = mp_mgr.list()

    workers: List[MP.Process] = []
    for ident in range(WORKER_COUNT):
        print(ident, end=" ", flush=True)
        mp_state[ident] = None
        w = MP.Process(
            target=process_data_worker,
            args=(mp_jobqueue, mp_state, mp_faileds),
        )
        w.start()
        workers.append(w)

    asyncio.run(fetch_data(mp_jobqueue))

    # 等待任务队列空且所有 worker 处于空闲状态
    safed_workers = 0
    while not mp_jobqueue.empty() or safed_workers < WORKER_COUNT:
        time.sleep(1.0)
        safed_workers = sum(1 for state in mp_state.values() if state == "IDLE")

    # 收集失败结果
    faileds = list(mp_faileds)

    # 关闭管理器
    mp_mgr.shutdown()
    mp_mgr.join()

    # 终止 worker
    [mp_jobqueue.put("DIE") for _ in workers]
    time.sleep(1.0)
    mp_jobqueue.close()
    [w.join() for w in workers]

如上所示,无法使用pool.map()收集失败结果,由此产生以下疑问:

  1. 从性能角度看,将mp_faileds从manager.list替换为Queue是否更优?场景仅需支持「添加元素」和「取出所有元素」操作。
  2. multiprocessing.queues.SimpleQueue的性能是否比普通Queue更优?能否确认这一点?

解答

1. Manager.list 替换为 Queue 的性能优势

是的,替换为Queue(或SimpleQueue)在性能上会显著更优,核心原因:

  • Manager.list依赖中间管理器进程实现跨进程共享,所有增删操作都需要通过管理器转发,额外的IPC(进程间通信)开销在高并发场景下会被放大。
  • Queue是Python多进程模块的原生IPC组件,基于管道或共享内存实现,直接在进程间传递数据,无需中间转发,操作延迟更低、吞吐量更高。

针对你的场景(仅需添加和批量取出元素),Queue完全适配:

  • 工作进程调用queue.put(failed_item)提交失败结果;
  • 主进程在任务全部完成后,通过循环get()结合非阻塞模式捕获Empty异常,一次性收集所有失败项。

2. SimpleQueue 的性能优势确认

SimpleQueue确实比普通Queue性能更优,它是Queue的轻量简化版本:

  • 移除了普通Queue中用于任务跟踪的复杂逻辑(比如task_done()、join()方法),减少了内部锁和状态维护的开销;
  • 仅保留核心的put()和get()操作,在不需要任务同步等待的场景中,能进一步降低IPC延迟。

需要注意的细节:

  • SimpleQueue没有empty()方法(Python 3.7+),主进程收集结果时需通过捕获MP.queues.Empty异常终止循环;
  • 它不支持close()和join_thread()方法,但在你的场景中,任务完成后直接终止工作进程即可,不影响使用。

替换后的核心代码示例

# 替换原有的 mp_faileds 定义
mp_faileds = MP.SimpleQueue()

# 工作进程中提交失败项(替换原 faileds_list.append(item))
mp_faileds.put(item)

# 主进程收集失败结果
faileds = []
while True:
    try:
        item = mp_faileds.get(block=False)
        faileds.append(item)
    except MP.queues.Empty:
        break

内容的提问来源于stack exchange,提问作者pepoluan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 08:30:36