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

Kafka反序列化异常时Offset仍自动提交的问题求助

Kafka反序列化异常时Offset仍自动提交的问题求助

大家好,我最近在测试Kafka的反序列化错误处理逻辑时遇到了一个棘手的问题,想请各位帮忙排查下:

我故意使用错误的反序列化器来测试Kafka对反序列化异常的处理,结果发现即便触发了反序列化错误,Kafka主题的Offset依然会自动提交,完全检测不到消费滞后。我已经尝试配置了ErrorHandlingDeserializer和DefaultErrorHandler,但问题还是没有解决,Offset还是会被自动提交。

以下是我的相关配置和代码:

application.yaml配置

server:
  port: 9292
spring:
  kafka:
    consumer:
      group-id: consumer-group-1
      enable-auto-commit: false
      bootstrap-servers: localhost:9092
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
    properties:
      spring.json.type.mapping: com.test.kafka_producer.dto.BankTransferEvent:com.test.kafka_consumer.dto.BankTransferEvent
      spring.deserializer.value.delegate.class: org.apache.kafka.common.serialization.StringDeserializer
      spring.json.trusted.packages: com.test.kafka_consumer.dto,com.test.kafka_producer.dto

Kafka消费者配置类

package com.test.kafka_consumer.configuration;

import com.test.kafka_consumer.dto.BankTransferEvent;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.util.backoff.FixedBackOff;

@Configuration
public class KafkaConsumerConfig {

    @Bean
    public DefaultErrorHandler errorHandler() {
        // 设置重试策略:间隔1秒,重试2次
        DefaultErrorHandler errorHandler = new DefaultErrorHandler(new FixedBackOff(1000L, 2L));
        // 恢复后不提交Offset
        errorHandler.setCommitRecovered(false);
        return errorHandler;
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, BankTransferEvent> kafkaListenerContainerFactory(
            ConsumerFactory<String, BankTransferEvent> consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<String, BankTransferEvent> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        // 启用手动确认模式
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
        // 给容器绑定错误处理器
        factory.setCommonErrorHandler(errorHandler());
        return factory;
    }
}

Kafka监听类

package com.test.kafka_consumer.service;

import com.test.kafka_consumer.dto.BankTransferEvent;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Service;

@Service
@Slf4j
public class KafkaConsumer {

    @KafkaListener(topics = "topic-json",groupId = "consumer-group-2", containerFactory = "kafkaListenerContainerFactory")
    public void consumeJsonEvent1(BankTransferEvent bankTransferEvent, Acknowledgment ack){
        log.info("Kafka Consumer aufgerufen");
        try {
            log.info("consumer-JSON-1 consume the event: {" + bankTransferEvent.toString() + "}");
            ack.acknowledge();
        } catch (Exception e){
            log.error("Error bei der Event-Verarbeitung: " + e.getMessage());
        }
    }
}

我已经确认enable-auto-commit设为了false,也开启了手动确认模式,并且在错误处理器里禁用了恢复后的Offset提交,但反序列化异常发生时,Offset还是会被自动提交,导致异常消息直接被跳过,既没有重试也没有保留消费滞后。有没有朋友能帮我看看哪里配置出问题了?

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 08:44:30