Spring Boot Kafka消费者遇ClassNotFoundException问题求助
问题描述
Spring Boot Kafka Producer通过Postman测试正常,发送<String, Complaint>类型消息,key使用StringSerializer序列化,value使用JsonSerializer序列化。但Consumer接收消息时抛出ClassNotFoundException,提示找不到ma.attijari.kafkacomplaintsproducer.models.Complaint类。Consumer端已创建结构完全一致的Complaint类,且两端均配置了类型映射(将Complaint映射至各自全类名),但问题仍未解决。
已尝试的Consumer端配置
- Config.java代码配置:
config.put(JsonDeserializer.TYPE_MAPPINGS, "Complaint:com.example.kafkaconsumertest.Complaint");
- application.properties配置:
spring.kafka.consumer.type-mappings=Complaint:com.example.kafkacomplaintsproducer.models.Complaint
错误日志
2024-07-25T22:40:51.749+01:00 ERROR 26752 --- [kafkaConsumerTest] [ntainer#0-0-C-1] o.a.k.c.c.internals.CompletedFetch : [Consumer clientId=consumer-group_id-1, groupId=group_id] Deserializers with error: Deserializers{keyDeserializer=org.apache.kafka.common.serialization.StringDeserializer@44fb7281, valueDeserializer=org.springframework.kafka.support.serializer.JsonDeserializer@efa868c} 2024-07-25T22:40:51.749+01:00 ERROR 26752 --- [kafkaConsumerTest] [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer : Consumer exception java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer at org.springframework.kafka.listener.DefaultErrorHandler.handleOtherException(DefaultErrorHandler.java:192) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1925) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1348) ~[spring-kafka-3.2.2.jar:3.2.2] at java.base/java.util.concurrent.CompletableFuture$AsyncRun.run(CompletableFuture.java:1804) ~[na:na] at java.base/java.lang.Thread.run(Thread.java:1583) ~[na:na] Caused by: org.apache.kafka.common.errors.RecordDeserializationException: Error deserializing key/value for partition complaints_test-0 at offset 1. If needed, please seek past the record to continue consumption. at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:331) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.CompletedFetch.fetchRecords(CompletedFetch.java:283) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.FetchCollector.fetchRecords(FetchCollector.java:168) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.FetchCollector.collectFetch(FetchCollector.java:134) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.Fetcher.collectFetch(Fetcher.java:145) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.pollForFetches(LegacyKafkaConsumer.java:666) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.poll(LegacyKafkaConsumer.java:617) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.LegacyKafkaConsumer.poll(LegacyKafkaConsumer.java:590) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:874) ~[kafka-clients-3.7.1.jar:na] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollConsumer(KafkaMessageListenerContainer.java:1625) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1600) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1405) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1296) ~[spring-kafka-3.2.2.jar:3.2.2] ... 2 common frames omitted Caused by: org.springframework.messaging.converter.MessageConversionException: failed to resolve class name. Class not found [ma.attijari.kafkacomplaintsproducer.models.Complaint] at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.getClassIdType(DefaultJackson2JavaTypeMapper.java:137) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.toJavaType(DefaultJackson2JavaTypeMapper.java:98) ~[spring-kafka-3.2.2.jar:3.2.2] at org.springframework.kafka.support.serializer.JsonDeserializer.deserialize(JsonDeserializer.java:571) ~[spring-kafka-3.2.2.jar:3.2.2] at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:73) ~[kafka-clients-3.7.1.jar:na] at org.apache.kafka.clients.consumer.internals.CompletedFetch.parseRecord(CompletedFetch.java:321) ~[kafka-clients-3.7.1.jar:na] ... 14 common frames omitted Caused by: java.lang.ClassNotFoundException: ma.attijari.kafkacomplaintsproducer.models.Complaint at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:641) ~[na:na] at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188) ~[na:na] at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:526) ~[na:na] at java.base/java.lang.Class.forName0(Native Method) ~[na:na] at java.base/java.lang.Class.forName(Class.java:534) ~[na:na] at java.base/java.lang.Class.forName(Class.java:513) ~[na:na] at org.springframework.util.ClassUtils.forName(ClassUtils.java:304) ~[spring-core-6.1.11.jar:6.1.11] at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.getClassIdType(DefaultJackson2JavaTypeMapper.java:133) ~[spring-kafka-3.2.2.jar:3.2.2] ... 18 common frames omitted
排查及解决步骤
1. 补全Producer端的类型映射配置
问题核心是Producer端未配置类型映射,导致序列化时将Producer端的全类名ma.attijari.kafkacomplaintsproducer.models.Complaint写入消息头,而非约定的别名Complaint,Consumer端的映射规则无法匹配原始全类名,因此抛出异常。
在Producer端添加以下配置:
- 方式一:application.properties配置
spring.kafka.producer.value-serializer=org.springframework.kafka.support.serializer.JsonSerializer spring.kafka.producer.type-mappings=Complaint:ma.attijari.kafkacomplaintsproducer.models.Complaint
- 方式二:Java Config类配置
@Bean public ProducerFactory<String, Complaint> producerFactory() { Map<String, Object> config = new HashMap<>(); // 补充其他Kafka基础配置(如bootstrap-servers等) config.put(JsonSerializer.TYPE_MAPPINGS, "Complaint:ma.attijari.kafkacomplaintsproducer.models.Complaint"); return new DefaultKafkaProducerFactory<>(config); }
2. 统一Consumer端的配置方式
同时在Config.java和application.properties中配置类型映射可能导致冲突,建议保留一种配置方式。例如仅保留application.properties的配置,并确保反序列化权限:
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.JsonDeserializer spring.kafka.consumer.type-mappings=Complaint:com.example.kafkaconsumertest.Complaint # 可选:信任所有包,避免反序列化权限拦截 spring.kafka.consumer.properties.spring.json.trusted.packages=*
3. 清理旧消息或重置Consumer偏移量
之前发送的消息已携带Producer端的全类名类型信息,即使配置修复,旧消息仍会触发错误。可选择以下方式处理:
- 测试环境直接删除Kafka主题的旧消息
- 重置Consumer的消费偏移量到最新位置,跳过旧消息:
@Autowired private KafkaListenerEndpointRegistry registry; public void resetConsumerOffset() { registry.getListenerContainers().forEach(container -> { container.stop(); container.seekToEnd(container.getAssignedPartitions()); container.start(); }); }
4. 配置ErrorHandlingDeserializer(可选)
根据日志提示,配置ErrorHandlingDeserializer可避免单条错误消息阻塞整个消费流程:
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer
内容的提问来源于stack exchange,提问作者khaoula baraka

