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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 07:35:23