多进程结果收集:用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()收集失败结果,由此产生以下疑问:
- 从性能角度看,将
mp_faileds从manager.list替换为Queue是否更优?场景仅需支持「添加元素」和「取出所有元素」操作。 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
相关产品推荐
相关产品推荐

