Spring Boot Kafka消费者ClassNotFoundException问题求助
问题原因与解决方案
这个问题我之前也碰到过,核心原因和Spring Kafka的JSON序列化机制有关,我来给你拆解一下并提供可行的解决方案:
问题根源
Spring Kafka默认的JsonSerializer在发送消息时,会把发送对象的全限定类名(比如你的com.test.demo.domain.Account)放到消息头的__TypeId__字段中。当消费者端的消息转换器(比如你用的StringJsonMessageConverter)尝试反序列化时,会优先读取这个__TypeId__头去加载对应类,但你的消费者用的是另一个包下的Account类,所以就会触发ClassNotFoundException或者转换失败的异常。
解决方案
这里提供两种常用的解决思路,你可以根据自己的场景选择:
方案一:自定义类型映射(推荐)
让生产者和消费者约定一个统一的类型标识(比如"Account"),代替默认的全限定类名,这样两边的类即使包名不同,也能通过这个标识匹配上。
生产者端修改代码
@Bean public ProducerFactory<String, Account> accountProducerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 自定义JsonSerializer的类型映射 JsonSerializer<Account> jsonSerializer = new JsonSerializer<>(); DefaultJackson2JavaTypeMapper typeMapper = new DefaultJackson2JavaTypeMapper(); Map<String, Class<?>> typeMap = new HashMap<>(); typeMap.put("Account", Account.class); // 设置自定义类型标识 typeMapper.setIdType(Jackson2JavaTypeMapper.IdType.NAME); typeMapper.setTypeMap(typeMap); jsonSerializer.setTypeMapper(typeMapper); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, jsonSerializer); return new DefaultKafkaProducerFactory<>(configProps); } @Bean public KafkaTemplate<String, Account> accountKafkaTemplate() { ProducerFactory<String, Account> factory = accountProducerFactory(); return new KafkaTemplate<>(factory); }
消费者端修改代码
public ConsumerFactory<String, com.kconsumer.accountconsumer.domain.Account> accountConsumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupName); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 自定义JsonDeserializer的类型映射,和生产者保持一致 JsonDeserializer<com.kconsumer.accountconsumer.domain.Account> jsonDeserializer = new JsonDeserializer<>(); DefaultJackson2JavaTypeMapper typeMapper = new DefaultJackson2JavaTypeMapper(); Map<String, Class<?>> typeMap = new HashMap<>(); typeMap.put("Account", com.kconsumer.accountconsumer.domain.Account.class); // 对应生产者的类型标识 typeMapper.setIdType(Jackson2JavaTypeMapper.IdType.NAME); typeMapper.setTypeMap(typeMap); jsonSerializer.setTypeMapper(typeMapper); jsonDeserializer.addTrustedPackages("com.kconsumer.accountconsumer.domain"); // 信任当前消费者的包 configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, jsonDeserializer); return new DefaultKafkaConsumerFactory<>(configProps); } @Bean public KafkaListenerContainerFactory<?> kafkaJsonListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, com.kconsumer.accountconsumer.domain.Account> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(accountConsumerFactory()); // 这里不需要再用StringJsonMessageConverter,JsonDeserializer会直接把消息转成目标对象 return factory; }
方案二:禁用__TypeId__消息头
如果生产者和消费者的Account类字段完全一致(只是包名/类名不同),可以让生产者不发送__TypeId__头,消费者直接根据目标类的字段结构反序列化JSON。
生产者端修改
@Bean public ProducerFactory<String, Account> accountProducerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); JsonSerializer<Account> jsonSerializer = new JsonSerializer<>(); // 重写类型映射逻辑,不添加__TypeId__头 jsonSerializer.setTypeMapper(new DefaultJackson2JavaTypeMapper() { @Override public void fromJavaType(JavaType javaType, Headers headers) { // 空实现,不添加类型头 } }); configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, jsonSerializer); return new DefaultKafkaProducerFactory<>(configProps); }
消费者端修改
public ConsumerFactory<String, com.kconsumer.accountconsumer.domain.Account> accountConsumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); configProps.put(ConsumerConfig.GROUP_ID_CONFIG, groupName); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); JsonDeserializer<com.kconsumer.accountconsumer.domain.Account> jsonDeserializer = new JsonDeserializer<>(com.kconsumer.accountconsumer.domain.Account.class); jsonDeserializer.addTrustedPackages("com.kconsumer.accountconsumer.domain"); jsonDeserializer.setUseTypeHeaders(false); // 禁用类型头检查 configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, jsonDeserializer); return new DefaultKafkaConsumerFactory<>(configProps); }
内容的提问来源于stack exchange,提问作者Joe P
相关产品推荐
相关产品推荐

