Spark:按Key重分区输出,咨询控制分区输出文件数量的方法
碰到这种小文件爆炸的情况真的头疼,我之前处理PB级数据时也踩过类似的坑,先给你说说大量小文件的负面影响,再分享几个靠谱的解决办法。
首先:大量小文件的负面影响
- 存储元数据压力大:HDFS的NameNode会把所有文件的元数据存在内存里,29万个文件会占用大量内存,严重时会拖垮NameNode,甚至导致集群不稳定。
- 后续读取性能极差:之后再读取这些数据时,Spark需要为每个小文件创建一个Task,调度这些Task的开销远大于实际计算的开销,跑起来慢到离谱。
- 集群资源浪费:每个小文件的写入都会占用一个Executor Slot,很多时候文件刚写完,Slot还没充分利用就被释放了,完全是浪费集群资源。
控制每个分区文件数量的解决方案
1. 从源头调整Shuffle分区数
Spark默认的spark.sql.shuffle.partitions是200,处理600GB数据时这个值太小了,会导致上游生成大量小分区,后续partitionBy后自然会分裂成更多小文件。你可以根据目标文件大小来计算合适的分区数:
比如希望每个输出文件大概128MB(和HDFS块大小匹配),600GB的数据大概需要 600*1024/128 = 4800 个Shuffle分区。设置方式:
spark.conf.set("spark.sql.shuffle.partitions", 4800)
这个调整要放在数据处理的最开始,让整个计算过程的分区数更合理。
2. 写入前按分区键合并分区
在write之前,用repartition指定每个(foo, bar)分区的文件数,这样可以直接控制每个分区的输出文件数量:
spark.createDataFrame(asRow, struct) // 每个(foo, bar)分区生成10个文件,可根据你的目标文件大小调整这个数字 .repartition($"foo", $"bar", 10) .write .partitionBy("foo", "bar") .format("text") .save("/some/output-path")
这里用repartition会触发Shuffle,但能精准控制每个分区的文件数;如果你的上游数据已经有较大的分区,也可以用coalesce来减少分区数(不会触发Shuffle,但只能减少不能增加)。
3. 用Spark参数动态控制文件大小
针对text格式(不像Parquet有自动合并机制),可以设置Spark的参数来限制每个文件的记录数或大小:
// 每个文件最多写入100万条记录,根据你的单条记录大小调整 spark.conf.set("spark.sql.files.maxRecordsPerFile", 1000000) // 调整文件打开的成本阈值,让Spark更倾向于合并小文件写入 spark.conf.set("spark.sql.files.openCostInBytes", 1024*1024*100) // 使用新版本的OutputCommitter,减少小文件生成 spark.conf.set("mapreduce.fileoutputcommitter.algorithm.version", 2)
这些参数可以配合前面的方法一起用,进一步优化文件数量。
4. 事后合并已生成的小文件
如果已经生成了大量小文件,也可以事后处理:用Spark重新读取数据,按分区合并后再写入:
spark.read.text("/some/output-path") .repartition($"foo", $"bar", 5) // 每个分区保留5个文件 .write .mode("overwrite") .partitionBy("foo", "bar") .format("text") .save("/some/output-path")
这种方式适合补救,但效率不如从源头控制,所以优先推荐前面的方法。
总结
最好的方式是提前规划:先根据数据量调整Shuffle分区数,再结合repartition或maxRecordsPerFile来控制每个分区的文件数,从源头避免小文件生成,比事后处理效率高很多。
内容的提问来源于stack exchange,提问作者minyo

