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

Java版Cloud Dataflow读取Parquet文件至BigQuery的实现咨询

解决Cloud Dataflow Java作业读取Parquet并写入BigQuery的问题

我完全理解你的困扰——Cloud Dataflow Java SDK确实没有直接支持Parquet读取的开箱即用组件,不过基于HadoopFileFormat自定义源是完全可行的,下面我给你详细的实现方案、代码示例,还有一些替代思路。

一、基于HadoopFileFormat的自定义Parquet源实现

这个方案的核心是利用Hadoop的ParquetInputFormat,结合Dataflow的HadoopFileFormat连接器来构建自定义读取源。

1. 先准备依赖

首先在你的Maven pom.xml中添加必要的依赖(注意版本兼容性,建议使用Dataflow的稳定版):

<dependencies>
    <!-- Dataflow核心依赖 -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-core</artifactId>
        <version>2.50.0</version>
    </dependency>
    <!-- Dataflow Hadoop文件系统支持 -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-io-hadoop-file-system</artifactId>
        <version>2.50.0</version>
    </dependency>
    <!-- Parquet Hadoop依赖 -->
    <dependency>
        <groupId>org.apache.parquet</groupId>
        <artifactId>parquet-hadoop</artifactId>
        <version>1.13.0</version>
        <exclusions>
            <!-- 排除冲突的Hadoop组件,使用Dataflow提供的版本 -->
            <exclusion>
                <groupId>org.apache.hadoop</groupId>
                <artifactId>hadoop-common</artifactId>
            </exclusion>
        </exclusions>
    </dependency>
    <!-- Dataflow BigQuery连接器 -->
    <dependency>
        <groupId>org.apache.beam</groupId>
        <artifactId>beam-sdks-java-io-google-cloud-platform</artifactId>
        <version>2.50.0</version>
    </dependency>
</dependencies>

2. 完整代码示例

假设你有一个业务数据类MyData对应Parquet文件中的结构,下面是完整的作业实现:

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.PipelineOptions;
import org.apache.beam.sdk.PipelineOptionsFactory;
import org.apache.beam.sdk.io.BigQueryIO;
import org.apache.beam.sdk.io.hadoop.HadoopFileFormat;
import org.apache.beam.sdk.io.hadoop.SerializableConfiguration;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.transforms.SimpleFunction;
import org.apache.beam.sdk.values.KV;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.io.NullWritable;
import org.apache.parquet.hadoop.ParquetInputFormat;
import org.apache.parquet.hadoop.api.ReadSupport;
import org.apache.parquet.io.api.Binary;
import org.apache.parquet.io.api.GroupConverter;
import org.apache.parquet.io.api.PrimitiveConverter;
import org.apache.parquet.io.api.RecordMaterializer;
import org.apache.parquet.io.api.VoidConverter;
import org.apache.parquet.schema.FileSchema;
import com.google.api.services.bigquery.model.TableFieldSchema;
import com.google.api.services.bigquery.model.TableRow;
import com.google.api.services.bigquery.model.TableSchema;
import java.util.Arrays;
import java.util.List;
import java.util.Map;

// 你的业务数据类,对应Parquet中的字段
public class MyData implements java.io.Serializable {
    private String id;
    private String name;
    private int value;

    // 构造函数、getter和setter
    public MyData() {}
    public String getId() { return id; }
    public void setId(String id) { this.id = id; }
    public String getName() { return name; }
    public void setName(String name) { this.name = name; }
    public int getValue() { return value; }
    public void setValue(int value) { this.value = value; }
}

// 自定义ReadSupport,将Parquet数据映射到MyData对象
public class MyDataReadSupport extends ReadSupport<MyData> implements java.io.Serializable {
    @Override
    public ReadContext init(InitContext context) {
        return new ReadContext(context.getConfiguration());
    }

    @Override
    public RecordMaterializer<MyData> prepareForRead(Configuration conf,
                                                     Map<String, String> keyValueMetaData,
                                                     FileSchema fileSchema,
                                                     ReadContext readContext) {
        return new RecordMaterializer<MyData>() {
            private MyData currentRecord;

            @Override
            public MyData getCurrentRecord() {
                return currentRecord;
            }

            @Override
            public GroupConverter getRootConverter() {
                return new GroupConverter() {
                    @Override
                    public void start() {
                        currentRecord = new MyData();
                    }

                    @Override
                    public void end() {}

                    @Override
                    public GroupConverter getConverter(int fieldIndex) {
                        // 根据Parquet文件的Schema字段名映射到MyData属性
                        String fieldName = fileSchema.getFields().get(fieldIndex).getName();
                        switch (fieldName) {
                            case "id":
                                return new PrimitiveConverter() {
                                    @Override
                                    public void addBinary(Binary value) {
                                        currentRecord.setId(value.toStringUsingUTF8());
                                    }
                                };
                            case "name":
                                return new PrimitiveConverter() {
                                    @Override
                                    public void addBinary(Binary value) {
                                        currentRecord.setName(value.toStringUsingUTF8());
                                    }
                                };
                            case "value":
                                return new PrimitiveConverter() {
                                    @Override
                                    public void addInt(int value) {
                                        currentRecord.setValue(value);
                                    }
                                };
                            default:
                                return new VoidConverter(); // 忽略未匹配的字段
                        }
                    }
                };
            }
        };
    }
}

