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属性导致,询问是否需要为二者单独编写序列化器,并请求示例代码。
问题根源
- LocalDateTime 序列化支持缺失:Jackson 默认不支持 JDK8 新增的
LocalDateTime等日期时间类型,直接序列化会触发类型不匹配异常。 - BDTO 序列化约束未满足:若
BDTO没有无参构造器、属性未提供 getter 方法,Jackson 无法自动解析其结构,导致序列化失败。 - 生产者配置错误: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
相关产品推荐
相关产品推荐

