PySpark读写Parquet文件后Schema不一致的原因及解决方法
问题原因
Spark读取Parquet目录时,默认会合并所有子文件/分区的Schema:当目录下存在Schema不一致的文件时,会生成包含所有字段的超集Schema。你一开始读取的是整个original_folder/目录,里面包含1月(旧Schema)和5月(新Schema)的文件,Spark自动合并两者Schema后生成了一个包含所有字段的DataFrame。筛选1月1日的数据时,那些只在5月Schema里存在的字段会被填充为null,最终df1的Schema是合并后的超集,而直接读取单个分区的df_orig用的是该分区原生的旧Schema,所以两者Schema不一致。
解决方案
针对你的需求,有几种可靠的处理方式:
1. 直接读取目标分区(最推荐)
跳过读取整个目录的步骤,直接读取指定的1月1日0点分区,Spark会直接使用该分区的原生Schema,处理后写入的文件Schema也会和原文件一致:
# 直接读取目标分区 df1 = spark.read.parquet("original_folder/datestr=20240101/hourstr=0/") # 调整分区并写入 df1.coalesce(80).write.mode("append").partitionBy("datestr","hourstr").option("parquet.block.size", 134217728).parquet("new_folder/")
2. 强制使用目标Schema读取整个目录
如果需要通过SQL筛选多个分区,但要保证目标分区的Schema正确,可以先获取目标分区的原生Schema,再强制Spark用这个Schema读取整个目录:
# 先获取目标分区的原生Schema target_schema = spark.read.parquet("original_folder/datestr=20240101/hourstr=0/").schema # 禁用Schema合并,强制使用目标Schema读取整个目录 df = spark.read.option("mergeSchema", "false").schema(target_schema).parquet("original_folder/") df.createOrReplaceTempView("all_records") # 筛选目标数据 df1 = spark.sql("select * from all_records where datestr='20240101' and hourstr = '0'") # 写入新文件夹 df1.coalesce(80).write.mode("append").partitionBy("datestr","hourstr").option("parquet.block.size", 134217728).parquet("new_folder/")
这种方式下,Spark会按目标Schema读取所有文件,不符合Schema的字段会被忽略或填充为null,但目标分区的数据会完整保留原Schema。
3. 分批次处理不同Schema的分区
如果需要处理多个不同Schema的分区,可以按Schema分组,分别读取和处理:
# 定义不同Schema对应的分区范围 old_schema_partitions = ["20240101/hourstr=0", "20240102/hourstr=1"] # 其他旧Schema分区 new_schema_partitions = ["20240501/hourstr=0", "20240502/hourstr=2"] # 新Schema分区 # 处理旧Schema分区 for partition in old_schema_partitions: date_part, hour_part = partition.split('/') df = spark.read.parquet(f"original_folder/datestr={date_part}/{hour_part}/") df.coalesce(80).write.mode("append").partitionBy("datestr","hourstr").option("parquet.block.size", 134217728).parquet("new_folder/") # 处理新Schema分区 for partition in new_schema_partitions: date_part, hour_part = partition.split('/') df = spark.read.parquet(f"original_folder/datestr={date_part}/{hour_part}/") df.coalesce(80).write.mode("append").partitionBy("datestr","hourstr").option("parquet.block.size", 134217728).parquet("new_folder/")
内容的提问来源于stack exchange,提问作者Rayne
相关产品推荐
相关产品推荐

