Spring Boot Kafka消费者初始化失败:未配置委托反序列化器
问题
开发了一个通知微服务,用于消费Kafka中的UserDTO类型对象。生产者已成功将数据写入topic(通过Offset Explorer验证),但启动消费者微服务时失败,报错显示无法构造Kafka Consumer,根本原因是未配置委托反序列化器。
消费者代码与配置
UserDTO.java
package com.example.notificationservice.model; public class UserDTO { private String userId; private String userName; private String email; // getters and setters }
KafkaMessageListener.java
package com.example.notificationservice.consumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Service; import com.example.notificationservice.model.UserDTO; @Service public class KafkaMessageListener { Logger log = LoggerFactory.getLogger(KafkaMessageListener.class); @KafkaListener(topics = "topic-example16", groupId = "group12") public void consume(UserDTO user) { log.info("consumer consume {}", user); } }
application.properties(消费者)
spring.kafka.consumer.bootstrap-servers=localhost:9092 spring.kafka.consumer.group-id=group12 spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.key.delegate=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.value.delegate=org.springframework.kafka.support.serializer.JsonDeserializer spring.kafka.consumer.properties.spring.json.trusted.packages=* spring.kafka.consumer.properties.spring.json.value.default.type=com.example.notificationservice.model.UserDTO
生产者代码与配置
ProducerController.java
package com.kafkaProducer.producer; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.web.bind.annotation.*; import com.kafkaProducer.UserDTO; import org.springframework.http.ResponseEntity; @RestController @RequestMapping("/kafka") public class ProducerController { private final KafkaTemplate<String, UserDTO> kafkaTemplate; @Autowired public ProducerController(KafkaTemplate<String, UserDTO> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } @PostMapping("/send") public ResponseEntity<String> sendMessage(@RequestBody UserDTO user) { kafkaTemplate.send("topic-example16", user); return ResponseEntity.ok("User data sent to Kafka."); } }
application.properties(生产者)
spring.kafka.producer.bootstrap-servers=localhost:9092 spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer
报错信息
ERROR[0;39m [35m30932[0;39m [2m---[0;39m [2m[ main][0;39m [36mo.s.boot.SpringApplication [0;39m [2m:[0;39m Application run failed org.springframework.context.ApplicationContextException: Failed to start bean 'org.springframework.kafka.config.internalKafkaListenerEndpointRegistry'; nested exception is org.apache.kafka.common.KafkaException: Failed to construct kafka consumer at org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:182) ~[spring-context-5.3.29.jar:5.3.29] ...(省略中间堆栈) Caused by: java.lang.IllegalStateException: No delegate deserializer configured at org.springframework.util.Assert.state(Assert.java:76) ~[spring-core-5.3.29.jar:5.3.29] at org.springframework.kafka.support.serializer.ErrorHandlingDeserializer.configure(ErrorHandlingDeserializer.java:138) ~[spring-kafka-2.8.11.jar:2.8.11] at org.apache.kafka.clients.consumer.KafkaConsumer.<init>(KafkaConsumer.java:716) ~[kafka-clients-3.1.2.jar:na] ... 28 common frames omitted
解决方案
问题出在消费者的配置错误:
- 为
key-deserializer指定了StringDeserializer,但配置委托反序列化器时,错误地将spring.deserializer.key.delegate设置成了ErrorHandlingDeserializer,导致ErrorHandlingDeserializer找不到正确的委托类,触发报错。
修正后的消费者application.properties配置:
spring.kafka.consumer.bootstrap-servers=localhost:9092 spring.kafka.consumer.group-id=group12 spring.kafka.consumer.auto-offset-reset=earliest # key直接使用StringDeserializer,不需要套ErrorHandlingDeserializer spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer # value使用ErrorHandlingDeserializer做错误处理,委托给JsonDeserializer执行实际反序列化 spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.value.delegate=org.springframework.kafka.support.serializer.JsonDeserializer spring.kafka.consumer.properties.spring.json.trusted.packages=* spring.kafka.consumer.properties.spring.json.value.default.type=com.example.notificationservice.model.UserDTO
补充说明
- 若要给key也增加错误处理能力,需调整为以下配置:
spring.kafka.consumer.key-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.key.delegate=org.apache.kafka.common.serialization.StringDeserializer
ErrorHandlingDeserializer仅做异常捕获与包装,必须指定实际执行序列化/反序列化的委托类才能正常工作。- 生产者的配置无问题,保持当前
JsonSerializer即可正确序列化UserDTO。
内容的提问来源于stack exchange,提问作者Navneet C.
相关产品推荐
相关产品推荐