// Dataflow主作业类
public class ParquetToBigQueryJob {
    public static void main(String[] args) {
        // 初始化Pipeline配置
        PipelineOptions options = PipelineOptionsFactory.fromArgs(args).create();
        Pipeline pipeline = Pipeline.create(options);

        // 配置Hadoop参数
        Configuration hadoopConf = new Configuration();
        // 指定自定义的ReadSupport类
        hadoopConf.setClass(ParquetInputFormat.READ_SUPPORT_CLASS, MyDataReadSupport.class, ReadSupport.class);

        pipeline.apply("读取Parquet文件", HadoopFileFormat.<NullWritable, MyData>read()
                        .withConfiguration(new SerializableConfiguration(hadoopConf))
                        .withInputFormat(ParquetInputFormat.class)
                        .withKeyClass(NullWritable.class)
                        .withValueClass(MyData.class)
                        .withFilePattern("gs://your-bucket/path/to/parquet-files/*.parquet"))
                // 提取KV中的MyData对象
                .apply("提取业务数据", MapElements.via(new SimpleFunction<KV<NullWritable, MyData>, MyData>() {
                    @Override
                    public MyData apply(KV<NullWritable, MyData> input) {
                        return input.getValue();
                    }
                }))
                // 写入BigQuery
                .apply("写入BigQuery", BigQueryIO.<MyData>write()
                        .to("your-gcp-project:your-dataset.your-table")
                        .withSchema(getBigQuerySchema())
                        .withFormatFunction(input -> new TableRow()
                                .set("id", input.getId())
                                .set("name", input.getName())
                                .set("value", input.getValue()))
                        .withWriteDisposition(BigQueryIO.Write.WriteDisposition.WRITE_APPEND)
                        .withCreateDisposition(BigQueryIO.Write.CreateDisposition.CREATE_IF_NEEDED));

        // 启动作业
        pipeline.run().waitUntilFinish();
    }

    // 定义BigQuery表结构
    private static TableSchema getBigQuerySchema() {
        List<TableFieldSchema> fields = Arrays.asList(
                new TableFieldSchema().setName("id").setType("STRING"),
                new TableFieldSchema().setName("name").setType("STRING"),
                new TableFieldSchema().setName("value").setType("INTEGER")
        );
        return new TableSchema().setFields(fields);
    }
}

3. 简化优化:如果Parquet是Avro生成的

如果你的Parquet文件是用Avro Schema生成的,不需要自定义ReadSupport,直接用Avro的内置支持即可:

// 替换Hadoop配置部分
import org.apache.parquet.hadoop.avro.AvroReadSupport;

hadoopConf.setClass(ParquetInputFormat.READ_SUPPORT_CLASS, AvroReadSupport.class, ReadSupport.class);
hadoopConf.set(AvroReadSupport.AVRO_READ_SCHEMA_KEY, MyAvroData.getClassSchema().toString());

这里的MyAvroData是用Avro工具生成的Java类,这样可以省去手写ReadSupport的麻烦。

二、替代方案:使用社区ParquetIO扩展

虽然官方核心SDK没有提供,但Apache Beam的扩展模块中有beam-sdks-java-io-parquet,可以直接用来读取Parquet,代码更简洁:

import org.apache.beam.sdk.io.parquet.ParquetIO;
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericRecord;

// 读取Parquet(假设是Avro格式的Parquet)
pipeline.apply("读取Parquet", ParquetIO.readGenericRecords(new Schema.Parser().parse(avroSchemaString))
                .from("gs://your-bucket/path/to/parquet-files/*.parquet"))
        // 转换为TableRow写入BigQuery
        .apply("转换为TableRow", MapElements.via(new SimpleFunction<GenericRecord, TableRow>() {
            @Override
            public TableRow apply(GenericRecord input) {
                return new TableRow()
                        .set("id", input.get("id").toString())
                        .set("name", input.get("name").toString())
                        .set("value", input.get("value"));
            }
        }))
        .apply("写入BigQuery", BigQueryIO.<TableRow>write()...);

需要添加对应依赖:

<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-parquet</artifactId>
    <version>2.50.0</version>
</dependency>

三、关键注意事项

  • 依赖冲突:Hadoop和Parquet的依赖容易和Dataflow自带的Hadoop组件冲突,一定要通过exclusions排除冲突的包
  • 序列化:所有自定义类(包括MyData、ReadSupport)必须实现Serializable,确保能在Dataflow的分布式环境中序列化传输
  • 性能调优:对于大文件,可以通过Hadoop配置调整分片大小,比如hadoopConf.setInt("mapreduce.input.fileinputformat.split.minsize", 64 * 1024 * 1024);设置64MB分片,提升并行度
  • Schema匹配:确保Parquet的字段类型和BigQuery的Schema完全匹配,或者在转换函数中做类型适配

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 08:33:25