Java中基于Apache Arrow写入Parquet文件的方法及替代方案咨询
将Arrow数据集写入Parquet文件的实现方案
一、基于现有Arrow数据集生成Parquet文件
你已经构建了Arrow的VectorSchemaRoot数据集,只需引入Arrow与Parquet的桥接依赖,替换写入逻辑即可生成Parquet文件。
1. 必要依赖
以Maven为例,添加arrow-parquet依赖(需与你使用的Arrow版本保持一致):
<dependency> <groupId>org.apache.arrow</groupId> <artifactId>arrow-parquet</artifactId> <version>${arrow.version}</version> </dependency>
2. 修改后的完整代码
将原代码中写入Arrow文件的部分替换为Parquet写入逻辑:
try (BufferAllocator allocator = new RootAllocator()) { Field name = new Field("name", FieldType.nullable(new ArrowType.Utf8()), null); Field age = new Field("age", FieldType.nullable(new ArrowType.Int(32, true)), null); Schema schemaPerson = new Schema(asList(name, age)); try (VectorSchemaRoot vectorSchemaRoot = VectorSchemaRoot.create(schemaPerson, allocator)) { VarCharVector nameVector = (VarCharVector) vectorSchemaRoot.getVector("name"); nameVector.allocateNew(3); nameVector.set(0, "David".getBytes()); nameVector.set(1, "Gladis".getBytes()); nameVector.set(2, "Juan".getBytes()); IntVector ageVector = (IntVector) vectorSchemaRoot.getVector("age"); ageVector.allocateNew(3); ageVector.set(0, 10); ageVector.set(1, 20); ageVector.set(2, 30); vectorSchemaRoot.setRowCount(3); // 写入Parquet文件 File parquetFile = new File("person.parquet"); try (FileOutputStream fos = new FileOutputStream(parquetFile); ParquetWriter<VectorSchemaRoot> parquetWriter = ParquetArrowWriter.builder() .withSchema(vectorSchemaRoot.getSchema()) .withWriteMode(ParquetFileWriter.Mode.OVERWRITE) .withAllocator(allocator) .build(fos.getChannel())) { parquetWriter.write(vectorSchemaRoot); System.out.println("成功写入Parquet文件,行数:" + vectorSchemaRoot.getRowCount()); } catch (IOException e) { e.printStackTrace(); } } }
二、两种生成Parquet文件的方式说明
1. 借助Apache Arrow生成
- 适合已有Arrow数据集的场景,无需重新定义Parquet Schema,直接通过桥接工具完成格式转换,代码复用性高。
- 底层基于Apache Parquet API实现,仅做格式适配,无需额外数据映射。
2. 直接使用Apache Parquet原生API生成
- 无需依赖Arrow,可直接定义Parquet Schema,将原始数据序列化写入文件。
- 示例代码:
// 定义Parquet Schema MessageType parquetSchema = MessageParser.parseMessageType( "message person { " + " required binary name (UTF8); " + " required int32 age; " + "}"); // 配置并写入Parquet文件 Path parquetPath = new Path("person_direct.parquet"); try (ParquetWriter<Group> writer = AvroParquetWriter.<Group>builder(parquetPath) .withSchema(parquetSchema) .withWriteMode(ParquetFileWriter.Mode.OVERWRITE) .withConf(new Configuration()) .build()) { // 构造单条数据并写入 Group person1 = new Group(parquetSchema) .append("name", "David") .append("age", 10); writer.write(person1); Group person2 = new Group(parquetSchema) .append("name", "Gladis") .append("age", 20); writer.write(person2); Group person3 = new Group(parquetSchema) .append("name", "Juan") .append("age", 30); writer.write(person3); System.out.println("直接写入Parquet文件完成"); } catch (IOException e) { e.printStackTrace(); }
- 该方式适合无Arrow数据集、直接从原始数据生成Parquet的场景,需手动完成数据到Parquet Group结构的映射。
内容的提问来源于stack exchange,提问作者jake wong
相关产品推荐
相关产品推荐

