Kafka报SerializationException异常如何配置忽略反序列化失败消息
问题根因
Kafka消费流程中,消息反序列化动作执行在消费者拉取消息之后、@KafkaListener标注的监听方法被调用之前,因此写在监听方法内部的try-catch无法捕获SerializationException。未做特殊配置时,该异常会导致消费线程反复重试拉取同一条格式错误的消息,直接阻断整个消费流程。
实现方案
1. 修改Kafka配置类,包装反序列化器
使用Spring Kafka内置的ErrorHandlingDeserializer分别包装原有的key、value反序列化器。反序列化失败时该包装类不会直接向外抛出阻断异常,而是将异常信息写入消息Header,将消息值置为null传递给后续流程。
修改后的KafkaConfig代码如下:
@Configuration @Slf4j public class KafkaConfig { private final KafkaProperties kafkaProperties; public KafkaConfig(KafkaProperties kafkaProperties) { this.kafkaProperties = kafkaProperties; } @Bean public ConcurrentKafkaListenerContainerFactory<String, ExternalCardServiceData> kafkaContainerFactory() { // 构造目标类型JSON反序列化器,配置信任所有包避免类型校验异常 JsonDeserializer<ExternalCardServiceData> jsonDeserializer = new JsonDeserializer<>(ExternalCardServiceData.class); jsonDeserializer.addTrustedPackages("*"); // 用ErrorHandlingDeserializer包装key、value的实际反序列化器 ErrorHandlingDeserializer<String> keyErrorDeserializer = new ErrorHandlingDeserializer<>(new StringDeserializer()); ErrorHandlingDeserializer<ExternalCardServiceData> valueErrorDeserializer = new ErrorHandlingDeserializer<>(jsonDeserializer); DefaultKafkaConsumerFactory<String, ExternalCardServiceData> consumerFactory = new DefaultKafkaConsumerFactory<>(kafkaProperties.buildConsumerProperties(), keyErrorDeserializer, valueErrorDeserializer); ConcurrentKafkaListenerContainerFactory<String, ExternalCardServiceData> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // 配置错误处理器:异常触发时记录日志,不重试直接提交偏移量跳过坏消息 factory.setCommonErrorHandler(new DefaultErrorHandler((consumerRecord, e) -> { log.error("消费消息失败,跳过当前消息,topic:{}, partition:{}, offset:{}, 异常信息:{}", consumerRecord.topic(), consumerRecord.partition(), consumerRecord.offset(), e.getMessage()); }, new FixedBackOff(0L, 0L))); // 配置消费完成单条消息后自动提交偏移量 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD); return factory; } }
注:如果使用的Spring Kafka版本低于2.8,将上述代码中的
DefaultErrorHandler替换为SeekToCurrentErrorHandler,构造参数逻辑完全一致。
2. 调整监听方法逻辑,兼容反序列化失败场景
反序列化失败时消息value为null,需要提前判空避免空指针,可按需记录错误信息:
@KafkaListener(topics = "#{'${topic}'}", groupId = "#{'${groupid}'}", autoStartup = "#{'${enabled}'}", containerFactory = "kafkaContainerFactory") public void updateExternalCardToken(ConsumerRecord<String, ExternalCardServiceData> record) { String key = record.key(); ExternalCardServiceData externalCardServiceData = record.value(); // 反序列化失败时value为null,直接记录错误后返回 if (externalCardServiceData == null) { log.error("消息反序列化失败,跳过处理,offset:{}", record.offset()); externalCardService.saveTokensErrors(key, "消息格式不合法,反序列化失败", "SerializationException"); return; } try { log.info("ExternalCardsListener. Received message: {}, offset={}", externalCardServiceData, record.offset()); externalCardService.updateToken(key, externalCardServiceData); } catch (Exception e) { log.error("Error processing message received from kafka. [Message={}]", externalCardServiceData); externalCardService.saveTokensErrors(key, externalCardServiceData.toString(), Arrays.toString(e.getStackTrace())); } }
配置说明
ErrorHandlingDeserializer为委托型反序列化器,本身不执行实际反序列化逻辑,仅捕获下层反序列化器抛出的所有异常,避免异常直接穿透到消费容器层导致消费线程卡死- 上述配置中
FixedBackOff(0L, 0L)代表异常触发后重试间隔0ms、最大重试次数0次,即失败后直接跳过,不会因为单条坏消息阻塞整个消费进度 - 反序列化失败的完整异常栈会写入消息Header,对应key为
ErrorHandlingDeserializer.VALUE_DESERIALIZER_EXCEPTION_HEADER,如果需要持久化完整异常信息,可以从该Header中读取序列化后的异常对象
内容的提问来源于stack exchange,提问作者Faik91
相关产品推荐
相关产品推荐

