为何无法将Queue作为参数传入multiprocessing的Pool模块?
问题:为何无法将Queue作为参数传入Pool模块?
无Queue的正常代码示例
创建文件mp_without_queue.py:
from multiprocessing import Pool import time # 补充原代码缺失的time导入 n = 3 def hello(i): print('hello') time.sleep(1) pool = Pool(processes = 2) for i in range(n): pool.apply_async(hello, args=(i,)) pool.close() pool.join()
执行命令及输出:
python3 mp_without_queue.py hello hello hello
添加Queue但无输出的代码示例
创建文件mp_with_queue.py:
from multiprocessing import Pool, Queue import time # 补充原代码缺失的time导入 q = Queue() n = 3 def hello(i, q): print('hello') time.sleep(1) pool = Pool(processes = 2) for i in range(n): pool.apply_async(hello, args=(i, q)) pool.close() pool.join()
执行上述代码后无任何输出,疑问:为何无法将Queue作为参数传入Pool模块?
原因及解决方法
核心原因:
multiprocessing.Queue对象不能直接传递给Pool的子进程。Pool依赖pickle序列化机制传递参数,但标准Queue无法被pickle序列化,子进程启动时会触发隐性错误,而apply_async的异步特性会掩盖这个错误,导致任务根本没执行,所以没有输出。正确方案:使用
multiprocessing.Manager().Queue()替代标准Queue。Manager创建的Queue是代理对象,可安全序列化并在Pool进程间共享。
修改后的可用代码:
from multiprocessing import Pool, Manager import time n = 3 def hello(i, q): print('hello') time.sleep(1) if __name__ == '__main__': # Windows平台下必须添加,避免进程重复初始化 q = Manager().Queue() pool = Pool(processes = 2) for i in range(n): pool.apply_async(hello, args=(i, q)) pool.close() pool.join()
执行后会正常输出三次hello。
额外提示:原代码缺失import time语句,运行时会触发NameError,建议补充;Windows系统下使用Pool必须将主逻辑放在if __name__ == '__main__':代码块中,否则会出现进程创建异常。
内容的提问来源于stack exchange,提问作者showkey
相关产品推荐
相关产品推荐

