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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 09:20:54