Spring环境下Kafka消费者如何反序列化生产者发送的DTO对象?
问题:Spring Kafka消费者无法反序列化跨微服务的DTO类
我使用Spring Kafka的KafkaTemplate发送消息,代码如下:
this.kafkaTemplate.send(applicationProperties.getKafka().getTopics().getNoticeGenerationEventsTopic(), japserName, noticeEventDTO);
但消费者端无法反序列化noticeEventDTO这个DTO,报错找不到生产者微服务中的类,错误栈如下:
Caused by: org.springframework.messaging.converter.MessageConversionException: failed to resolve class name. Class not found [com.xyz.abc.registration.service.dto.NoticeEventDTO] at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.getClassIdType(DefaultJackson2JavaTypeMapper.java:137) at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.toJavaType(DefaultJackson2JavaTypeMapper.java:98) at org.springframework.kafka.support.converter.JsonMessageConverter.determineJavaType(JsonMessageConverter.java:135) at org.springframework.kafka.support.converter.JsonMessageConverter.extractAndConvertValue(JsonMessageConverter.java:107) at org.springframework.kafka.support.converter.MessagingMessageConverter.toMessage(MessagingMessageConverter.java:193) at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.toMessagingMessage(MessagingMessageListenerAdapter.java:346) at org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter.onMessage(RecordMessagingMessageListenerAdapter.java:83) at org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter.onMessage(RecordMessagingMessageListenerAdapter.java:53) at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doInvokeOnMessage(KafkaMessageListenerContainer.java:2857) ... 12 common frames omitted Caused by: java.lang.ClassNotFoundException: com.abc.xyz.registration.service.dto.NoticeEventDTO at java.base/jdk.internal.loader.BuiltinClassLoader.loadClass(BuiltinClassLoader.java:641) at java.base/jdk.internal.loader.ClassLoaders$AppClassLoader.loadClass(ClassLoaders.java:188) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:520) at java.base/java.lang.Class.forName0(Native Method) at java.base/java.lang.Class.forName(Class.java:467) at org.springframework.boot.devtools.restart.classloader.RestartClassLoader.loadClass(RestartClassLoader.java:121) at java.base/java.lang.ClassLoader.loadClass(ClassLoader.java:520) at java.base/java.lang.Class.forName0(Native Method) at java.base/java.lang.Class.forName(Class.java:467) at org.springframework.util.ClassUtils.forName(ClassUtils.java:304) at org.springframework.kafka.support.mapping.DefaultJackson2JavaTypeMapper.getClassIdType(DefaultJackson2JavaTypeMapper.java:133) ... 20 common frames omitted
解决方案
1. 共享DTO公共模块
将生产者的NoticeEventDTO类抽取为独立的公共Jar模块,让生产者和消费者项目都依赖该模块。需保证DTO的包名、类名、字段定义(包括@JsonProperty等序列化注解)完全一致,这样消费者端就能直接找到对应类完成反序列化。
2. 配置自定义类型映射
如果不想共享模块,可通过配置Jackson类型映射,将生产者的全限定类名替换为自定义标识,映射到消费者本地的DTO类:
生产者端配置
@Bean public JsonMessageConverter jsonMessageConverter() { MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter(); DefaultJackson2JavaTypeMapper typeMapper = new DefaultJackson2JavaTypeMapper(); // 自定义类型标识,替换默认的全限定类名 Map<String, Class<?>> typeMap = new HashMap<>(); typeMap.put("notice-event", NoticeEventDTO.class); typeMapper.setIdClassMapping(typeMap); typeMapper.setTypePrecedence(Jackson2JavaTypeMapper.TypePrecedence.TYPE_ID); converter.setJavaTypeMapper(typeMapper); return converter; } @Bean public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> producerFactory, JsonMessageConverter jsonMessageConverter) { KafkaTemplate<String, Object> template = new KafkaTemplate<>(producerFactory); template.setMessageConverter(jsonMessageConverter); return template; }
消费者端配置
@Bean public JsonMessageConverter consumerJsonMessageConverter() { MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter(); DefaultJackson2JavaTypeMapper typeMapper = new DefaultJackson2JavaTypeMapper(); Map<String, Class<?>> typeMap = new HashMap<>(); // 用同一个自定义标识映射到消费者本地的DTO类 typeMap.put("notice-event", LocalNoticeEventDTO.class); typeMapper.setIdClassMapping(typeMap); typeMapper.setTypePrecedence(Jackson2JavaTypeMapper.TypePrecedence.TYPE_ID); converter.setJavaTypeMapper(typeMapper); return converter; } @Bean public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory( ConsumerFactory<String, Object> consumerFactory, JsonMessageConverter consumerJsonMessageConverter) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setMessageConverter(consumerJsonMessageConverter); return factory; }
3. 禁用类型信息头,指定固定反序列化类型
如果不需要自动类型推断,可关闭Spring Kafka默认的类型信息头,消费者端直接指定目标类型:
生产者端配置
@Bean public JsonMessageConverter jsonMessageConverter() { MappingJackson2MessageConverter converter = new MappingJackson2MessageConverter(); // 关闭自动生成类型信息头 converter.setGenerateTypeHeaders(false); return converter; }
消费者端处理
在@KafkaListener中直接指定接收的DTO类型:
@KafkaListener(topics = "${kafka.topics.notice-generation-events}", containerFactory = "kafkaListenerContainerFactory") public void handleNoticeEvent(@Payload LocalNoticeEventDTO noticeEvent) { // 业务处理逻辑 }
也可以使用StringDeserializer接收JSON字符串,再手动用Jackson反序列化到目标类:
@KafkaListener(topics = "${kafka.topics.notice-generation-events}") public void handleNoticeEvent(@Payload String noticeEventJson) { ObjectMapper objectMapper = new ObjectMapper(); LocalNoticeEventDTO noticeEvent = objectMapper.readValue(noticeEventJson, LocalNoticeEventDTO.class); // 业务处理逻辑 }
内容的提问来源于stack exchange,提问作者OwlR
相关产品推荐
相关产品推荐

