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

Spring Kafka 2.2.0重试时RecoveryCallback重复执行问题求助

Spring Kafka 2.2.0 重试后转DLT异常问题解决

问题描述

在Spring Kafka 2.2.0中配置主主题消息重试n次,达到最大重试次数后将消息发送至DLT(死信主题)。预期流程为:n次重试→触发RecoveryCallback→抛出自定义异常→SeekToCurrentErrorHandler捕获异常并转发消息到DLT。但实际流程是:n次重试→RecoveryCallback→再次n次重试→RecoveryCallback→ErrorHandler,RecoveryCallback被执行两次,不符合预期。

原配置代码

public class KafkaQueueConfiguration {
    private static Logger log = LoggerFactory.getLogger(KafkaQueueConfiguration.class);

    @Autowired
    private StringRedisTemplate redisTemplate;

    @Value("${bootstrap.ip}")
    private String consumerBootstrap;

    @Value("${group.id}")
    private String consumerGroupId;

    @Value("${consumer.offset.type}")
    private String autoOffsetType;

    @Value("${kafka.backoff.interval}")
    private Long fixedInterval;

    @Value("${kafka.maximum.poll.records}")
    private Integer maxPollRecords;

    @Value("${bank-notification-sasl-jaas-config}")
    private String saslJaasConfig;

    @Value("${bank-notification-json-type-mapping}")
    private String jsonTypeMapping;

    @Value("${bank-notification-dejson-delegate}")
    private String deJsonDelegate;

    @Value("${kafka.topics.bank.notification.dlt}")
    private String dlqName;

    @Value("${kafka-consumer-saslMechanism}")
    private String saslMechanism;

    @Value("${kafka-consumer-securityProtocol}")
    private String securityProtocol;

    @Value(("${heartbeat-interval-ms}"))
    private String heartbeatInterval;

    @Value(("${session-timeout-ms}"))
    private String sessionTimeout;

    @Value(("${max-poll-interval-ms}"))
    private Integer maxPollInterval;

    @Value(("${max-retry-attempt-ms}"))
    private Integer maxRetryAttempts;

    @Value(("${retry-interval-ms}"))
    private long retryInterval;

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumerBootstrap);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetType);
        props.put(ConsumerConfig.RETRY_BACKOFF_MS_CONFIG,fixedInterval);
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG,maxPollInterval);
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG,maxPollRecords);
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionTimeout);
        props.put(SaslConfigs.SASL_JAAS_CONFIG,saslJaasConfig);
        props.put(JsonDeserializer.TYPE_MAPPINGS,jsonTypeMapping);
        props.put(SaslConfigs.SASL_MECHANISM,saslMechanism);
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG,securityProtocol);
        props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS,deJsonDelegate);
        props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG,heartbeatInterval);
        props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG ,SslConfigs.DEFAULT_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM);
        return props;
    }
    @Bean
    public ConsumerFactory<Object, Object> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactory(
            ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
            ConsumerFactory<Object, Object> kafkaConsumerFactory, KafkaTemplate<Object, Object> template) {

        ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();

        factory.setConsumerFactory(consumerFactory());
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
        factory.getContainerProperties().setAckOnError(false);
        factory.setRetryTemplate(retryTemplate());

        factory.setRecoveryCallback((context->{
            log.info("Recovery callback");
            Acknowledgment ack= (Acknowledgment) context.getAttribute(RetryingMessageListenerAdapter.CONTEXT_ACKNOWLEDGMENT);
            
            log.info("Manually committed..{}",countRecover);
            log.info("RetryCount :: {} and last retry exception :: {}",context.getRetryCount(),context.getLastThrowable().getCause().getMessage());
            log.info("Retry records :: {}",context.getAttribute(RetryingMessageListenerAdapter.CONTEXT_RECORD));
            ConsumerRecord consumerRecord= (ConsumerRecord) context.getAttribute(RetryingMessageListenerAdapter.CONTEXT_RECORD);
            EftTransactionDetail eftTransactionDetail= (EftTransactionDetail) consumerRecord.value();
            log.info("Retrying records ::{}",eftTransactionDetail);
            redisTemplate.opsForValue().setIfAbsent(BankNotificationsConstants.SMS_REDIS_PREFIX + eftTransactionDetail.getAtdEntryId(),eftTransactionDetail.getAtdEntryId());
            ack.acknowledge();
            throw new ValidationException(SystemTag.HG,context.getLastThrowable().getCause().getMessage(),context.getLastThrowable().getCause().getMessage());
        }));
        factory.setErrorHandler(errorHandler(publisher(template)));
        configurer.configure(factory, kafkaConsumerFactory);
        return factory;
    }

    @Bean
    public SeekToCurrentErrorHandler errorHandler(DeadLetterPublishingRecoverer deadLetterPublishingRecoverer) {
        
        return new SeekToCurrentErrorHandler(deadLetterPublishingRecoverer,0);//set maxfailure 0 attempts

    }
    @Bean
    public DeadLetterPublishingRecoverer publisher(KafkaTemplate<Object, Object> template) {
        return new DeadLetterPublishingRecoverer(template);
    }

    public RetryTemplate retryTemplate(){

        RetryTemplate retryTemplate = new RetryTemplate();
        retryTemplate.setRetryPolicy(new SimpleRetryPolicy(maxRetryAttempts));
        FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy();
        backOffPolicy.setBackOffPeriod(retryInterval);
        retryTemplate.setBackOffPolicy(backOffPolicy);
        return retryTemplate;
    }
}

