Databricks中Pandas转Spark DF后collect/存表报PickleException如何解决
报错诱因
该错误的核心原因是你自定义的转换函数仅完成了Spark侧字段类型的声明,没有处理Pandas DataFrame内部的实际值类型:Pandas的数值、布尔等类型默认是numpy原生类型(如np.int64、np.float64、np.bool_),Spark在执行collect()、write等需要全量序列化分区数据的操作时,会用Pickle序列化numpy类型对象,而numpy类型对象的反序列化需要传入构造参数,不符合Spark预期的零参数构造规则,因此抛出序列化异常。display(sdf)能正常展示的原因是该算子仅会拉取少量分区的前若干行数据,在Driver端做轻量转换后直接渲染,不会触发全量数据的跨节点序列化逻辑,因此不会暴露类型问题。
修复方案
方案1:修改自定义转换函数,增加numpy类型转原生Python类型的逻辑
调整后的pandas_to_spark函数参考如下:
import numpy as np from pyspark.sql.types import * def equivalent_type(f): if f == 'datetime64[ns]': return TimestampType() elif f == 'int64': return LongType() elif f == 'int32': return IntegerType() elif f == 'float64': return FloatType() elif f == 'bool': return BooleanType() else: return StringType() def define_structure(string, format_type): try: typo = equivalent_type(format_type) except: typo = StringType() return StructField(string, typo) def pandas_to_spark(pandas_df): columns = list(pandas_df.columns) types = list(pandas_df.dtypes) struct_list = [] for column, typo in zip(columns, types): struct_list.append(define_structure(column, typo)) p_schema = StructType(struct_list) # 新增:将numpy类型值转换为Python原生类型 pandas_df = pandas_df.applymap(lambda x: x.item() if isinstance(x, np.generic) else x) # 新增:将numpy空值替换为Spark可识别的None pandas_df = pandas_df.replace({np.nan: None}) return sqlContext.createDataFrame(pandas_df, p_schema)
方案2:直接使用Spark原生转换能力(更推荐)
Spark 2.3及以上版本(Databricks运行环境默认满足版本要求)已经内置了Pandas DataFrame转Spark DataFrame的自动类型适配逻辑,无需自定义类型映射,直接调用原生方法即可避免该问题:
sdf = spark.createDataFrame(pdf)
内容的提问来源于stack exchange,提问作者Rens
相关产品推荐
相关产品推荐

