PySpark加载Parquet:StructType字段缺失致整列NULL的问题咨询
我有一个存储在Parquet文件中的data列,该列内容示例为:"{\"col1\":123,\"col2\":123,\"col3\":13,\"col4\":565.0, \"col5\":565.0}"。当数据中缺失col5字段时,把数据加载到DataFrame后,整个data列会变成NULL,哪怕col1、col2、col3和col4都有有效数据。
执行以下代码:
df = spark.read.parquet("/path/to/parquet_file") df.printSchema() df.show()
输出的Schema为:
|-- data: struct (nullable = true)
| |-- col1: long (nullable = true)
| |-- col2: long (nullable = true)
| |-- col3: long (nullable = true)
| |-- col4: double (nullable = true)
| |-- col5: string (nullable = true)
我尝试指定schema把data列转为string类型,但一直遇到类型转换错误:
df = spark.read.schema(incomingSchema).parquet("/path/to/parquet_file")
这个问题的核心是Parquet文件里的data列被存成了struct类型,而Parquet的struct类型要求所有字段必须存在,只要某条记录缺了struct里的任意字段,Spark就会把整个struct标记成NULL。要保留存在的字段、忽略缺失的字段,得换个方式读取数据:
方法1:先读二进制再转字符串解析JSON
from pyspark.sql.functions import from_json, col from pyspark.sql.types import StructType, StructField, LongType, DoubleType, StringType # 定义最终需要的字段schema target_schema = StructType([ StructField("col1", LongType(), nullable=True), StructField("col2", LongType(), nullable=True), StructField("col3", LongType(), nullable=True), StructField("col4", DoubleType(), nullable=True), StructField("col5", StringType(), nullable=True) ]) # 先把data列转成字符串,再解析成JSON结构,最后展开字段 df = spark.read.parquet("/path/to/parquet_file") \ .withColumn("data_str", col("data").cast("string")) \ .withColumn("data_parsed", from_json(col("data_str"), target_schema)) \ .drop("data", "data_str") \ .select("data_parsed.*") df.printSchema() df.show()
方法2:强制指定字符串schema读取
如果直接指定字符串schema报错,大概率是因为Parquet文件的元数据已经把data列标记成了struct类型。可以关闭schema合并参数,再强制按字符串类型读取:
# 关闭schema合并,避免Spark自动沿用文件元数据的struct类型 spark.conf.set("spark.sql.parquet.mergeSchema", "false") # 定义schema,把data列设为字符串类型 incoming_schema = StructType([ StructField("data", StringType(), nullable=True) ]) df = spark.read.schema(incoming_schema).parquet("/path/to/parquet_file") # 解析JSON并展开字段 df = df.withColumn("data_parsed", from_json(col("data"), target_schema)) \ .drop("data") \ .select("data_parsed.*")
关键提示
- Parquet是强类型列式存储,struct类型对字段完整性要求严格,缺字段就会导致整个struct为NULL;
- 先转字符串再用
from_json解析,能灵活处理缺失字段,只保留存在的字段值; - 如果Parquet元数据已经固定了
data的类型,直接改schema会冲突,关闭合并schema或者先读二进制是可行的解决办法。
内容的提问来源于stack exchange,提问作者user3364894

