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

使用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对应Avro long,TIMESTAMP(微秒)按需转为Avro的timestamp-millis或timestamp-micros。
  • 空值处理:若BigQuery字段允许为空,Avro Schema需设置可空类型(如["null", "string"]),转换时需处理null场景。
  • Schema Registry:使用KafkaAvroSerializer必须配置schema.registry.url,否则无法自动注册或获取Schema。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 15:05:15