如何向.arrow/.arrows文件写入多批次数据并实现读取?
Arrow文件多批次读写实现方案
一、多批次写入实现
你现有代码的问题在于VectorSchemaRoot被包裹在try-with-resources中,返回时已被自动关闭,且仅生成单个数据批次。要实现多批次写入,需借助ArrowFileWriter循环写入多个批次,具体实现如下:
核心步骤
- 定义统一Schema(所有批次的Schema必须完全一致)
- 初始化
ArrowFileWriter,绑定输出流与Schema - 循环生成多个数据批次,逐个写入文件
- 完成写入后关闭Writer释放资源
示例代码
import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.IntVector; import org.apache.arrow.vector.VarCharVector; import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.arrow.vector.types.ArrowType; import org.apache.arrow.vector.types.Field; import org.apache.arrow.vector.types.FieldType; import org.apache.arrow.vector.types.pojo.Schema; import org.apache.arrow.vector.ipc.ArrowFileWriter; import java.io.FileOutputStream; import java.io.IOException; import java.util.List; public class ArrowMultiBatchWriter { private static final RootAllocator allocator = new RootAllocator(Long.MAX_VALUE); // 生成单个批次的数据 private static VectorSchemaRoot createBatch(int batchSize, int batchOffset) { Field age = new Field("age", FieldType.nullable(new ArrowType.Int(32, true)), null); Field name = new Field("name", FieldType.nullable(new ArrowType.Utf8()), null); Schema schema = new Schema(List.of(age, name)); VectorSchemaRoot root = VectorSchemaRoot.create(schema, allocator); IntVector ageVector = (IntVector) root.getVector("age"); VarCharVector nameVector = (VarCharVector) root.getVector("name"); ageVector.allocateNew(batchSize); nameVector.allocateNew(batchSize); for (int i = 0; i < batchSize; i++) { int value = batchOffset + i; ageVector.set(i, value * 10); nameVector.set(i, ("John " + value * 10).getBytes()); } root.setRowCount(batchSize); return root; } // 多批次写入文件 public static void writeMultiBatch(String filePath) throws IOException { // 用第一个批次初始化Writer与Schema VectorSchemaRoot firstBatch = createBatch(100, 0); try (FileOutputStream fos = new FileOutputStream(filePath); ArrowFileWriter writer = new ArrowFileWriter(firstBatch, null, fos.getChannel())) { writer.start(); // 写入第一个批次 writer.writeBatch(); firstBatch.close(); // 循环写入后续4个批次 for (int i = 1; i < 5; i++) { VectorSchemaRoot batch = createBatch(100, i * 100); writer.writeBatch(batch); batch.close(); } writer.end(); } } public static void main(String[] args) throws IOException { writeMultiBatch("multi_batch_data.arrow"); allocator.close(); } }
二、多批次读取实现
读取多批次数据时,使用ArrowFileReader循环加载每个批次,直到文件末尾:
示例代码
import org.apache.arrow.memory.RootAllocator; import org.apache.arrow.vector.IntVector; import org.apache.arrow.vector.VarCharVector; import org.apache.arrow.vector.VectorSchemaRoot; import org.apache.arrow.vector.ipc.ArrowFileReader; import java.io.FileInputStream; import java.io.IOException; public class ArrowMultiBatchReader { private static final RootAllocator allocator = new RootAllocator(Long.MAX_VALUE); public static void readMultiBatch(String filePath) throws IOException { try (FileInputStream fis = new FileInputStream(filePath); ArrowFileReader reader = new ArrowFileReader(fis.getChannel(), allocator)) { VectorSchemaRoot root = reader.getVectorSchemaRoot(); // 循环加载下一个批次,直到无数据可读 while (reader.loadNextBatch()) { IntVector ageVector = (IntVector) root.getVector("age"); VarCharVector nameVector = (VarCharVector) root.getVector("name"); System.out.println("读取到批次,行数:" + root.getRowCount()); // 遍历当前批次数据 for (int i = 0; i < root.getRowCount(); i++) { System.out.printf("age: %d, name: %s%n", ageVector.get(i), new String(nameVector.get(i))); } } root.close(); } } public static void main(String[] args) throws IOException { readMultiBatch("multi_batch_data.arrow"); allocator.close(); } }
三、关键概念梳理
- Record Batch(批次):Arrow中数据的基本组织单元,由一组同Schema的列向量组成,包含固定行数的数据,单个Arrow文件可包含多个批次。
- VectorSchemaRoot:代表一个批次的完整数据集合,包含列向量与Schema,读写时每个批次对应一个
VectorSchemaRoot实例(也可复用实例)。 - ArrowFileWriter/Reader:专门处理Arrow文件的I/O组件,原生支持多批次数据的写入与读取。
四、学习资料推荐
- Apache Arrow官方Java文档:重点查看I/O操作章节,详细讲解文件读写的API与执行流程。
- Apache Arrow官方GitHub仓库Java示例:仓库内有大量可运行的示例代码,覆盖多批次读写、内存管理等核心场景。
- 《Apache Arrow in Action》书籍:系统讲解Arrow的核心概念、多语言实现与最佳实践,适合深入学习底层原理。
内容的提问来源于stack exchange,提问作者Don_Quijote
相关产品推荐
相关产品推荐

