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

Spring Boot微服务间Kafka DTO反序列化问题:类不在信任包

Spring Boot Kafka 反序列化问题解决方案

问题描述

我正在开发基于Spring Boot的微服务架构应用,认证服务生产消息,用户服务消费消息,但消费端反序列化时遇到问题。

当前在两个事件场景生产消息:

  1. 用户在认证服务注册时发送UserDTO:
public record UserDTO(
     UUID id,
     String username,
     String email,
     Role role
) {
}
// Note: I use same name for records across service
  1. 用户修改邮箱时发送EmailChangeDTO:
public record EmailChangeDTO(
        String username,
        @NonNull String email
) {
}

我使用String序列化Key、Json序列化Value,但用户服务反序列化时出现以下错误:

Caused by: java.lang.IllegalArgumentException: The class 'com.example.authservice.Dto.EmailChangeDTO' is not in the trusted packages: [java.util, java.lang, com.example.userservice.DTO]. If you believe this class is safe to deserialize, please provide its name. If the serialization is only done by a trusted source, you can also enable trust all (*).

用户注册场景也有相同错误。

我的消费者配置如下:

package com.example.userservice.Config;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.kafka.support.serializer.ErrorHandlingDeserializer;
import org.springframework.kafka.support.serializer.JsonDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.annotation.EnableKafka;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

import java.util.HashMap;
import java.util.Map;

@Configuration
@EnableKafka
public class KafkaConsumerConfig {
    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public ConsumerFactory<String,Object> consumerFactory(){
        Map<String, Object> props = new HashMap<>();

        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "user-service-group");

        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);

        props.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class);
        props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class);
        props.put(JsonDeserializer.TRUSTED_PACKAGES,"com.userservice.userservice.DTO");
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Object> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

尝试添加认证服务包到信任包后,出现以下错误:

Caused by: java.lang.ClassNotFoundException: com.example.authservice.Dto.UserDTO

但将消息反序列化为String后,用ObjectMapper手动映射到用户服务DTO可正常工作:

@KafkaListener(topics = "customer-registration", groupId = "user-service-group")
    public void consumeUserRegistration(String message) {
        try {
            log.info("Received user registration event: {}", message);
            UserDTO userDTO = objectMapper.readValue(message, UserDTO.class);
            userService.createUser(userDTO);
        } catch (JsonProcessingException e) {
            log.error("Error processing user registration event", e);
        }
    }

请问如何配置Kafka的JsonDeserializer以信任多包DTO类?手动用ObjectMapper反序列化是否最优?


解决方案

1. 配置JsonDeserializer信任多包并正确反序列化

问题核心原因:

  • 信任包路径配置错误:原配置中com.userservice.userservice.DTO是重复路径,需改为用户服务实际的DTO包路径(如com.example.userservice.DTO)
  • 认证服务的DTO类不在用户服务类路径中,直接信任认证服务包无意义,因为用户服务不存在该类

正确配置步骤:

步骤1:修正信任包配置

在消费者配置中,将信任包设置为用户服务自身的DTO包,多个包用逗号分隔;内部可信环境可直接信任所有包:

@Bean
public ConsumerFactory<String,Object> consumerFactory(){
    Map<String, Object> props = new HashMap<>();

    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "user-service-group");

    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);

    props.put(ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS, StringDeserializer.class);
    props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, JsonDeserializer.class);
    
    // 信任用户服务的DTO包,多包用逗号分隔
    props.put(JsonDeserializer.TRUSTED_PACKAGES, "com.example.userservice.DTO");
    // 内部环境可信任所有包
    // props.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
    
    return new DefaultKafkaConsumerFactory<>(props);
}

步骤2:在KafkaListener中明确指定目标DTO类型

在消费方法中,通过@Payload直接指定要反序列化的本地DTO类,Spring Kafka会自动匹配:

@KafkaListener(topics = "customer-registration", groupId = "user-service-group")
public void consumeUserRegistration(@Payload UserDTO userDTO) {
    log.info("Received user registration event: {}", userDTO);
    userService.createUser(userDTO);
}

@KafkaListener(topics = "email-change", groupId = "user-service-group")
public void consumeEmailChange(@Payload EmailChangeDTO emailChangeDTO) {
    log.info("Received email change event: {}", emailChangeDTO);
    // 执行邮箱修改业务逻辑
}

步骤3:优化生产者配置(可选)

为避免生产者写入自身类全限定名导致消费者混淆,可在生产者端关闭类型信息头:

props.put(JsonSerializer.USE_TYPE_INFO_HEADERS, false);

2. 手动用ObjectMapper反序列化是否最优?

手动反序列化不是最优方案,原因如下:

  • 冗余代码:Spring Kafka的JsonDeserializer已封装ObjectMapper逻辑,手动处理属于重复造轮子
  • 配置不统一:每个消费方法单独处理反序列化,难以统一配置ObjectMapper的全局行为(如日期格式化、空值处理)
  • 错误处理繁琐:ErrorHandlingDeserializer可统一捕获反序列化异常,手动处理需每个方法编写try-catch块

仅在消息格式不固定、需要自定义复杂反序列化逻辑时,手动处理才有价值。当前场景下,使用Spring Kafka提供的JsonDeserializer并正确配置,是更简洁、可维护的选择。


内容的提问来源于stack exchange,提问作者Anish

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 17:44:56