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
相关产品推荐
相关产品推荐

