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

Kafka序列化含BDTO列表与LocalDateTime的ADTO对象报错求助

Kafka 对象序列化异常排查与解决

问题描述

刚接触Kafka,使用KafkaTemplate发送自定义对象时持续抛出序列化异常:

Failed to execute CommandLineRunner... Caused by SerializationException: Error when serializing RecipeDTO to byte[]

相关代码结构如下:

实体类

class ADTO{
    /* 其他属性 */
    List<BDTO> bDTOs; // 注:原代码中List[]为笔误,修正为泛型列表
    LocalDateTime localDate;
}

class BDTO{
    /* 其他属性 */
}

自定义序列化器

public class ADTOSerializer implements Serializer<ADTO> {
    private final ObjectMapper objectMapper = new ObjectMapper();

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
    }

    @Override
    public byte[] serialize(String s, ADTO aDTO) {
        try {
            if (aDTO == null){
                return null;
            }
            return objectMapper.writeValueAsBytes(aDTO);
        } catch (Exception e) {
            throw new SerializationException("Error when serializing RecipeDTO to byte[]");
        }
    }

    @Override
    public void close() {
    }
}

生产者配置

public Map<String, Object> producerConfig(){
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ADTOSerializer.class);
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ADTOSerializer.class);
    return props;
}

对象发送代码

@Bean
CommandLineRunner commandLineRunner(KafkaTemplate<String, ADTO> kafkaTemplate){
    return args -> {
        ADTO aDTO = new ADTO(
            /* 属性初始化 */
        );
        kafkaTemplate.send("topic", aDTO);
    };
}

推测异常由BDTO列表或LocalDateTime属性导致,询问是否需要为二者单独编写序列化器,并请求示例代码。

问题根源

  1. LocalDateTime 序列化支持缺失:Jackson 默认不支持 JDK8 新增的LocalDateTime等日期时间类型,直接序列化会触发类型不匹配异常。
  2. BDTO 序列化约束未满足:若BDTO没有无参构造器、属性未提供 getter 方法,Jackson 无法自动解析其结构,导致序列化失败。
  3. 生产者配置错误:KEY 类型为String,却配置了ADTOSerializer,这会导致 KEY 序列化时额外报错。

解决方案

无需为BDTO单独编写序列化器,只需配置 Jackson 支持日期类型、修正实体类约束,并调整生产者配置即可。

1. 配置 Jackson 支持 LocalDateTime

修改自定义序列化器,添加 JDK8 时间模块:

public class ADTOSerializer implements Serializer<ADTO> {
    private final ObjectMapper objectMapper;

    public ADTOSerializer() {
        objectMapper = new ObjectMapper();
        // 注册JDK8时间模块,支持LocalDateTime序列化
        objectMapper.registerModule(new JavaTimeModule());
        // 可选:将日期序列化为ISO格式字符串,而非时间戳数组
        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
    }

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
    }

    @Override
    public byte[] serialize(String s, ADTO aDTO) {
        try {
            if (aDTO == null){
                return null;
            }
            return objectMapper.writeValueAsBytes(aDTO);
        } catch (Exception e) {
            // 打印具体异常栈,便于排查细节
            e.printStackTrace();
            throw new SerializationException("Error when serializing ADTO to byte[]", e);
        }
    }

    @Override
    public void close() {
    }
}

2. 修正 BDTO 序列化约束

确保BDTO具备无参构造器和属性 getter 方法:

class BDTO{
    /* 其他属性 */
    
    // 必须添加无参构造器(Jackson默认依赖)
    public BDTO() {}

    // 为所有需要序列化的属性添加getter方法
    // 示例:private String name;
    // public String getName() { return name; }
}

3. 修正生产者配置

KEY 类型为String,需改用StringSerializer:

public Map<String, Object> producerConfig(){
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
    // KEY为String类型,使用StringSerializer
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ADTOSerializer.class);
    return props;
}

4. 配套消费者反序列化器

消费者接收ADTO时,需要对应的反序列化器:

public class ADTODeserializer implements Deserializer<ADTO> {
    private final ObjectMapper objectMapper;

    public ADTODeserializer() {
        objectMapper = new ObjectMapper();
        objectMapper.registerModule(new JavaTimeModule());
        objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
    }

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
    }

    @Override
    public ADTO deserialize(String s, byte[] bytes) {
        try {
            if (bytes == null){
                return null;
            }
            return objectMapper.readValue(bytes, ADTO.class);
        } catch (Exception e) {
            e.printStackTrace();
            throw new SerializationException("Error when deserializing byte[] to ADTO", e);
        }
    }

    @Override
    public void close() {
    }
}

同时配置消费者工厂:

@Bean
public ConsumerFactory<String, ADTO> consumerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
    props.put(ConsumerConfig.GROUP_ID_CONFIG, "groupID");
    props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ADTODeserializer.class);
    return new DefaultKafkaConsumerFactory<>(props);
}

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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 05:33:17