为何无法在ProcessPoolExecutor中使用multiprocessing.Queue?如何解决?
问题描述
运行以下代码:
from concurrent.futures import ProcessPoolExecutor, as_completed from multiprocessing import Queue q = Queue() def my_task(x, queue): queue.put("Task Complete") return x with ProcessPoolExecutor() as executor: tasks = [executor.submit(my_task, i, q) for i in range(10)] for task in as_completed(tasks): print(task.result())
出现错误:
concurrent.futures.process._RemoteTraceback: """ Traceback (most recent call last): File "/usr/lib/python3.10/multiprocessing/queues.py", line 244, in _feed obj = _ForkingPickler.dumps(obj) File "/usr/lib/python3.10/multiprocessing/reduction.py", line 51, in dumps cls(buf, protocol).dump(obj) File "/usr/lib/python3.10/multiprocessing/queues.py", line 58, in __getstate__ context.assert_spawning(self) File "/usr/lib/python3.10/multiprocessing/context.py", line 373, in assert_spawning raise RuntimeError( RuntimeError: Queue objects should only be shared between processes through inheritance """ The above exception was the direct cause of the following exception: Traceback (most recent call last): File "/tmp/nn.py", line 14, in <module> print(task.result()) File "/usr/lib/python3.10/concurrent/futures/_base.py", line 451, in result return self.__get_result() File "/usr/lib/python3.10/concurrent/futures/_base.py", line 403, in __get_result raise self._exception File "/usr/lib/python3.10/multiprocessing/queues.py", line 244, in _feed obj = _ForkingPickler.dumps(obj) File "/usr/lib/python3.10/multiprocessing/reduction.py", line 51, in dumps cls(buf, protocol).dump(obj) File "/usr/lib/python3.10/multiprocessing/queues.py", line 58, in __getstate__ context.assert_spawning(self) File "/usr/lib/python3.10/multiprocessing/context.py", line 373, in assert_spawning raise RuntimeError( RuntimeError: Queue objects should only be shared between processes through inheritance
疑问:
- 既然无法在多进程场景中使用
multiprocessing.Queue,那它的设计用途是什么? - 如何修改代码使其正常运行?
- 实际业务中需要每个工作进程频繁更新队列反馈任务状态,供另一个线程读取更新进度条,该怎么处理?
解答
1. multiprocessing.Queue的设计用途
multiprocessing.Queue是专门用于进程间安全通信的组件,但它的使用有严格限制:只能通过进程继承的方式共享。也就是说,必须在父进程中创建Queue,然后子进程通过fork(Unix系统)或继承(Windows系统下的spawn模式,父进程传递必要资源)的方式直接使用这个Queue,而不能通过参数传递的方式将Queue对象序列化(pickle)后传给子进程——multiprocessing.Queue本身不支持被pickle序列化,这就是报错的核心原因。
正确使用示例(基于multiprocessing.Process):
from multiprocessing import Process, Queue def my_task(x, queue): queue.put(f"Task {x} Complete") if __name__ == "__main__": q = Queue() processes = [Process(target=my_task, args=(i, q)) for i in range(10)] for p in processes: p.start() for p in processes: p.join() # 读取队列内容 while not q.empty(): print(q.get())
2. 修改原代码的两种方法
方法一:使用multiprocessing.Manager().Queue()
Manager创建的Queue是代理对象,底层通过独立服务器进程管理队列,支持被pickle序列化,可安全传递给ProcessPoolExecutor的任务函数:
from concurrent.futures import ProcessPoolExecutor, as_completed from multiprocessing import Manager def my_task(x, queue): queue.put(f"Task {x} Complete") return x if __name__ == "__main__": with Manager() as manager: q = manager.Queue() with ProcessPoolExecutor() as executor: tasks = [executor.submit(my_task, i, q) for i in range(10)] for task in as_completed(tasks): print(task.result()) # 读取状态信息 while not q.empty(): print(q.get())
方法二:改用multiprocessing.Process(适合简单场景)
如果不需要线程池复用能力,直接创建子进程让其继承父进程Queue:
from multiprocessing import Process, Queue def my_task(x, queue): queue.put(f"Task {x} Complete") if __name__ == "__main__": q = Queue() processes = [Process(target=my_task, args=(i, q)) for i in range(10)] for p in processes: p.start() for p in processes: p.join() while not q.empty(): print(q.get())
3. 实际业务中进度条的处理方案
方案一:用Manager().Queue传递状态信息
工作进程向Queue写入阶段状态,主线程启动单独线程读取并更新进度条:
from concurrent.futures import ProcessPoolExecutor, as_completed from multiprocessing import Manager import threading import time from tqdm import tqdm def worker_task(x, queue): time.sleep(0.1) queue.put(1) # 传递进度增量 return x def update_progress(queue, pbar, total): completed = 0 while completed < total: completed += queue.get() pbar.update(1) pbar.close() if __name__ == "__main__": total_tasks = 10 with Manager() as manager: q = manager.Queue() pbar = tqdm(total=total_tasks) # 启动进度更新线程 progress_thread = threading.Thread(target=update_progress, args=(q, pbar, total_tasks)) progress_thread.start() with ProcessPoolExecutor() as executor: tasks = [executor.submit(worker_task, i, q) for i in range(total_tasks)] for task in as_completed(tasks): pass progress_thread.join()
方案二:用共享计数器(更轻量)
若仅需统计完成数,用multiprocessing.Value创建共享计数器,比Queue更高效:
from concurrent.futures import ProcessPoolExecutor, as_completed from multiprocessing import Value, Lock import threading import time from tqdm import tqdm def worker_task(x, counter, lock): time.sleep(0.1) with lock: counter.value += 1 return x def update_progress(counter, pbar, total): while counter.value < total: pbar.n = counter.value pbar.refresh() time.sleep(0.05) pbar.n = total pbar.refresh() pbar.close() if __name__ == "__main__": total_tasks = 10 counter = Value('i', 0) lock = Lock() pbar = tqdm(total=total_tasks) progress_thread = threading.Thread(target=update_progress, args=(counter, pbar, total_tasks)) progress_thread.start() with ProcessPoolExecutor() as executor: tasks = [executor.submit(worker_task, i, counter, lock) for i in range(total_tasks)] for task in as_completed(tasks): pass progress_thread.join()
内容的提问来源于stack exchange,提问作者chrislamp
相关产品推荐
相关产品推荐

