You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.23 01:50:13