Spring Boot Kafka多主题JSON自动反序列化问题咨询
之前用Spring Kafka消费products主题时,通过全局配置指定默认反序列化目标类com.example.inventoryservice.core.entity.Product,能正常将消息转为Product对象。现在需要同时消费products和users主题,分别对应Product和User实体类,希望实现无需全局指定默认类型,让框架自动根据@KafkaListener方法的参数类型完成反序列化,或者用其他更灵活的方式。
方案1:启用类型头信息(推荐)
让生产者在发送消息时,把目标类的类型信息放到Kafka消息头里,消费者通过读取这个头信息自动匹配反序列化类型,同时配合方法参数类型完成绑定。
配置修改
- 消费者端
application.yml:
spring: kafka: bootstrap-servers: localhost:9094 consumer: key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: spring: json: use: type: headers: true # 开启读取类型头 # 无需再配置全局default.type
- 生产者端(若由你负责)对应配置,确保发送时带上类型头:
spring: kafka: producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer properties: spring: json: add: type: headers: true # 发送时添加类型头
配置完成后,消费者收到消息时会自动通过头里的类型信息反序列化成对应对象,你的@KafkaListener方法可以直接用Product和User作为参数,无需额外配置。
方案2:为每个Listener配置专属的容器工厂
如果无法修改生产者配置(比如生产者是第三方服务,不会发送类型头),可以创建多个ConcurrentKafkaListenerContainerFactory,每个工厂对应不同的反序列化目标类,然后在@KafkaListener里指定使用对应的工厂。
配置类示例
@Configuration public class KafkaConfig { @Value("${spring.kafka.bootstrap-servers}") private String bootstrapServers; // 针对Product的容器工厂 @Bean public ConcurrentKafkaListenerContainerFactory<String, Product> productKafkaContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Product> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(productConsumerFactory()); return factory; } private ConsumerFactory<String, Product> productConsumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, Product.class.getName()); props.put(JsonDeserializer.USE_TYPE_HEADERS, false); return new DefaultKafkaConsumerFactory<>(props); } // 针对User的容器工厂 @Bean public ConcurrentKafkaListenerContainerFactory<String, User> userKafkaContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, User> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(userConsumerFactory()); return factory; } private ConsumerFactory<String, User> userConsumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class); props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, User.class.getName()); props.put(JsonDeserializer.USE_TYPE_HEADERS, false); return new DefaultKafkaConsumerFactory<>(props); } }
修改@KafkaListener注解
public class KafkaMessagingService implements MessagingService { @Override @KafkaListener(id = "product_consumer", topics = "products", containerFactory = "productKafkaContainerFactory") public void processProductAdded(Product product) { System.out.println(product); } @Override @KafkaListener(id = "user_consumer", topics = "users", containerFactory = "userKafkaContainerFactory") public void processUserAdded(User user) { System.out.println(user); } }
方案3:使用@KafkaHandler处理多类型消息
如果希望用同一个Listener容器处理多个主题的消息,可以把多个处理方法放在同一个类里,用@KafkaHandler标注每个方法,框架会根据消息类型自动匹配对应的方法。
代码示例
@Component @KafkaListener(id = "inventory_service_consumer", topics = {"products", "users"}) public class KafkaMessagingService implements MessagingService { @KafkaHandler public void processProductAdded(Product product) { System.out.println("Received Product: " + product); } @KafkaHandler public void processUserAdded(User user) { System.out.println("Received User: " + user); } }
对应配置
这种方式同样需要生产者发送类型头,消费者配置开启读取类型头(同方案1的消费者配置),这样框架才能识别消息类型并匹配到对应的@KafkaHandler方法。
内容的提问来源于stack exchange,提问作者Djole Pi

