PySpark不同写入方式致Parquet文件大小异常问题求助
为什么两种写入方式会导致Parquet文件大小差异10倍?
答案是肯定的,你遇到的这个Parquet文件大小异常,大概率就是DataFrame.insertInto()和Spark SQL INSERT OVERWRITE这两种写入方式的差异造成的。下面具体拆解两者的核心区别以及对文件大小的影响:
1. 分区写入的逻辑差异
你提到其他日期都是用INSERT OVERWRITE TABLE ... PARTITION(crawled_at_ds = '{}')明确指定了分区,这种写法会直接定位到目标分区进行覆盖写入,Spark SQL会针对单分区场景做针对性优化:
- 自动控制并行度,保证输出文件数量(你这里是16个)和大小均匀;
- 针对分区内的数据做合理的shuffle和分区,让每个输出文件的数据分布更适合Parquet的压缩。
而2019-08-03用的insertInto(table_name, overwrite=True)写法,如果你的DataFrame包含分区列crawled_at_ds,它会触发动态分区写入,但这种模式下:
- 动态分区的并行度控制逻辑和指定分区的SQL写法不同,可能导致数据在输出时的分布不够均匀;
- 如果
overwrite=True是覆盖整个表而非单个分区,Spark会重新扫描全表数据(虽然你这里的DataFrame是当天的数据),但底层的写入流程会更复杂,可能打乱了原本的优化逻辑。
2. Parquet压缩与序列化的差异
Parquet的文件大小很大程度取决于压缩率,而两种写入方式在压缩处理上可能存在差异:
- 压缩Codec默认值不同:Spark SQL的
INSERT OVERWRITE可能默认启用了高效的压缩(比如Snappy或Gzip),而insertInto方法可能因为配置继承问题,使用了更低效的压缩甚至无压缩。你可以通过查看文件的元信息验证这一点; - 行组与页大小的优化:SQL写法会自动适配Parquet的行组大小(默认128MB),让数据在写入时按行组组织,最大化压缩效率;而
insertInto可能因为执行计划的不同,行组大小设置不合理,导致压缩率骤降; - 数据序列化方式:
insertInto需要先对齐表的列顺序(你用了select(cols)),这一步可能改变了数据的序列化格式,而Parquet对不同序列化格式的压缩效率差异很大。
3. 执行计划的优化差异
Spark SQL的INSERT OVERWRITE语句会经过Spark Catalyst优化器的全量优化,比如:
- 自动过滤掉不必要的列(虽然你这里是全列写入,但优化器会做其他调整);
- 对分区数据做预排序,而Parquet对有序数据的压缩率远高于无序数据;
- 避免不必要的shuffle操作,让数据写入更高效。
而DataFrame.insertInto()的执行计划相对简单,很多SQL层的优化不会被触发,比如不会自动对数据做排序,导致Parquet的压缩效果大打折扣——这也是为什么记录数相近,但文件体积差10倍的核心原因之一。
如何验证?
你可以做几个简单的测试来确认:
- 用Parquet工具(比如
parquet-tools)查看异常日期和正常日期的文件元信息,对比压缩Codec、行组大小等参数; - 分别打印两种写入方式的执行计划:
对比两者的shuffle、排序、写入步骤差异;# 查看insertInto的执行计划 df.select(cols).write.insertInto(table_name, overwrite=True).explain() # 查看SQL写法的执行计划 sparkSession.sql(insertSQL.format(table_name, date_str)).explain() - 检查Spark的Parquet相关配置,比如
spark.sql.parquet.compression.codec,确认两种写入方式的配置是否一致。
内容的提问来源于stack exchange,提问作者Y.Su
相关产品推荐
相关产品推荐

