PySpark RDD中NumPy ndarray子类自定义属性丢失问题求助
PySpark在分布式计算中会对RDD中的对象进行**序列化(Worker节点端)和反序列化(Driver端)**操作,默认使用Pickle作为序列化工具。但numpy的ndarray子类的自定义属性(比如你这里的extra)不会被numpy默认的Pickle逻辑自动序列化——numpy在序列化ndarray时,只处理其核心的数组数据、shape、dtype等内置属性,子类新增的属性会被忽略。
所以你在Worker节点的map操作中创建的MyArray对象确实有extra属性,但序列化发送回Driver端时,这个属性没被保存;反序列化后得到的MyArray对象只是恢复了numpy内置的属性,自定义的extra就丢失了,自然触发AttributeError。
给MyArray类自定义Pickle序列化/反序列化方法
实现__getstate__和__setstate__方法,手动把自定义属性加入序列化字典:import numpy as np class MyArray(np.ndarray): def __new__(cls, input_array, extra=None): obj = np.asarray(input_array).view(cls) obj.extra = extra return obj def __getstate__(self): # 保存numpy数组的状态和自定义属性 state = self.__dict__.copy() state['numpy_state'] = super().__reduce__()[2] return state def __setstate__(self, state): # 先恢复numpy数组的状态 super().__setstate__(state['numpy_state']) # 再恢复自定义属性 self.__dict__.update(state)这样Pickle就能正确处理自定义属性的序列化,Driver端反序列化后就能正常访问
extra。用封装类替代numpy子类
避免继承ndarray,而是创建一个封装类,把数组和自定义属性作为类的成员:class ArrayWrapper: def __init__(self, array, extra=None): self.array = np.asarray(array) self.extra = extra这种方式的序列化逻辑完全由Python的Pickle处理,不会出现属性丢失的问题,缺点是需要修改后续代码中对数组的访问方式(比如从
obj.shape改成obj.array.shape)。单独存储自定义属性
如果不需要把属性和数组绑定在同一个对象里,可以把MyArray的数组部分和extra属性拆分成RDD的二元组,比如map操作返回(my_array, extra),后续处理时直接访问二元组的第二个元素即可。
内容的提问来源于stack exchange,提问作者MikeMayer67

