Spring Boot Kafka配置trusted.packages后消费者仍遇反序列化错误
Kafka反序列化错误排查与解决
已完成Kafka消费者配置,生产者可正常发送消息,但消费者配置trusted.packages后仍持续触发反序列化错误。
生产者服务代码
@Service public class KafkaMessagePublisher { private final KafkaTemplate<String,Object> template; public KafkaMessagePublisher(KafkaTemplate<String,Object> template){ this.template = template; } public void sendEventsToTopic(Customer customer){ CompletableFuture<SendResult<String, Object>> fut = template.send(topic, customer); fut.whenComplete(((result, ex) -> { if(ex== null){ log.info("sent Message=[ {} ] with offset=[ {} ]" ,customer.toString(),result.getRecordMetadata().offset()); } else{ log.error("Unable to send the message=[ {} ] due to =[{}]" ,customer.toString(),ex.getMessage()); } })); } }
生产者配置
spring: kafka: producer: bootstrap-server: localhost:9092 key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.springframework.kafka.support.serializer.JsonSerializer # properties: # spring: # json: # trusted: # packages: com.example.kafka.producer.dto.* app: kafka: topic-name: TestTopic partitions: 2 replicationFactor: 1
消费者服务代码
@Service @Slf4j public class KafkaMessageListner { @KafkaListener(topics = "${app.kafka.topic-name}",groupId = "con") public void consume(Customer customer){ log.info("Consumer consume the message >>>>: {}",customer.toString()); } }
消费者配置
server: port: 9292 spring: kafka: consumer: bootstrap-server: localhost:9092 key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer properties: spring: json: trusted: packages: com.example.kafka.consumer.dto.* app: kafka: topic-name: TestTopic
曾尝试的无效配置
properties: spring: json: trusted: packages: java.util, java.lang, com.example.kafka.consumer.dto.*
Customer实体类(生产者、消费者端字段一致)
import lombok.Data; import lombok.ToString; @Data @ToString public class Customer { private String id; private String name; private String email; private String contactNo; }
核心报错信息
Caused by: java.lang.IllegalArgumentException: The class 'com.example.kafka.producer.dto.Customer' is not in the trusted packages: [java.util, java.lang, com.example.kafka.consumer.dto.*]. If you believe this class is safe to deserialize, please provide its name. If the serialization is only done by a trusted source, you can also enable trust all (*).
解决方案
1. 将生产者的包添加到消费者可信列表
报错明确显示,序列化后的类全限定名是com.example.kafka.producer.dto.Customer,但消费者只信任自己的com.example.kafka.consumer.dto.*包,因此需要把生产者的包加入可信配置:
spring: kafka: consumer: properties: spring: json: trusted: packages: com.example.kafka.consumer.dto.*, com.example.kafka.producer.dto.*
2. 统一实体类的包路径(推荐)
将生产者和消费者的Customer类移到同一个公共包下(比如com.example.kafka.common.dto),两边代码都引用这个包的实体类,然后消费者配置信任该公共包:
spring: kafka: consumer: properties: spring: json: trusted: packages: com.example.kafka.common.dto.*
3. 临时开启信任所有包(仅测试环境)
如果是测试环境且确认消息来源安全,可以直接信任所有包,但生产环境不建议使用:
spring: kafka: consumer: properties: spring: json: trusted: packages: "*"
内容的提问来源于stack exchange,提问作者Kishore Jetty
相关产品推荐
相关产品推荐

