Spark读写Parquet时如何避免分区产生大量小文件?
解决Spark写入Parquet时大量小文件问题(不降低计算并行度)
首先得明确问题根源:你用coalesce(200)是全局合并分区,但partitionBy(date)会把这200个分区里的数据按date拆分到不同的输出分区中——如果每个全局分区都包含多种date的数据,那每个date分区就会生成200个小文件,这就是问题所在。想要不降低计算阶段的并行度,只在写入阶段针对每个分区合并文件,可以试试下面几种方案:
方案1:使用maxRecordsPerFile参数控制单文件记录数
这是最省心的方法,不需要修改数据处理逻辑,只需要在写入时添加参数,让Spark自动将小文件合并到指定记录数上限:
df.write .option("maxRecordsPerFile", 1000000) // 每个文件最多100万条记录,可根据单条记录大小调整 .partitionBy("date") .parquet("foo")
这个参数会在写入阶段生效:Spark会为每个date分区维护输出流,当某个流的记录数达到上限时,就会关闭当前文件并新建一个。这样既保留了前面计算的高并行度(计算阶段还是用原来的分区数),又能有效减少小文件的数量。
方案2:按分区键+指定分区数重分区(适合提前知道分区键分布的场景)
如果你清楚date的分布情况(比如每天的数据量大致相当),可以在写入前先按date重分区,同时指定全局总分区数:
// 假设你希望全局总共有1000个分区,同一个date的数据会被哈希到固定数量的分区中 df.repartition(col("date"), 1000) .write .partitionBy("date") .parquet("foo")
repartition(col("date"), numPartitions)会按date进行哈希分区,确保同一个date的数据只会落在几个指定的分区里,这样写入时每个date分区的文件数就等于该date对应的重分区数量,不会出现大量小文件。而且重分区是在计算完成后进行的,前面的计算阶段依然保持原来的高并行度。
方案3:使用Spark SQL的动态分区插入(适合Hive表场景)
如果你是要写入Hive兼容的分区表,可以用Spark SQL的方式,配合动态分区参数和文件记录数限制:
SET hive.exec.dynamic.partition=true; SET hive.exec.dynamic.partition.mode=nonstrict; SET spark.sql.files.maxRecordsPerFile=1000000; INSERT OVERWRITE DIRECTORY 'foo' PARTITION (date) SELECT col1, col2, date FROM your_hive_table;
这种方式和方案1类似,但更贴合Hive生态,同样能在不影响计算并行度的前提下控制文件数量。
注意事项
- 不要在计算阶段就用
coalesce或者repartition减少分区数,那样会直接降低并行度,影响计算性能; maxRecordsPerFile的数值建议和HDFS块大小匹配:比如HDFS块是256MB,就按单条记录大小算出对应记录数,避免文件过大或过小;- 如果你的
date分区数据量差异极大(比如某些日期数据量特别小),可以额外加逻辑判断,对小分区单独做coalesce合并,但这种情况属于特例。
内容的提问来源于stack exchange,提问作者Georg Heiler
相关产品推荐
相关产品推荐

