高效合并大Parquet文件:解决parquet-tools内存占用过高问题
解决Parquet合并时内存溢出问题的几种方案
我之前在处理大规模Parquet文件合并时也踩过parquet-tools的内存坑,它默认的合并逻辑确实会把最终输出的整个文件内容都加载到内存里,这在Hadoop容器环境里很容易触发OOM被kill。下面分享几个亲测有效的解决思路:
一、用Spark替代parquet-tools(最推荐)
Spark对Parquet的分布式处理支持非常成熟,它会自动把数据分片到多个executor处理,内存管理也更智能,不会一次性加载所有数据。只需要几行代码就能完成合并:
Scala版本示例:
val spark = SparkSession.builder() .appName("ParquetMerge") .getOrCreate() // 读取所有待合并的Parquet文件 val inputDF = spark.read.parquet("hdfs://your/input/path/*") // 写入合并后的文件,mode根据需求选append/overwrite inputDF.write.mode("overwrite").parquet("hdfs://your/output/path") spark.stop()
PySpark版本示例:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ParquetMerge").getOrCreate() input_df = spark.read.parquet("hdfs://your/input/path/*") input_df.write.mode("overwrite").parquet("hdfs://your/output/path") spark.stop()
提交作业时可以根据数据规模调整executor内存,比如:
spark-submit --executor-memory 4g --driver-memory 2g --num-executors 8 your_merge_script.py
二、调整parquet-tools的JVM与Hadoop容器参数(应急方案)
如果一定要用parquet-tools,你可以通过调整Hadoop Map任务的内存配置来缓解OOM问题。关键是给Map容器分配足够的内存,同时合理设置JVM堆大小(一般是容器内存的70%-80%,留一部分给容器本身的开销):
hadoop jar parquet-tools.jar merge \ -Dmapreduce.map.memory.mb=8192 \ # 给Map容器分配8GB内存 -Dmapreduce.map.java.opts=-Xmx6144m \ # JVM堆设为6GB hdfs://your/input/path/* \ hdfs://your/output/path/merged.parquet
不过这个方法只能处理中等规模的文件合并,当输出文件太大时还是会OOM,因为parquet-tools的底层逻辑没有改变。
三、自定义Parquet合并工具(更灵活)
如果需要更精细的内存控制,可以基于Parquet官方API写一个简单的合并程序,通过设置行组大小(Row Group Size)和页大小(Page Size)来限制内存占用,避免一次性加载所有数据:
Java示例框架:
import org.apache.parquet.avro.AvroParquetReader; import org.apache.parquet.avro.AvroParquetWriter; import org.apache.parquet.hadoop.ParquetReader; import org.apache.parquet.hadoop.ParquetWriter; import org.apache.parquet.hadoop.metadata.CompressionCodecName; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.avro.generic.GenericRecord; import org.apache.avro.Schema; import java.util.List; public class ParquetMerger { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Schema schema = new Schema.Parser().parse(new Path(args[0])); // 从第一个输入文件读取Schema Path outputPath = new Path(args[1]); List<Path> inputPaths = List.of(args).subList(2, args.length).stream().map(Path::new).toList(); // 初始化Writer,设置行组大小为128MB,页大小为1MB,使用Snappy压缩 ParquetWriter<GenericRecord> writer = AvroParquetWriter .<GenericRecord>builder(outputPath) .withSchema(schema) .withConf(conf) .withRowGroupSize(128 * 1024 * 1024) .withPageSize(1 * 1024 * 1024) .withCompressionCodec(CompressionCodecName.SNAPPY) .build(); // 遍历所有输入文件,逐记录写入 for (Path inputPath : inputPaths) { ParquetReader<GenericRecord> reader = AvroParquetReader.<GenericRecord>builder(inputPath).withConf(conf).build(); GenericRecord record; while ((record = reader.read()) != null) { writer.write(record); } reader.close(); } writer.close(); } }
编译打包后,用Hadoop jar命令提交运行,这样内存占用会被控制在你设置的行组大小范围内,不会出现和输出文件大小相当的内存消耗。
四、其他小技巧
- 提前检查待合并的Parquet文件是否有相同的Schema,Schema不一致会导致额外的内存开销甚至合并失败;
- 如果是大量小Parquet文件合并,Spark会自动优化输出文件数量(可以通过
spark.sql.shuffle.partitions调整),避免生成过多小文件。
内容的提问来源于stack exchange,提问作者TLD
相关产品推荐
相关产品推荐

