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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:06:02