You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何向.arrow/.arrows文件写入多批次数据并实现读取?

Arrow文件多批次读写实现方案

一、多批次写入实现

你现有代码的问题在于VectorSchemaRoot被包裹在try-with-resources中,返回时已被自动关闭,且仅生成单个数据批次。要实现多批次写入,需借助ArrowFileWriter循环写入多个批次,具体实现如下:

核心步骤

  1. 定义统一Schema(所有批次的Schema必须完全一致)
  2. 初始化ArrowFileWriter,绑定输出流与Schema
  3. 循环生成多个数据批次,逐个写入文件
  4. 完成写入后关闭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组件,原生支持多批次数据的写入与读取。

四、学习资料推荐

  1. Apache Arrow官方Java文档:重点查看I/O操作章节,详细讲解文件读写的API与执行流程。
  2. Apache Arrow官方GitHub仓库Java示例:仓库内有大量可运行的示例代码,覆盖多批次读写、内存管理等核心场景。
  3. 《Apache Arrow in Action》书籍:系统讲解Arrow的核心概念、多语言实现与最佳实践,适合深入学习底层原理。

内容的提问来源于stack exchange,提问作者Don_Quijote

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.11 07:50:49