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; } }
问题原因
- 双重重试逻辑冲突:同时配置了listener层面的
RetryTemplate和容器层面的SeekToCurrentErrorHandler,两者都会处理异常重试。当RetryTemplate耗尽重试次数触发RecoveryCallback并抛出异常后,SeekToCurrentErrorHandler会捕获异常并基于偏移量回退再次触发消费,进而再次触发RetryTemplate的重试流程,最终导致RecoveryCallback执行两次。 - ACK逻辑矛盾:RecoveryCallback中手动提交偏移量,但
SeekToCurrentErrorHandler仍会尝试回退偏移量,导致消息被重新拉取消费。 - 反序列化配置错误:重复设置
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
相关产品推荐
相关产品推荐

