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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 09:25:40