如何使用multiprocessing.Pool与Queue实现指定内容打印?
问题描述
我原本通过multiprocessing.Process逐个从共享队列中取出元素并打印,现在想改用multiprocessing.Pool来实现同样的逻辑,但过程中碰到了问题,该怎么调整呢?
原代码如下:
from multiprocessing import Pool import multiprocessing import queue import time def myTask(queue): value = queue.get() print("Process {} Popped {} from the shared Queue".format(multiprocessing.current_process().pid, value)) queue.task_done() def main(): m = multiprocessing.Manager() sharedQueue = m.Queue() sharedQueue.put(2) sharedQueue.put(3) sharedQueue.put(4) process1 = multiprocessing.Process(target=myTask, args=(sharedQueue,)) process1.start() process1.join() process2 = multiprocessing.Process(target=myTask, args=(sharedQueue,)) process2.start() process2.join() process3 = multiprocessing.Process(target=myTask, args=(sharedQueue,)) process3.start() process3.join() if __name__ == '__main__': main() print("Done")
解决方案
用Pool结合共享队列时,有两个核心点需要注意:
- 必须用
multiprocessing.Manager()创建队列,这样才能在Pool的子进程间安全共享; - 要保证任务数量和队列元素匹配,同时等待所有任务完成以及队列的
task_done信号处理完毕。
这里给你两种可行的实现方式:
方式一:使用pool.map(适合元素数量固定的场景)
这种方式代码简洁,适合提前知道队列元素数量的情况:
from multiprocessing import Pool import multiprocessing def myTask(queue): value = queue.get() print(f"Process {multiprocessing.current_process().pid} Popped {value} from the shared Queue") queue.task_done() def main(): # 用with语句自动管理Manager资源 with multiprocessing.Manager() as m: sharedQueue = m.Queue() # 往队列中放入元素 for num in [2, 3, 4]: sharedQueue.put(num) # 创建进程池,进程数和队列元素数量一致 with Pool(processes=3) as pool: # 给每个元素提交一个任务,传入共享队列(队列是引用类型,所有任务用同一个队列) pool.map(myTask, [sharedQueue]*3) # 等待队列中所有取出的元素都完成task_done sharedQueue.join() print("Done") if __name__ == '__main__': main()
方式二:使用pool.apply_async(更灵活,适合动态元素数量)
如果队列元素数量不确定或者动态生成,用异步提交任务的方式更灵活:
from multiprocessing import Pool import multiprocessing def myTask(queue): value = queue.get() print(f"Process {multiprocessing.current_process().pid} Popped {value} from the shared Queue") queue.task_done() def main(): with multiprocessing.Manager() as m: sharedQueue = m.Queue() items = [2, 3, 4] for num in items: sharedQueue.put(num) with Pool(processes=3) as pool: # 异步提交所有任务,收集结果对象 task_results = [pool.apply_async(myTask, args=(sharedQueue,)) for _ in items] # 等待每个任务执行完成 for res in task_results: res.get() sharedQueue.join() print("Done") if __name__ == '__main__': main()
关键说明
- 用
with语句管理Manager和Pool,会自动处理资源释放,不用手动调用close()和join(); - 共享队列必须通过
Manager创建,普通的multiprocessing.Queue无法在Pool的进程间共享; sharedQueue.join()是必要的,它会阻塞主进程,直到所有从队列取出的元素都调用了task_done(),确保任务真正执行完毕。
内容的提问来源于stack exchange,提问作者LeeChanghwa
相关产品推荐
相关产品推荐

