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

如何用Apache Beam将Kafka记录写入Avro文件?

如何使用Apache Beam的AvroIO将GenericRecord写入Avro文件?

问题背景

我已经实现通过Apache Beam从Kafka主题读取GenericRecord记录并打印到控制台的功能,现在想了解如何使用AvroIO将这些GenericRecord写入Avro文件。以下是当前可运行的代码:

public class BeamConsumer {

    public static void main(String[] args) throws IOException {

        PipelineOptions options = PipelineOptionsFactory.create();
        Pipeline pipeline = Pipeline.create(options);
        Schema schema =
                new Schema.Parser()
                        .parse(new File("schema.avsc"));

        PTransform<PBegin, PCollection<KafkaRecord<GenericRecord, GenericRecord>>> input =
                KafkaIO.<GenericRecord, GenericRecord>read()
                        .withBootstrapServers(
                                "${kafkaserveraddress}")
                        .withTopic("my-topic") // use
                        // withTopics(List<String>) to read from multiple topics.
                        .withKeyDeserializer(
                                ConfluentSchemaRegistryDeserializerProvider.of(
                                        "${schemaregistryaddress}",
                                        "schemaregistrysubjectkey"))
                        .withValueDeserializer(
                                ConfluentSchemaRegistryDeserializerProvider.of(
                                        "${schemaregistryaddress}",
                                        "schemaregistrysubjectvalue"))
                        .withConsumerConfigUpdates(
                                ImmutableMap.of(
                                        ConsumerConfig.GROUP_ID_CONFIG,
                                        "my-group-id",
                                        CommonClientConfigs.SECURITY_PROTOCOL_CONFIG,
                                        "SSL",
                                        SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG,
                                        "/truststore.jks",
                                        SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG,
                                        "******",
                                        SslConfigs.SSL_KEYSTORE_LOCATION_CONFIG,
                                        "/keystore.jks",
                                        SslConfigs.SSL_KEYSTORE_PASSWORD_CONFIG,
                                        "*******",
                                        SslConfigs.SSL_KEY_PASSWORD_CONFIG,
                                        "*******",
                                        ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
                                        "latest",
                                        ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG,
                                        600000));

        pipeline
                .apply(input)
                .apply(
                        "ExtractRecord",
                        ParDo.of(
                                new DoFn<
                                        KafkaRecord<GenericRecord, GenericRecord>,
                                        KafkaRecord<GenericRecord, GenericRecord>>() {
                                    @DoFn.ProcessElement
                                    public void processElement(ProcessContext c) {
                                        KafkaRecord<GenericRecord, GenericRecord> record =
                                                (KafkaRecord<GenericRecord, GenericRecord>) c.element();
                                        KV<GenericRecord, GenericRecord> log = record.getKV();
                                        System.out.println("Key Obtained: " + log.getKey());
                                        System.out.println("Value Obtained: " + log.getValue().toString());
                                        c.output(record);
                                    }
                                }));
        //            .apply("WriteToAvro",
        // AvroIO.writeGenericRecords(schema).to("/Users/mjain34/code/avroutils/src/main/resources/file.avro"));
        //
        PipelineResult run = pipeline.run();
        run.waitUntilFinish(Duration.standardSeconds(1000));
    }
}

注:已修改配置信息以保护隐私

修改后的实现代码

以下是调整后可以将Kafka中的GenericRecord写入Avro文件的完整代码:

import org.apache.avro.Schema;
import org.apache.avro.generic.GenericRecord;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.PipelineResult;
import org.apache.beam.sdk.io.AvroIO;
import org.apache.beam.sdk.io.kafka.KafkaIO;
import org.apache.beam.sdk.io.kafka.KafkaRecord;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.PipelineOptionsFactory;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.KV;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.config.CommonClientConfigs;
import org.apache.kafka.common.config.SslConfigs;
import com.google.common.collect.ImmutableMap;
import java.io.File;
import java.io.IOException;
import org.joda.time.Duration;

