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

spring-kafka 2.7.0中@RetryableTopic引发监听器无限循环问题

Spring Kafka 2.7.0 @RetryableTopic 无限循环问题

在Spring Kafka 2.7.0版本中使用@RetryableTopic注解时,触发异常后监听器陷入无限循环。重试逻辑能将消息推送到正确的重试主题,但会持续重复读取同一条消息。尝试调整配置中注释掉的相关参数,问题仍未解决。

配置代码(KafkaConfiguration.java)

@Configuration
//@EnableKafkaRetryTopic
@EnableKafka
@Slf4j
@ConditionalOnProperty(value = "kafka.config.enabled", havingValue = "true", matchIfMissing = false)
public class KafkaConfiguration {
  
  @Autowired
  ReformerKafkaProperties kafkaProperties;
  
  @Bean
  public ProducerFactory<String, String> producerFactory() {
    Map<String, Object> configProps = new HashMap<>();
    configProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers());
    configProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    configProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    configProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, Object.class);
    configProps.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
    return new DefaultKafkaProducerFactory<>(configProps);
  }
  
  @Bean
  public KafkaTemplate<String, String> kafkaTemplate() {
    return new KafkaTemplate<>(producerFactory());
  }

  @Bean
  public ConsumerFactory<String, String> consumerFactory() {
    Map<String, Object> consumerProps = new HashMap<>();
    consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaProperties.getBootstrapServers());
//    consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true);
//    consumerProps.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 1000);
    consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaProperties.getConsumerGroup());
    consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
    consumerProps.put(JsonDeserializer.VALUE_DEFAULT_TYPE, Object.class);
    consumerProps.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
    return new DefaultKafkaConsumerFactory<>(consumerProps);
  }

  @Bean
  public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
    factory.setConsumerFactory(consumerFactory());
//    factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD);
//    factory.setErrorHandler(new SeekToCurrentErrorHandler(new FixedBackOff(0L, 0)));
    return factory;
  }

}

监听器代码(AsyncReformationKafkaListener.java)

@Slf4j
@Component
@ConditionalOnProperty(value = "kafka.config.enabled", havingValue = "true", matchIfMissing = false)
public class AsyncReformationKafkaListener {

  @RetryableTopic(
      backoff = @Backoff(value = 30000),
      attempts = "3",
      autoCreateTopics = "true",
      include = RuntimeException.class)
  @KafkaListener(
      topics = "${kafka.config.async-topic}", 
      groupId = "${kafka.config.consumer-group}")
  public void executeReformation(String messageString, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) {
    log.info("Kafka msg at topic {}", topic);
    throw new RuntimeException();
  }

}

主题偏移量信息

仅在reformer-async-local主题生产了一条消息,但重试主题中的消息数量呈指数级增长:

a13400412@MAC-a13400412 bin % sh kafka-get-offsets.sh --bootstrap-server=localhost:9092 --topic=reformer-async-local-retry-1
reformer-async-local-retry-1:0:1492
a13400412@MAC-a13400412 bin % sh kafka-get-offsets.sh --bootstrap-server=localhost:9092 --topic=reformer-async-local-retry-0
reformer-async-local-retry-0:0:175
a13400412@MAC-a13400412 bin % sh kafka-get-offsets.sh --bootstrap-server=localhost:9092 --topic=reformer-async-local        
reformer-async-local:0:1
a13400412@MAC-a13400412 bin % 

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 07:55:08