Kafka:为不同主题配置不同的反序列化器
实现不同Kafka主题使用不同反序列化器
完全可以实现,Spring Kafka支持为不同的@KafkaListener指定专属的消费者配置,从而让原有Avro主题继续用KafkaAvroDeserializer,新JSON主题用JsonDeserializer。具体实现步骤如下:
1. 配置多个ConsumerFactory
创建两个独立的ConsumerFactory,分别对应Avro和JSON格式的消息反序列化:
@Configuration public class KafkaConsumerConfig { // Avro消费者配置 @Bean("avroConsumerFactory") public ConsumerFactory<String, SpecificRecord> avroConsumerFactory(KafkaProperties kafkaProperties) { Map<String, Object> configProps = new HashMap<>(kafkaProperties.buildConsumerProperties()); // 配置Avro反序列化器 configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, KafkaAvroDeserializer.class); configProps.put(KafkaAvroDeserializerConfig.SCHEMA_REGISTRY_URL_CONFIG, "你的Schema Registry地址"); configProps.put(KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG, true); return new DefaultKafkaConsumerFactory<>(configProps); } // JSON消费者配置 @Bean("jsonConsumerFactory") public ConsumerFactory<String, YourJsonDto> jsonConsumerFactory(KafkaProperties kafkaProperties) { Map<String, Object> configProps = new HashMap<>(kafkaProperties.buildConsumerProperties()); // 配置JSON反序列化器 configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); // 指定JSON对应的实体类 configProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, YourJsonDto.class.getName()); // 可选:关闭类型校验,允许所有包下的实体类 configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); return new DefaultKafkaConsumerFactory<>(configProps); } }
2. 配置对应的ListenerContainerFactory
为每个ConsumerFactory创建对应的ConcurrentKafkaListenerContainerFactory,供@KafkaListener指定使用:
@Configuration public class KafkaListenerContainerConfig { @Bean("avroListenerContainerFactory") public ConcurrentKafkaListenerContainerFactory<String, SpecificRecord> avroListenerContainerFactory( @Qualifier("avroConsumerFactory") ConsumerFactory<String, SpecificRecord> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, SpecificRecord> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); return factory; } @Bean("jsonListenerContainerFactory") public ConcurrentKafkaListenerContainerFactory<String, YourJsonDto> jsonListenerContainerFactory( @Qualifier("jsonConsumerFactory") ConsumerFactory<String, YourJsonDto> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, YourJsonDto> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); return factory; } }
3. 在@KafkaListener中指定对应的容器工厂
在消费不同主题的监听器上,通过containerFactory属性指定对应的工厂:
@Component public class KafkaMessageListener { // 消费原有Avro主题,使用avroListenerContainerFactory @KafkaListener(topics = "原有Avro主题名", containerFactory = "avroListenerContainerFactory") public void consumeAvroMessage(SpecificRecord avroMessage) { // 处理Avro格式消息 } // 消费新JSON主题,使用jsonListenerContainerFactory @KafkaListener(topics = "新JSON主题名", containerFactory = "jsonListenerContainerFactory") public void consumeJsonMessage(YourJsonDto jsonMessage) { // 处理JSON格式消息 } }
注意事项
- 全局的
spring.kafka.consumer.value-deserializer配置可以不用设置,避免和局部配置冲突; YourJsonDto是你自定义的对应JSON消息的实体类,需要和JSON字段一一对应;- Avro配置中的
SCHEMA_REGISTRY_URL_CONFIG要替换为你实际的Schema Registry地址; - 如果需要配置消费者的分组ID,可以在对应的
ConsumerFactory配置中添加,或者直接在@KafkaListener的groupId属性中指定。
内容的提问来源于stack exchange,提问作者Jardo
相关产品推荐
相关产品推荐

