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

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ý

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 04:23:22