spring-kafka消息重试时DELIVERY_ATTEMPT投递次数不递增问题
spring-kafka RetryTopic模式下投递计数始终为1问题排查
问题场景
使用spring-kafka 2.8.6版本,结合RetryTopicConfiguration实现非阻塞消息重试能力,初始监听器代码如下:
@KafkaListener( topics = "...", groupId = "...", containerFactory = "kafkaListenerContainerFactory") public void listenWithHeaders(final @Valid @Payload Event event, @Header(KafkaHeaders.DELIVERY_ATTEMPT) final int deliveryAttempt) { }
已配置公共错误处理器,同时开启了投递次数请求头功能,容器工厂配置代码如下:
@Bean public ConcurrentKafkaListenerContainerFactory<String, Event> kafkaListenerContainerFactory(@Qualifier("ConsumerFactory") final ConsumerFactory<String, Event> consumerFactory) { final ConcurrentKafkaListenerContainerFactory<String, Event> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.setCommonErrorHandler(new DefaultErrorHandler(new ExponentialBackOff(kafkaProperties.getExponentialBackoffInitialInterval(), kafkaProperties.getExponentialBackoffMultiplier()))); LOGGER.info("setup ConcurrentKafkaListenerContainerFactory"); factory.getContainerProperties().setDeliveryAttemptHeader(true); return factory; }
异常现象
触发消息重试时,从消息头中获取的投递次数始终为1,不会随重试次数递增。
补充配置
使用如下重试主题配置实现非阻塞重试:
@Bean public RetryTopicConfiguration retryableTopicKafkaTemplate(@Qualifier("kafkaTemplate") KafkaTemplate<String, Event> kafkaTemplate) { return RetryTopicConfigurationBuilder .newInstance() .exponentialBackoff( properties.getExponentialBackoffInitialInterval(), properties.getExponentialBackoffMultiplier(), properties.getExponentialBackoffMaxInterval()) .autoCreateTopics(properties.isRetryTopicAutoCreateTopics(), properties.getRetryTopicAutoCreateNumPartitions(), properties.getRetryTopicAutoCreateReplicationFactor()) .maxAttempts(properties.getMaxAttempts()) .notRetryOn(...) .retryTopicSuffix(properties.getRetryTopicSuffix()) .dltSuffix(properties.getDltSuffix()) .create(kafkaTemplate); }
根因说明
KafkaHeaders.DELIVERY_ATTEMPT是阻塞重试场景(即DefaultErrorHandler在当前消费者线程内直接重试消费)使用的投递计数头。使用RetryTopic非阻塞重试机制时,重试逻辑为将消息转发到独立的重试主题完成重新投递,投递次数计数存储在RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS对应的消息头中,原有代码绑定了错误的header key,因此无法获取到随重试递增的正确计数值。
另外该重试计数头在消息首次投递到主主题时不存在,需将参数类型设为包装类型Integer,同时设置header非必须,避免首次消费时触发参数绑定异常。
修复代码
调整监听器中获取投递次数的Header配置,调整后功能恢复正常:
@KafkaListener( topics = "...", groupId = "...", containerFactory = "kafkaListenerContainerFactory") public void listenWithHeaders(final @Valid @Payload Event event, @Header(value = RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTS, required = false) final Integer deliveryAttempt) { ... }
内容的提问来源于stack exchange,提问作者lansing
相关产品推荐
相关产品推荐