问题原因

  1. 双重重试逻辑冲突:同时配置了listener层面的RetryTemplate和容器层面的SeekToCurrentErrorHandler,两者都会处理异常重试。当RetryTemplate耗尽重试次数触发RecoveryCallback并抛出异常后,SeekToCurrentErrorHandler会捕获异常并基于偏移量回退再次触发消费,进而再次触发RetryTemplate的重试流程,最终导致RecoveryCallback执行两次。
  2. ACK逻辑矛盾:RecoveryCallback中手动提交偏移量,但SeekToCurrentErrorHandler仍会尝试回退偏移量,导致消息被重新拉取消费。
  3. 反序列化配置错误:重复设置VALUE_DESERIALIZER_CLASS_CONFIG,先指定ErrorHandlingDeserializer后又覆盖为JsonDeserializer,导致ErrorHandlingDeserializer未生效,可能引发额外异常处理问题。

解决方案

核心调整

移除listener层面的重试配置,统一由SeekToCurrentErrorHandler处理重试和DLT转发,修正反序列化配置,调整ACK模式简化偏移量管理。

修改后的配置代码

public class KafkaQueueConfiguration {
    private static Logger log = LoggerFactory.getLogger(KafkaQueueConfiguration.class);

    @Autowired
    private StringRedisTemplate redisTemplate;

    @Value("${bootstrap.ip}")
    private String consumerBootstrap;

    @Value("${group.id}")
    private String consumerGroupId;

    @Value("${consumer.offset.type}")
    private String autoOffsetType;

    @Value("${kafka.backoff.interval}")
    private Long fixedInterval;

    @Value("${kafka.maximum.poll.records}")
    private Integer maxPollRecords;

    @Value("${bank-notification-sasl-jaas-config}")
    private String saslJaasConfig;

    @Value("${bank-notification-json-type-mapping}")
    private String jsonTypeMapping;

    @Value("${bank-notification-dejson-delegate}")
    private String deJsonDelegate;

    @Value("${kafka.topics.bank.notification.dlt}")
    private String dlqName;

    @Value("${kafka-consumer-saslMechanism}")
    private String saslMechanism;

    @Value("${kafka-consumer-securityProtocol}")
    private String securityProtocol;

    @Value("${heartbeat-interval-ms}")
    private String heartbeatInterval;

    @Value("${session-timeout-ms}")
    private String sessionTimeout;

    @Value("${max-poll-interval-ms}")
    private Integer maxPollInterval;

    @Value("${max-retry-attempt-ms}")
    private Integer maxRetryAttempts;

    @Value("${retry-interval-ms}")
    private long retryInterval;

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, consumerBootstrap);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 修正反序列化配置:用ErrorHandlingDeserializer包裹JsonDeserializer
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ErrorHandlingDeserializer.class);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroupId);
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetType);
        props.put(ConsumerConfig.RETRY_BACKOFF_MS_CONFIG, fixedInterval);
        props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, maxPollInterval);
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, maxPollRecords);
        props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, sessionTimeout);
        props.put(SaslConfigs.SASL_JAAS_CONFIG, saslJaasConfig);
        props.put(JsonDeserializer.TYPE_MAPPINGS, jsonTypeMapping);
        props.put(SaslConfigs.SASL_MECHANISM, saslMechanism);
        props.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol);
        props.put(ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS, deJsonDelegate);
        props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, heartbeatInterval);
        props.put(SslConfigs.SSL_ENDPOINT_IDENTIFICATION_ALGORITHM_CONFIG, SslConfigs.DEFAULT_SSL_ENDPOINT_IDENTIFICATION_ALGORITHM);
        return props;
    }

    @Bean
    public ConsumerFactory<Object, Object> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<Object, Object> kafkaListenerContainerFactory(
            ConcurrentKafkaListenerContainerFactoryConfigurer configurer,
            ConsumerFactory<Object, Object> kafkaConsumerFactory, KafkaTemplate<Object, Object> template) {

        ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();

        factory.setConsumerFactory(consumerFactory());
        // 调整ACK模式为RECORD,简化偏移量管理
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD);
        factory.getContainerProperties().setAckOnError(false);
        // 移除listener层面的RetryTemplate和RecoveryCallback
        factory.setErrorHandler(errorHandler(publisher(template)));
        configurer.configure(factory, kafkaConsumerFactory);
        return factory;
    }

    @Bean
    public SeekToCurrentErrorHandler errorHandler(DeadLetterPublishingRecoverer deadLetterPublishingRecoverer) {
        // 配置重试后退策略,替代原RetryTemplate的逻辑
        FixedBackOff backOff = new FixedBackOff(retryInterval, maxRetryAttempts - 1);
        // 初始化ErrorHandler,指定DLT转发器和重试规则
        return new SeekToCurrentErrorHandler(deadLetterPublishingRecoverer, backOff);
    }

    @Bean
    public DeadLetterPublishingRecoverer publisher(KafkaTemplate<Object, Object> template) {
        // 自定义DLT主题映射,确保消息发送到指定死信主题
        return new DeadLetterPublishingRecoverer(template, (record, ex) -> new TopicPartition(dlqName, record.partition()));
    }
}

额外说明

  • 若需在重试耗尽后执行自定义逻辑(如原RecoveryCallback中的Redis存储操作),可扩展DeadLetterPublishingRecoverer,或在SeekToCurrentErrorHandler构造时传入自定义RecoveryCallback。
  • SeekToCurrentErrorHandler的maxAttempts包含首次消费,配置时需注意:若需n次重试,应设置maxRetryAttempts - 1作为后退策略的最大尝试次数。

内容的提问来源于stack exchange,提问作者Deepanshu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 01:15:38