Как сделать, чтобы заходило в обработчик 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()


Ответы (0 шт):