Как после consumer rebalance процесса, сбросить offset, что бы Consumer мог прочитать топик сначала?

Есть ситуация, когда Consumer1 читает сообщения из кафка топика. При подключении второго Consumer2 с таким же groupId, происходит перераспределение партиций, можно ли как то сбросить offset, что бы после процесса перераспределения, оба Consumer'а читали топик с начала?

Мне нужно что бы во время процесса перерапределения партиций, выполнялась определенная логика, в том числе что бы после каждого перераспределения партиций сбрасывалось значение offset, что бы после перераспределения информация равномерно распределилась между потребителями. Я нашел, что процесс ре-балансировки партиций модно отследить реализовав интерфейс ConsumerRebalanceListener

@Service
public class KafkaRebalanceListenerHandler implements ConsumerRebalanceListener {

    @Override
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {

    }

    @Override
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {

    }
}

И далее указать переопределенный KafkaRebalanceListenerHandler при настройке KafkaConsumer

    @Bean
    @RefreshScope
    public ConcurrentKafkaListenerContainerFactory<String, String> listenerFactory() {
        final ConcurrentKafkaListenerContainerFactory<String, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setErrorHandler(kafkaListenerErrorHandler);
        factory.getContainerProperties().setConsumerRebalanceListener(kafkaRebalanceListenerHandler);
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

с этим вроде тоже разобрался. Остается неясность, как привязавшись к перераспределению партиций, застравить вычитывать топик заново.


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

Автор решения: Roman Konoval

По идее достаточно реализовать ConsumerSeekAware в вашем listener-e, а именно метод onPartitionsAssigned, чтоб он сбрасывал offset-ы на начало, когда партиция назначается данному Consumer-у:

public class KafkaMessageListener implements ConsumerSeekAware {
    @KafkaListener(topics = "my.topic")
    public void listen(byte[] payload) {
        // ...
    }

    @Override
    public void registerSeekCallback(ConsumerSeekCallback callback) {
    }

    @Override
    public void onPartitionsAssigned(Map<TopicPartition, Long> assignments, ConsumerSeekCallback callback) {
        assignments.forEach((t, o) -> callback.seekToBeginning(t.topic(), t.partition()));
    }

    @Override
    public void onIdleContainer(Map<TopicPartition, Long> assignments, ConsumerSeekCallback callback) {
    }
}
→ Ссылка