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

Java写入Parquet文件:已废弃API的替代方案咨询

Java写入Parquet的废弃API替代方案

当org.apache.parquet.hadoop.ParquetWriter和org.apache.parquet.avro.AvroParquetWriter被废弃后,有以下几种成熟的替代方案:

1. 官方推荐的Builder模式API

官方已将原直接实例化方式替换为构建器模式,通过ParquetWriter.Builder及其子类(如AvroParquetWriter.Builder)创建写入器,这是最直接的替代方案,完全兼容原有业务逻辑。

示例代码(Avro数据写入Parquet):

import org.apache.avro.Schema;
import org.apache.avro.generic.GenericRecord;
import org.apache.parquet.avro.AvroParquetWriter;
import org.apache.parquet.hadoop.ParquetFileWriter;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.conf.Configuration;

// 定义Avro Schema
Schema schema = new Schema.Parser().parse("{\n" +
        "  \"type\": \"record\",\n" +
        "  \"name\": \"User\",\n" +
        "  \"fields\": [\n" +
        "    {\"name\": \"id\", \"type\": \"int\"},\n" +
        "    {\"name\": \"name\", \"type\": \"string\"}\n" +
        "  ]\n" +
        "}");

// 通过Builder创建AvroParquetWriter
try (AvroParquetWriter<GenericRecord> writer = AvroParquetWriter.<GenericRecord>builder(new Path("user_data.parquet"))
        .withSchema(schema)
        .withWriteMode(ParquetFileWriter.Mode.OVERWRITE)
        .withConf(new Configuration())
        .build()) {
    // 构造并写入数据
    GenericRecord user = new org.apache.avro.generic.GenericData.Record(schema);
    user.put("id", 1);
    user.put("name", "Alice");
    writer.write(user);
}

2. Apache Spark DataFrame API

如果是大数据批量处理场景,Spark的DataFrame API提供了极简的Parquet读写能力,无需关注底层Parquet细节,同时支持分区、压缩等高级特性。

示例代码:

import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.StructType;
import org.apache.spark.sql.SaveMode;
import java.util.Arrays;

SparkSession spark = SparkSession.builder()
        .appName("ParquetWriteExample")
        .master("local[*]")
        .getOrCreate();

// 构造数据与Schema
StructType schema = new StructType()
        .add("id", DataTypes.IntegerType)
        .add("name", DataTypes.StringType);
Row row1 = org.apache.spark.sql.RowFactory.create(1, "Alice");
Row row2 = org.apache.spark.sql.RowFactory.create(2, "Bob");

// 创建DataFrame并写入Parquet
spark.createDataFrame(Arrays.asList(row1, row2), schema)
        .write()
        .mode(SaveMode.OVERWRITE)
        .parquet("spark_output.parquet");

spark.stop();

3. Apache Arrow Parquet Writer

对于追求高性能、低延迟的场景,Apache Arrow提供了基于内存列存的Parquet实现,读写效率远超传统API,适合处理大量小文件或实时数据的场景。

示例代码:

import org.apache.arrow.memory.BufferAllocator;
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.pojo.Field;
import org.apache.arrow.vector.types.pojo.Schema;
import org.apache.parquet.arrow.ArrowParquetWriter;
import org.apache.hadoop.fs.Path;
import java.util.List;

try (BufferAllocator allocator = new RootAllocator()) {
    // 构造Arrow Schema
    Schema arrowSchema = new Schema(List.of(
            Field.nullable("id", org.apache.arrow.vector.types.Types.MinorType.INT.getType()),
            Field.nullable("name", org.apache.arrow.vector.types.Types.MinorType.VARCHAR.getType())
    ));

    // 构造数据批次
    try (VectorSchemaRoot root = VectorSchemaRoot.create(arrowSchema, allocator)) {
        IntVector idVector = (IntVector) root.getVector("id");
        VarCharVector nameVector = (VarCharVector) root.getVector("name");

        idVector.allocateNew(2);
        nameVector.allocateNew(2);
        idVector.set(0, 1);
        nameVector.set(0, "Alice".getBytes());
        idVector.set(1, 2);
        nameVector.set(1, "Bob".getBytes());
        root.setRowCount(2);

        // 写入Parquet
        try (ArrowParquetWriter writer = new ArrowParquetWriter(
                new Path("arrow_output.parquet"),
                arrowSchema,
                null,
                null,
                0,
                null,
                new org.apache.hadoop.conf.Configuration())) {
            writer.write(root);
        }
    }
}

如果是流处理场景,Flink提供了专门的Parquet Sink,支持实时写入Parquet文件,同时兼容Avro、Protobuf等序列化格式。

示例代码(Avro流写入):

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.formats.parquet.avro.ParquetAvroWriters;
import org.apache.flink.core.fs.Path;
import org.apache.flink.core.fs.FileSystem;
import org.apache.avro.Schema;

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 定义Avro Schema
Schema schema = new Schema.Parser().parse("{\n" +
        "  \"type\": \"record\",\n" +
        "  \"name\": \"User\",\n" +
        "  \"fields\": [\n" +
        "    {\"name\": \"id\", \"type\": \"int\"},\n" +
        "    {\"name\": \"name\", \"type\": \"string\"}\n" +
        "  ]\n" +
        "}");

// 构造数据流并写入Parquet
env.fromElements(
        new org.apache.avro.generic.GenericData.Record(schema) {{
            put("id", 1);
            put("name", "Alice");
        }},
        new org.apache.avro.generic.GenericData.Record(schema) {{
            put("id", 2);
            put("name", "Bob");
        }}
).sinkTo(ParquetAvroWriters.forGenericRecord(schema)
        .withPath(new Path("flink_output.parquet"))
        .withWriteMode(FileSystem.WriteMode.OVERWRITE)
        .build());

env.execute("Flink Parquet Write");

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 02:55:19