如何用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
相关产品推荐
相关产品推荐

