能否使用Hadoop API将大型Parquet文件拆分为小文件
核心结论
- 你之前的代码未生效的原因:
FileSplit是MapReduce等计算框架用于标记输入分片的逻辑对象,本身不包含任何对存储层文件的切割、写入操作,仅实例化这个对象不会对实际文件产生任何修改。 - 直接按字节偏移量切割Parquet文件的方案完全不可行:Parquet是自描述列式存储格式,文件尾部存储全文件的元数据索引,内部按行组、数据页组织存储,数据块默认采用压缩编码,直接切字节会破坏压缩块结构、丢失元数据,生成的文件无法被正常解析。
- 无Spark环境下完全可以实现大Parquet文件拆分:仅需引入Parquet官方原生Java依赖,配合Hadoop FileSystem API,通过流式读取+滚动写入的方式即可完成拆分,不需要引入Spark等重型计算组件。
实现步骤
- 首先读取源Parquet文件的Schema、压缩编码配置,保证拆分后的文件和源文件配置完全一致,避免兼容性问题
- 以流式方式逐条读取源文件记录,不要将全量数据加载到内存,避免大文件场景下出现OOM
- 提前设定拆分阈值(可按记录数/文件大小设置,建议单文件大小不小于128MB,对齐Parquet默认行组大小)
- 写入数据量达到阈值时,关闭当前的Parquet写入器,创建新的小文件继续写入,直到源文件全部读取完成
核心实现代码
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.parquet.avro.AvroParquetReader; import org.apache.parquet.avro.AvroParquetWriter; import org.apache.parquet.hadoop.ParquetFileReader; import org.apache.parquet.hadoop.ParquetReader; import org.apache.parquet.hadoop.ParquetWriter; import org.apache.parquet.hadoop.metadata.CompressionCodecName; import org.apache.parquet.hadoop.metadata.ParquetMetadata; import org.apache.avro.generic.GenericRecord; import java.io.IOException; public class ParquetSplitter { // 可根据单条记录大小调整单文件记录数阈值,对应单文件大小约1GB private static final long MAX_RECORDS_PER_FILE = 20_000_000; public static void splitParquetFile(String sourcePathStr, String targetDir, String targetFilePrefix) throws IOException { Configuration conf = new Configuration(); Path sourcePath = new Path(sourcePathStr); // 读取源文件元数据,获取压缩配置 ParquetMetadata sourceMeta = ParquetFileReader.readFooter(conf, sourcePath); CompressionCodecName sourceCodec = sourceMeta.getBlocks().get(0).getColumns().get(0).getCodec(); // 初始化流式读取器,不加载全量数据到内存 try (ParquetReader<GenericRecord> reader = AvroParquetReader.<GenericRecord>builder(sourcePath) .withConf(conf) .build()) { GenericRecord currentRecord; long currentRecordCount = 0; int fileIndex = 0; ParquetWriter<GenericRecord> currentWriter = null; while ((currentRecord = reader.read()) != null) { // 达到单文件阈值时,关闭旧写入流,创建新的目标文件 if (currentWriter == null || currentRecordCount >= MAX_RECORDS_PER_FILE) { if (currentWriter != null) { currentWriter.close(); } Path targetPath = new Path(targetDir + "/" + targetFilePrefix + "_" + fileIndex + ".parquet"); currentWriter = AvroParquetWriter.<GenericRecord>builder(targetPath) .withSchema(currentRecord.getSchema()) .withCompressionCodec(sourceCodec) .withConf(conf) .withRowGroupSize(ParquetWriter.DEFAULT_BLOCK_SIZE) .build(); fileIndex++; currentRecordCount = 0; } currentWriter.write(currentRecord); currentRecordCount++; } // 关闭最后一个文件的写入流 if (currentWriter != null) { currentWriter.close(); } } } public static void main(String[] args) throws IOException { // 入参1为源Parquet路径,入参2为拆分后文件存放目录 splitParquetFile(args[0], args[1], "split_part"); } }
注意事项
- 处理100GB级文件时,建议将任务部署在HDFS集群节点上运行,优先走本地短路径读取数据,降低跨网络传输的带宽开销
- 不要将拆分后单文件设置过小(比如小于64MB),大量小文件会增加HDFS NameNode的元数据压力,同时降低Parquet后续查询效率
- 运行前确保引入的Parquet依赖版本和源文件生成的Parquet版本兼容,避免出现Schema解析错误
内容的提问来源于stack exchange,提问作者Shankar Kumar
相关产品推荐
相关产品推荐

