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); } } }
4. Apache Flink Parquet Sink
如果是流处理场景,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
相关产品推荐
相关产品推荐

