使用Apache Beam将BigQuery表数据发送为Kafka Avro消息
将BigQuery TableRow 转换为 Avro 格式并写入 Kafka
步骤1:定义Avro Schema或实体类
先根据BigQuery表结构,创建匹配的Avro Schema(或用Avro工具生成Java实体类)。示例如下:
假设BigQuery表包含id(INT64)、name(STRING)、create_time(TIMESTAMP)字段,Avro Schema可写为:
{ "type": "record", "name": "User", "namespace": "com.example", "fields": [ {"name": "id", "type": "long"}, {"name": "name", "type": "string"}, {"name": "create_time", "type": {"type": "long", "logicalType": "timestamp-millis"}} ] }
若用Java实体类,可通过Avro注解生成:
package com.example; import org.apache.avro.reflect.AvroName; import org.apache.avro.reflect.AvroSchema; @AvroName("User") @AvroSchema("{\"type\":\"record\",\"name\":\"User\",\"namespace\":\"com.example\",\"fields\":[{\"name\":\"id\",\"type\":\"long\"},{\"name\":\"name\",\"type\":\"string\"},{\"name\":\"create_time\",\"type\":{\"type\":\"long\",\"logicalType\":\"timestamp-millis\"}}]}") public class User { private long id; private String name; private long createTime; // 构造函数、getter、setter方法 }
步骤2:编写TableRow到Avro记录的转换逻辑
在Beam Pipeline中添加ParDo转换,实现TableRow到Avro格式的映射:
方式一:转换为Avro GenericRecord
import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.values.TableRow; import java.io.File; // 加载Avro Schema文件 Schema avroSchema = new Schema.Parser().parse(new File("user.avsc")); PCollection<GenericRecord> avroRecords = rows.apply(ParDo.of(new DoFn<TableRow, GenericRecord>() { @ProcessElement public void processElement(ProcessContext c) { TableRow row = c.element(); GenericRecord record = new GenericData.Record(avroSchema); // 逐个字段映射,注意类型转换 record.put("id", row.getLong("id")); record.put("name", row.getString("name")); // BigQuery TIMESTAMP为微秒级,转Avro毫秒级时间戳 long timestampMicros = row.getLong("create_time"); record.put("create_time", timestampMicros / 1000); c.output(record); } }));
方式二:转换为自定义Avro实体类
import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.values.TableRow; PCollection<User> avroUsers = rows.apply(ParDo.of(new DoFn<TableRow, User>() { @ProcessElement public void processElement(ProcessContext c) { TableRow row = c.element(); User user = new User(); user.setId(row.getLong("id")); user.setName(row.getString("name")); long timestampMicros = row.getLong("create_time"); user.setCreateTime(timestampMicros / 1000); c.output(user); } }));
步骤3:配置KafkaIO写入Avro数据
根据转换后的类型,配置KafkaIO的序列化规则:
针对GenericRecord
使用Confluent的KafkaAvroSerializer,需指定Schema Registry地址:
import io.confluent.kafka.serializers.KafkaAvroSerializer; import org.apache.kafka.common.serialization.LongSerializer; avroRecords.apply(KafkaIO.<Long, GenericRecord>write() .withBootstrapServers("kafka:29092") .withTopic("test") .withKeySerializer(LongSerializer.class) .withValueSerializer(KafkaAvroSerializer.class) .withProducerConfigProperty("schema.registry.url", "http://schema-registry:8081") );
Key可根据业务选择,比如用BigQuery表的主键。
针对自定义实体类
同样使用KafkaAvroSerializer,配置Schema Registry地址即可:
import io.confluent.kafka.serializers.KafkaAvroSerializer; import org.apache.kafka.common.serialization.LongSerializer; avroUsers.apply(KafkaIO.<Long, User>write() .withBootstrapServers("kafka:29092") .withTopic("test") .withKeySerializer(LongSerializer.class) .withValueSerializer(KafkaAvroSerializer.class) .withProducerConfigProperty("schema.registry.url", "http://schema-registry:8081") );
注意事项
- 类型匹配:确保BigQuery与Avro的类型对应,比如BigQuery
INT64对应Avrolong,TIMESTAMP(微秒)按需转为Avro的timestamp-millis或timestamp-micros。 - 空值处理:若BigQuery字段允许为空,Avro Schema需设置可空类型(如
["null", "string"]),转换时需处理null场景。 - Schema Registry:使用
KafkaAvroSerializer必须配置schema.registry.url,否则无法自动注册或获取Schema。
内容的提问来源于stack exchange,提问作者vkt
相关产品推荐
相关产品推荐

