Spring Kafka 2.8.x:如何用CommonErrorHandler跳过反序列化失败消息
解决Spring Kafka 2.8.x跳过反序列化失败消息的方案
核心思路
反序列化失败会触发CommonErrorHandler的handleDeserializationException方法(而非你之前重写的handleOtherException)。我们可以通过该方法从RecordDeserializationException中直接获取异常消息的元数据(分区、偏移量等),然后手动提交偏移量跳过这条错误消息,无需解析日志内容。
修改后的配置代码
@Bean public ConsumerFactory<String, Farewell> farewellConsumerFactory() { groupId = LocalTime.now().toString(); Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress); props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId); props.put(JsonDeserializer.TRUSTED_PACKAGES, "*"); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); // 关闭自动提交,避免自动提交错误偏移量 props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new JsonDeserializer<>(Farewell.class)); } @Bean public ConcurrentKafkaListenerContainerFactory<String, Farewell> farewellKafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, Farewell> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setCommonErrorHandler(new CommonErrorHandler() { @Override public void handleDeserializationException(ListenerExecutionFailedException ex, Consumer<?, ?> consumer, MessageListenerContainer container, boolean batchListener) { // 从异常链中提取反序列化异常 RecordDeserializationException deserializationEx = (RecordDeserializationException) ex.getCause(); ConsumerRecord<?, ?> failedRecord = deserializationEx.getRecord(); // 记录错误信息 System.err.printf("反序列化失败,跳过消息:topic=%s, partition=%d, offset=%d%n", failedRecord.topic(), failedRecord.partition(), failedRecord.offset()); // 手动提交下一个偏移量,跳过当前错误消息 Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>(); offsets.put(new TopicPartition(failedRecord.topic(), failedRecord.partition()), new OffsetAndMetadata(failedRecord.offset() + 1)); consumer.commitSync(offsets); } }); // 设置监听器为手动提交模式 factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL); factory.setConsumerFactory(farewellConsumerFactory()); return factory; }
监听器类调整
需要在监听器方法中手动提交正常消息的偏移量:
@KafkaListener(topics = "${topicId}", containerFactory = "farewellKafkaListenerContainerFactory") public void farewellListener(Farewell message, Acknowledgment acknowledgment) { System.out.println("Received Message in group " + groupId + "| " + message); // 手动提交正常消息的偏移量 acknowledgment.acknowledge(); }
关键说明
- 关闭自动提交:避免Kafka自动提交错误的偏移量,导致重复消费或误跳过正常消息。
- 分场景提交偏移量:正常消息通过
Acknowledgment接口提交,错误消息直接提交offset+1的偏移量,确保跳过当前失败消息。 - 直接获取元数据:通过
RecordDeserializationException的getRecord()方法直接拿到失败消息的topic、分区、偏移量,完全不需要解析日志。
内容的提问来源于stack exchange,提问作者Knowledge_seeker
相关产品推荐
相关产品推荐

