PySpark读取ORC去重后输出文件体积异常增大原因排查
问题背景
我有一个存储ORC文件的目录,通过du -sh验证原文件总大小为3.9G,包含10626个约128KB的小文件,且所有文件schema一致。因数据存在大量重复,使用以下PySpark代码执行全列去重:
# script.py import pyspark.sql import pyspark.sql.functions instance_conf = pyspark.SparkConf() instance_conf.set('spark.local.dir', '/pool/spark_temp') instance_conf.set("spark.sql.autoBroadcastJoinThreshold", -1) spark = pyspark.sql.SparkSession.builder.config(conf=instance_conf).getOrCreate() dirname = "/pool/in" dirname_out = "/pool/out" dataframe = spark.read.orc(dirname) dataframe_out = dataframe.dropDuplicates() # 全列去重 dataframe_out.count() dataframe_out.write.mode("overwrite").orc(dirname_out)
提交命令为:
/usr/local/spark-3.4.1-bin-hadoop3/bin/spark-submit --driver-memory 4g --executor-memory 2g --executor-cores 2 script.py
去重已生效:输出数据行数从467271623降至401149568,且id的distinct计数与原数据一致(均为21174972),但输出目录总大小达12G,仅包含201个文件,体积远超原文件。
体积异常增大的原因解析
结合相关技术问题的核心逻辑,导致去重后ORC文件总体积暴涨的原因主要有以下几点:
1. 列存储压缩效率骤降
ORC作为列存储格式,压缩效率高度依赖列内数据的重复度和局部相关性。原目录的大量128KB小ORC文件中,每个文件内部的局部数据重复度高,压缩算法能充分发挥作用;而dropDuplicates()会触发Shuffle操作,数据被全列哈希重新分区,原本集中的重复数据被打散到不同分区,列内的重复模式被破坏,压缩比大幅下降。即使总行数减少,但每行数据的平均存储体积上升,最终总大小反超原文件。
2. Shuffle分区与文件尺寸的影响
dropDuplicates()本质是基于全列的聚合去重,必须通过Shuffle将相同哈希值的数据分到同一分区处理。Spark默认的Shuffle分区数(与输出的201个文件对应)会让每个分区的数据量远大于原128KB的小文件:
- 大文件内的数据分布更分散,列内重复度不足,压缩效果远不如小文件组合;
- 原小文件可能是上游写入时特意控制的小尺寸,配合局部数据重复实现了高压缩比,合并后的大文件失去了这种优势。
3. ORC写入默认配置的适配问题
Spark写入ORC时的默认配置(如stripe大小、压缩算法参数)未针对去重后的数据优化。原小文件的stripe尺寸可能更小,更适配局部重复数据的压缩;而默认的大stripe尺寸在处理分散的数据时,无法有效利用列存储的压缩特性,进一步降低了压缩效率。
内容的提问来源于stack exchange,提问作者user19695124

