Как сделать, чтобы заходило в обработчик ConsumerRecordRecoverer, после исчерпания всех повторных попыток обработки сообщений в кафке?
@Bean
public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> createKafkaListenerContainerFactory(Properties properties) {
ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
ConsumerFactory<String, String> consumerFactory = new DefaultKafkaConsumerFactory<>(consumerCreatePrepaidConfigs());
factory.setConsumerFactory(consumerFactory);
factory.setConcurrency(properties.getConcurrency);
ConsumerRecordRecoverer consumerRecordRecoverer = createConsumerRecordRecoverer();
factory.setErrorHandler(new SeekToCurrentErrorHandler(
consumerRecordRecoverer,
new FixedBackOff(
properties.getRetryInterval(),
properties.getMaxRetryAttempts())
)
);
return factory;
}
private static ConsumerRecordRecoverer createConsumerRecordRecoverer(){
return (consumerRecord, e) -> {
System.out.println("Число попыток исчерпано");
// Логика обработки
};
}
Я пытаюсь сделать повторные обработки сообщений, если по каким-то причинам произошёл Exception. Если количество попыток исчерпано, то запускается некая логика (например, отправка сообщения пользователю)
После исчерпания всех попыток не заходит в createConsumerRecordRecoverer()