You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Boot Kafka消费者遇ClassNotFoundException问题求助

Spring Boot Kafka Consumer ClassNotFoundException 排查方案

问题描述

Spring Boot Kafka Producer通过Postman测试正常,发送<String, Complaint>类型消息,key使用StringSerializer序列化,value使用JsonSerializer序列化。但Consumer接收消息时抛出ClassNotFoundException,提示找不到ma.attijari.kafkacomplaintsproducer.models.Complaint类。Consumer端已创建结构完全一致的Complaint类,且两端均配置了类型映射(将Complaint映射至各自全类名),但问题仍未解决。

已尝试的Consumer端配置

  1. Config.java代码配置:
config.put(JsonDeserializer.TYPE_MAPPINGS, "Complaint:com.example.kafkaconsumertest.Complaint");
  1. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.20 11:45:55