Spring for Apache Kafka JSON反序列化ClassNotFound异常排查
问题描述
从Kafka主题消费消息时,无明显诱因抛出如下错误:
2022-06-28 14:17:52.044 INFO 1 --- [ntainer#0-0-C-1] o.a.k.clients.consumer.KafkaConsumer : [Consumer clientId=consumer-api1-1, groupId=api1] Seeking to offset 1957 for partition ActiveProxySources-0 2022-06-28T14:17:52.688451744Z 2022-06-28 14:17:52.687 ERROR 1 --- [ntainer#0-0-C-1] o.s.kafka.listener.DefaultErrorHandler : Backoff none exhausted for ActiveProxySources-0@1957 2022-06-28T14:17:52.688499949Z 2022-06-28T14:17:52.688511943Z org.springframework.kafka.listener.ListenerExecutionFailedException: Listener failed; nested exception is org.springframework.kafka.support.serializer.DeserializationException: failed to deserialize; nested exception is org.springframework.messaging.converter.MessageConversionException: failed to resolve class name. Class not found [com.freeproxy.parser.model.kafka.KafkaMessage]; nested exception is java.lang.ClassNotFoundException: com.freeproxy.parser.model.kafka.KafkaMessage 2022-06-28 14:17:52.688964715Z
其余使用相同配置进行Kafka消息生产、消费的应用均可正常运行,唯独当前应用出现该问题。预期需求是从两个Kafka主题消费消息(两个主题内的消息结构完全一致,存储相同类型的对象),但即使尝试仅消费单个主题,也会抛出上述错误。
部署环境为Docker-Compose运行Kafka及所有关联应用,已编写自定义Kafka消费者配置类,但无论是消费双主题还是单主题,错误仍然存在。
相关配置
消息实体与消费逻辑
class KafkaMessage { String id IdStatus status } @Service @Slf4j class ConsumerService { Set<String> activeProxies = [] int getActiveProxiesNumber() { activeProxies.size() } Set<String> activeProxySources = [] int getActiveProxySourcesNumber() { activeProxySources.size() } @KafkaListener(topics = "ActiveProxies") public void consumeProxyId(KafkaMessage message) { log.info("Consuming ${message.id}: ${message.status}") if (message.status == IdStatus.ADD) activeProxies.add(message.id) if (message.status == IdStatus.DELETE) activeProxies.remove(message.id) } @KafkaListener(topics = "ActiveProxySources") public void consumeProxySourceId(KafkaMessage message) { log.info("Consuming ${message.id}: ${message.status}") if (message.status == IdStatus.ADD) activeProxySources.add(message.id) if (message.status == IdStatus.DELETE) activeProxySources.remove(message.id) } }
Topic配置类
@Configuration public class TopicConfig { @Value(value = "kafka:9092") private String bootstrapAddress @Value(value = "ActiveProxies") private String activeProxies @Value(value = "ActiveProxySources") private String activeProxySources @Bean public NewTopic ActiveProxiesTopic() { return TopicBuilder.name(activeProxies).partitions(1).replicas(1) .config(org.apache.kafka.common.config.TopicConfig.RETENTION_MS_CONFIG, "60000").build() } @Bean public NewTopic ActiveProxySourcesTopic() { return TopicBuilder.name(activeProxySources).partitions(1).replicas(1) .config(org.apache.kafka.common.config.TopicConfig.RETENTION_MS_CONFIG, "60000").build() } }
application.properties配置
server.port=30329 spring.data.mongodb.database=free-proxy-engine spring.kafka.bootstrap-servers=kafka:9092 spring.kafka.consumer.group-id=consumer-Api1 spring.kafka.consumer.properties.spring.json.trusted.packages=* spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.consumer.enable-auto-commit=false 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.value.delegate.class=org.springframework.kafka.support.serializer.JsonDeserializer
故障原因
核心报错为java.lang.ClassNotFoundException: com.freeproxy.parser.model.kafka.KafkaMessage,根因是Spring自带的JsonDeserializer默认会读取消息头中存储的生产端消息实体全限定类名,尝试在当前消费端类路径下加载同名类完成反序列化。
当前消费端的KafkaMessage类所在包路径,和生产端序列化时写入消息头的com.freeproxy.parser.model.kafka路径不匹配,因此类加载失败触发反序列化异常。配置的spring.json.trusted.packages=*仅用于放开反序列化的包信任范围,不会修改类加载匹配逻辑,因此无法解决该问题。
解决方案
选择以下任意一种方案即可修复:
- 方案1:将当前消费端的
KafkaMessage类移动到com.freeproxy.parser.model.kafka包下,保证和生产端实体类的全限定名完全一致,重启应用即可正常消费。 - 方案2:修改消费者配置,关闭从消息头读取类名的默认行为,直接指定反序列化目标类型。在application.properties中追加如下配置:
# 替换为你当前项目中KafkaMessage类的实际全限定名 spring.kafka.consumer.properties.spring.json.value.default.type=com.yourpackage.KafkaMessage spring.kafka.consumer.properties.spring.json.use.type.headers=false
配置后反序列化会直接将消息体转换为指定类型,不再依赖消息头中存储的生产端类名。
- 方案3:如果不需要修改全局消费者配置,可在
@KafkaListener注解上单独为对应监听器指定反序列化规则和目标类型,不影响项目内其他消费者的运行。
内容的提问来源于stack exchange,提问作者Oleksandr Popov
相关产品推荐
相关产品推荐

