PySpark写入Parquet遇PickleException:numpy.dtype序列化问题排查
PySpark写入Parquet报PickleException的原因与解决方法
问题根源:完全和数据类型相关。你的
outlets_geolocationDataFrame中存在numpy数据类型(比如numpy.int64、numpy.float64或numpy.dtype对象等)。调用head()时数据被拉取到本地Python环境,numpy类型能正常解析显示;但写入Parquet时,Spark需要对数据做分布式序列化处理,底层使用的razorvine pickle库不支持带参数构造的numpy.dtype对象,因此抛出该异常。解决方法:将DataFrame中的所有numpy类型转换为Spark兼容的类型:
- 自定义匹配区县的UDF时,确保返回值使用原生Python类型(比如把
np.int64(result)改为int(result)) - 对已有列执行类型转换,使用Spark的
cast方法:from pyspark.sql.types import IntegerType, DoubleType, StringType # 示例:将numpy整数类型列转为Spark整数类型 df = outlets_geolocation.withColumn("district_id", outlets_geolocation["district_id"].cast(IntegerType())) # 将numpy浮点类型列转为Spark浮点类型 df = df.withColumn("longitude", df["longitude"].cast(DoubleType())) - 如果是从pandas DataFrame转换而来,指定Spark Schema避免numpy类型被带入:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, DoubleType, StringType # 定义对应Schema schema = StructType([ StructField("outlet_id", StringType(), True), StructField("longitude", DoubleType(), True), StructField("latitude", DoubleType(), True) ]) spark = SparkSession.builder.getOrCreate() # 转换时指定Schema spark_df = spark.createDataFrame(pandas_outlets_df, schema=schema)
- 自定义匹配区县的UDF时,确保返回值使用原生Python类型(比如把
内容的提问来源于stack exchange,提问作者Minura Punchihewa
相关产品推荐
相关产品推荐

