如何使用KafkaAvroDeserializer创建Kafka Streams消费流?
使用KafkaAvroDeserializer创建Kafka Streams流的解决方案
既然你的输入主题键和值都用了KafkaAvroSerializer序列化,而且后续要对接Confluent JDBC Sink连接器(它只认Avro序列化的数据),那咱们得把消费逻辑改成用对应的KafkaAvroDeserializer来处理,具体步骤如下:
1. 配置Kafka Streams的Avro Serdes
首先要在Streams配置里指定Schema Registry地址,并且初始化Avro专用的Serde(这里推荐用SpecificAvroSerde,因为它能绑定你生成的Avro实体类,类型更安全):
import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.StreamsBuilder; import org.apache.kafka.streams.kstream.Produced; import io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde; import java.util.Collections; import java.util.Map; import java.util.Properties; // 初始化Streams配置 Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "avro-processing-app"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-broker:9092"); // 关键:指定Schema Registry的地址 props.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://your-schema-registry:8081"); // 为键和值创建SpecificAvroSerde实例 SpecificAvroSerde<YourKeyAvroClass> keyAvroSerde = new SpecificAvroSerde<>(); SpecificAvroSerde<Document> valueAvroSerde = new SpecificAvroSerde<>(); // 配置Serde,注意键的Serde要把第二个参数设为true Map<String, String> serdeConfig = Collections.singletonMap( AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://your-schema-registry:8081" ); keyAvroSerde.configure(serdeConfig, true); // true表示这是键的Serde valueAvroSerde.configure(serdeConfig, false); // false表示这是值的Serde
2. 创建使用Avro反序列化的KStream
之前你用的是KStream<String, Document>,现在要改成对应Avro类型的泛型,并且在stream()方法里指定咱们配置好的Serdes:
StreamsBuilder builder = new StreamsBuilder(); // 用Avro Serdes消费输入主题 final KStream<YourKeyAvroClass, Document> avroStream = builder.stream( "your-input-topic", Consumed.with(keyAvroSerde, valueAvroSerde) ); // 后续的处理逻辑(比如过滤、转换等) // ... // 最后发送到JDBC Sink连接器监听的主题,直接用同样的Avro Serdes即可 avroStream.to("jdbc-sink-target-topic", Produced.with(keyAvroSerde, valueAvroSerde));
注意事项
- 如果你没有生成具体的Avro实体类,可以用
GenericAvroSerde替代,此时泛型会是KStream<GenericRecord, GenericRecord>,但类型安全会弱一些,需要手动解析字段。 - 确保你的输入主题的键和值的Avro Schema已经在Schema Registry中注册(既然生产时用了
KafkaAvroSerializer,这一步应该已经完成了)。 - JDBC Sink连接器那边也要对应配置Avro Converter,比如在连接器配置里加上:
key.converter=io.confluent.connect.avro.AvroConverter key.converter.schema.registry.url=http://your-schema-registry:8081 value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://your-schema-registry:8081
内容的提问来源于stack exchange,提问作者Arturo Knight
相关产品推荐
相关产品推荐

