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

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、行组大小等参数;
  • 分别打印两种写入方式的执行计划:
    # 查看insertInto的执行计划
    df.select(cols).write.insertInto(table_name, overwrite=True).explain()
    # 查看SQL写法的执行计划
    sparkSession.sql(insertSQL.format(table_name, date_str)).explain()
    
    对比两者的shuffle、排序、写入步骤差异;
  • 检查Spark的Parquet相关配置,比如spark.sql.parquet.compression.codec,确认两种写入方式的配置是否一致。

内容的提问来源于stack exchange,提问作者Y.Su

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 07:58:32