Не все потребители кафки получают назначение на разделы - PullRequest
1 голос
/ 14 февраля 2020

У меня 10 потребителей и 10 разделов. Я беру количество разделов

int partitionCount = getPartitionCount(kafkaUrl);

и создаю одинаковое количество потребителей с тем же group.id .

    public void listen() {
        try {
            String kafkaUrl = getKafkaUrl();
            int partitionCount = getPartitionCount(kafkaUrl);
            Stream.iterate(0, i -> i + 1)
                    .limit(partitionCount)
                    .forEach(index -> executorService.execute(() ->
                            consumerTask.invokeKafkaConsumerTask(prepareConsumerConfig(index, kafkaUrl), INPUT_TOPIC)));
        } catch (Exception exception) {
            logger.error("Cannot receive event from kafka ", exception);
        }


    public void invokeKafkaConsumerTask(Properties properties, String topicName) {
        try(KafkaConsumer<String, String> consumer = new KafkaConsumer<>(properties)) {
            consumer.subscribe(Collections.singletonList(topicName));
            logger.info("[KAFKA] consumer created");
            invokeKafkaConsumer(consumer);
        } catch (IllegalArgumentException exception) {
            logger.error("Cannot create kafka consumer ", exception);
        }
    }

    private void invokeKafkaConsumer(KafkaConsumer<String, String> consumer) {
        try {
            while (true) {
                ConsumerRecords<String, String> consumerRecords = consumer.poll(Duration.ofSeconds(4));
                if (consumerRecords.count() > 0) {
                    consumeRecords(consumerRecords);
                    consumer.commitSync();
                }
            }
        } catch (Exception e) {
            logger.error("Error while receiving records ", e);
        }
    }

метод getPartitionCount возвращает 10 разделов, чтобы он работал правильно

config выглядит следующим образом

Properties properties = new Properties();
properties.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaUrl);
properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
properties.put(ConsumerConfig.GROUP_ID_CONFIG, CONSUMER_CLIENT_ID);
properties.put(ConsumerConfig.CLIENT_ID_CONFIG, CONSUMER_CLIENT_ID + index);
properties.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
properties.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "300000");
properties.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, "10000");
properties.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, String.valueOf(Integer.MAX_VALUE));
properties.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.RoundRobinAssignor");

что я вижу после назначения потребителей разделу

TOPIC      PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG             CLIENT-ID                                                        
topicName  1          89391           89391           0               consumer0
topicName  3          88777           88777           0               consumer1
topicName  5          89280           89280           0               consumer2
topicName  4          88776           88776           0               consumer2
topicName  0          4670991         4670991         0               consumer0
topicName  9          23307           89343           66036           consumer4
topicName  7          89610           89610           0               consumer3
topicName  8          88167           88167           0               consumer4
topicName  2          89138           89138           0               consumer1
topicName  6          88967           88967           0               consumer3

только половина потребители были назначены на перегородки, почему это произошло? В соответствии с документацией должен быть один потребитель на раздел. Я делаю что-то неправильно? kafka версия 2.1.1 .

Я также нахожу несколько таких журналов ->

Setting newly assigned partitions:[empty]

Ответы [ 2 ]

0 голосов
/ 18 февраля 2020

[решение] интересный случай Я изменил group.id и partition.assignment.strategy, добавил auto.offset.reset = самое раннее и похоже, что он работает ...

0 голосов
/ 14 февраля 2020

Вы подписываетесь на коллекцию topi c name или java Pattern?
Если вы подписываетесь на Pattern, измените partition.assignment.strategy на RoundRobinAssignor или StickyAssignor.

Добро пожаловать на сайт PullRequest, где вы можете задавать вопросы и получать ответы от других членов сообщества.
...