public class BeamConsumer {

    public static void main(String[] args) throws IOException {

        PipelineOptions options = PipelineOptionsFactory.create();
        Pipeline pipeline = Pipeline.create(options);
        Schema schema =
                new Schema.Parser()
                        .parse(new File("schema.avsc")); // 确保该Schema与Kafka中Value的Schema完全一致

        PTransform<PBegin, PCollection<KafkaRecord<GenericRecord, GenericRecord>>> input =
                KafkaIO.<GenericRecord, GenericRecord>read()
                        .withBootstrapServers("${kafkaserveraddress}")
                        .withTopic("my-topic")
                        .withKeyDeserializer(
                                ConfluentSchemaRegistryDeserializerProvider.of(
                                        "${schemaregistryaddress}",
                                        "schemaregistrysubjectkey"))
                        .withValueDeserializer(
                                ConfluentSchemaRegistryDeserializerProvider.of(
                                        "${schemaregistryaddress}",
                                        "schemaregistrysubjectvalue"))
                        .withConsumerConfigUpdates(
                                ImmutableMap.of(
                                        ConsumerConfig.GROUP_ID_CONFIG, "my-group-id",
                                        CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SSL",
                                        SslConfigs.SSL_TRUSTSTORE_LOCATION_CONFIG, "/truststore.jks",
                                        SslConfigs.SSL_TRUSTSTORE_PASSWORD_CONFIG, "******",
                                        SslConfigs.SSL_KEYSTORE_LOCATION_CONFIG, "/keystore.jks",
                                        SslConfigs.SSL_KEYSTORE_PASSWORD_CONFIG, "*******",
                                        SslConfigs.SSL_KEY_PASSWORD_CONFIG, "*******",
                                        ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest",
                                        ConsumerConfig.REQUEST_TIMEOUT_MS_CONFIG, 600000));

        pipeline
                .apply(input)
                .apply(
                        "ExtractValueRecord",
                        ParDo.of(
                                new DoFn<KafkaRecord<GenericRecord, GenericRecord>, GenericRecord>() {
                                    @ProcessElement
                                    public void processElement(ProcessContext c) {
                                        KafkaRecord<GenericRecord, GenericRecord> record = c.element();
                                        KV<GenericRecord, GenericRecord> log = record.getKV();
                                        // 保留打印日志用于调试
                                        System.out.println("Key Obtained: " + log.getKey());
                                        System.out.println("Value Obtained: " + log.getValue().toString());
                                        // 输出GenericRecord类型的Value,供AvroIO写入
                                        c.output(log.getValue());
                                    }
                                }))
                .apply(
                        "WriteToAvro",
                        AvroIO.writeGenericRecords(schema)
                                .to("/Users/mjain34/code/avroutils/src/main/resources/output") // 指定输出目录
                                .withSuffix(".avro") // 确保文件带有Avro后缀
                                .withNumShards(1)); // 本地调试时强制生成单个文件,分布式运行可移除该配置

        PipelineResult run = pipeline.run();
        run.waitUntilFinish(Duration.standardSeconds(1000));
    }
}

关键修改说明

  • 调整DoFn输出类型:将原DoFn的输出从KafkaRecord<GenericRecord, GenericRecord>改为GenericRecord,直接输出Kafka记录中的Value(如需写入Key可改为log.getKey()),因为AvroIO的writeGenericRecords需要接收PCollection<GenericRecord>类型的输入。
  • 配置AvroIO写入参数:
    • to()指定输出目录而非单个文件,Beam会根据分片策略生成对应文件
    • withSuffix(".avro")确保生成的文件带有正确的Avro格式后缀
    • withNumShards(1)用于本地调试,强制生成单个文件便于查看;分布式运行时可移除该配置,由Beam自动管理分片数量
  • Schema一致性:确保加载的schema.avsc与Kafka中Value的Schema完全匹配,避免写入时出现Schema不兼容的错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 19:05:04