使用Python multiprocessing的ProcessPoolExecutor时队列大小始终为0是什么原因?
问题原因分析
- 核心原因是进程地址空间隔离,你当前使用的
Queue默认是独立进程内存空间内的对象,无法跨主进程和ProcessPoolExecutor启动的子进程共享。
你在主进程定义的全局queue变量,在ProcessPoolExecutor启动子进程时(Python默认spawn启动模式下),子进程会重新导入主模块,生成自己的全新queue实例,和主进程的queue完全无关。子进程调用put写入的是自己进程内的队列,主进程自然读取不到任何数据,所以队列始终为空。 - 补充问题:你的代码中存在拼写错误
ennumerate,正确写法是enumerate,这个错误会直接导致代码运行抛出异常,也是潜在的运行失败原因。
修复方案
方案1:使用跨进程共享的Manager队列
将普通Queue替换为multiprocessing.Manager()创建的共享队列,修改后代码示例:
import os from multiprocessing import Manager from concurrent.futures import ProcessPoolExecutor def produce(i, item, queue): data = process(i, item) queue.put(data) def process(i, item): # 保留原有处理逻辑 data = do_processing(i, item) return data if __name__ == '__main__': # 创建跨进程共享的队列 with Manager() as manager: queue = manager.Queue() records = load_records() with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: print('produce items') # 修正拼写错误,同时把共享队列传入任务 for i, item in enumerate(records.items()): executor.submit(produce, i, item, queue) print('queue size:{}'.format(queue.qsize())) while not queue.empty(): save(queue.get())
方案2:直接收集Future返回结果(更推荐,无需队列)
因为你本来就是等所有任务处理完才统一消费,完全不需要额外用队列,直接收集每个任务的返回结果即可,代码更简洁也没有额外的进程间通信开销:
import os from concurrent.futures import ProcessPoolExecutor def process(i, item): data = do_processing(i, item) return data if __name__ == '__main__': records = load_records() results = [] with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: print('produce items') futures = [executor.submit(process, i, item) for i, item in enumerate(records.items())] # 收集所有任务返回结果 for future in futures: results.append(future.result()) print('总处理结果数:{}'.format(len(results))) for data in results: save(data)
内容的提问来源于stack exchange,提问作者Exploring
相关产品推荐
相关产品推荐

