Python3.6 Mac环境下,如何将大尺寸多维Numpy数组传入多进程队列?
解决Numpy大数组传入multiprocessing队列失效的问题
我之前在Mac上用Python3.6处理大尺寸numpy数组时,也碰到过一模一样的问题——小数组能正常在队列里传递,一旦换成大张量就直接“罢工”。核心原因是默认的pickle序列化机制对大对象支持不佳,再加上大数组本身内存占用高,直接放进队列会导致内存复制开销爆炸,甚至触发系统隐性的内存限制。
下面给你两个可行的解决方案:
方案一:用共享内存传递(推荐)
这种方式不需要复制整个数组,而是让父子进程共享同一块内存区域,只传递内存句柄和数组的元数据(形状、数据类型),既省内存又彻底避免序列化问题。因为你用的是Python3.6(还没有shared_memory模块),我们可以用multiprocessing.Array结合numpy来实现:
import numpy as np from multiprocessing import JoinableQueue, Process, Array import ctypes def numpy_to_shared(arr): # 匹配numpy dtype对应的ctypes类型 dtype_map = { np.float64: ctypes.c_double, np.float32: ctypes.c_float, np.int32: ctypes.c_int32, np.int64: ctypes.c_int64 } ctype = dtype_map.get(arr.dtype) if not ctype: raise ValueError(f"不支持的dtype类型: {arr.dtype}") # 创建共享内存数组 shared_arr = Array(ctype, arr.size) # 将numpy数组的数据拷贝到共享内存 np_shared = np.frombuffer(shared_arr.get_obj(), dtype=arr.dtype).reshape(arr.shape) np_shared[:] = arr[:] return shared_arr, arr.shape, arr.dtype def shared_to_numpy(shared_arr, shape, dtype): # 从共享内存还原numpy数组 return np.frombuffer(shared_arr.get_obj(), dtype=dtype).reshape(shape) class Writer(Process): def __init__(self, que): super().__init__() self.queue = que def run(self): for i in range(10): # 生成你需要的4D大张量 data = np.random.randn(100, 1, 16, 12000) # 转换为共享内存格式放入队列 shared_data = numpy_to_shared(data) self.queue.put(shared_data) print(f"已发送第{i+1}个数据,形状: {data.shape}") # 发送结束信号 self.queue.put(None) self.queue.join() class Reader(Process): def __init__(self, que): super().__init__() self.queue = que def run(self): while True: item = self.queue.get() if item is None: self.queue.task_done() break # 还原numpy数组 shared_arr, shape, dtype = item data = shared_to_numpy(shared_arr, shape, dtype) print(f"已接收数据,形状: {data.shape}") self.queue.task_done() if __name__ == "__main__": que = JoinableQueue() writer = Writer(que) reader = Reader(que) writer.start() reader.start() writer.join() reader.join()
方案二:升级pickle协议
如果不想改共享内存的逻辑,也可以尝试用更高版本的pickle协议来序列化大数组。Python3.4+支持协议4,对大对象的序列化效率和兼容性更好:
import numpy as np from multiprocessing import JoinableQueue, Process import pickle class Writer(Process): def __init__(self, que): super().__init__() self.queue = que def run(self): for i in range(10): data = np.random.randn(100, 1, 16, 12000) # 用协议4序列化数组 serialized_data = pickle.dumps(data, protocol=4) self.queue.put(serialized_data) print(f"已发送第{i+1}个数据") self.queue.put(None) self.queue.join() class Reader(Process): def __init__(self, que): super().__init__() self.queue = que def run(self): while True: item = self.queue.get() if item is None: self.queue.task_done() break # 反序列化还原数组 data = pickle.loads(item) print(f"已接收数据,形状: {data.shape}") self.queue.task_done() if __name__ == "__main__": que = JoinableQueue() writer = Writer(que) reader = Reader(que) writer.start() reader.start() writer.join() reader.join()
额外注意事项
- 你的4D张量每个大概占146MB(float64类型),10个就是1.46GB,直接传递会重复占用内存,共享内存方案能避免这个问题。
- MacOS的Python3.6默认用
fork启动子进程,但大数组的pickle在fork环境下依然可能有隐性问题,共享内存是更稳妥的选择。
内容的提问来源于stack exchange,提问作者Petr Nejedlý
相关产品推荐
相关产品推荐

