Примеры использования KafkaConsumer#

API потребителя предоставляет определенную гибкость для большого количества сценариев использования. В подразделах приведены примеры таких кейсов.

Автоматическая фиксация смещения (offset)#

Пример ниже показывает, как работает API потребителя Corax при автоматической фиксации смещения:

Properties props = new Properties();
     props.setProperty("bootstrap.servers", "localhost:9092");
     props.setProperty("group.id", "test");
     props.setProperty("enable.auto.commit", "true");
     props.setProperty("auto.commit.interval.ms", "1000");
     props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
     props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
     KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
     consumer.subscribe(Arrays.asList("foo", "bar"));
     while (true) {
         ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
         for (ConsumerRecord<String, String> record : records)
             System.out.printf("offset = %d, key = %s, value = %s%n", record.offset(), record.key(), record.value());
     }

Соединение с кластером инициализируется при помощи указания списка из одного или нескольких брокеров в конфигурационном параметре bootstrap.servers. Данный список используется для обнаружения других брокеров в кластере. Он не должен быть исчерпывающим списком серверных узлов в кластере, однако есть возможность указать больше одного серверного узла на случай, если при подключении клиента к кластеру некоторые серверы будут offline.

Настройка enable.auto.commit означает, что смещения фиксируются автоматически с частотой, которая контролируется конфигурационным параметром auto.commit.interval.ms.

В данном примере потребитель подписывается на топики foo и bar как часть группы потребителей, имеющей название test, согласно параметру group.id.

Настройки десериализатора определяют, каким образом необходимо трансформировать байты в объекты. Например, при передаче строковых десериализаторов пары «ключ-значение» станут простыми строками.

Ручная фиксация смещений#

С помощью ручной фиксации смещений пользователь может самостоятельно определять момент, когда запись считается потребленной, и фиксировать ее смещение. Это может быть полезным в случаях, когда потребление сообщений объединено с логикой обработки и, таким образом, сообщение не может считаться потребленным, пока не завершена его обработка:

Properties props = new Properties();
     props.setProperty("bootstrap.servers", "localhost:9092");
     props.setProperty("group.id", "test");
     props.setProperty("enable.auto.commit", "false");
     props.setProperty("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
     props.setProperty("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
     KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
     consumer.subscribe(Arrays.asList("foo", "bar"));
     final int minBatchSize = 200;
     List<ConsumerRecord<String, String>> buffer = new ArrayList<>();
     while (true) {
         ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
         for (ConsumerRecord<String, String> record : records) {
             buffer.add(record);
         }
         if (buffer.size() >= minBatchSize) {
             insertIntoDb(buffer);
             consumer.commitSync();
             buffer.clear();
         }
     }

В данном примере потребляется определенный набор записей и перемещается в память. Как только в памяти соберется достаточно большой набор записей, они будут перемещены в БД. Если бы фиксация происходила автоматически, как в предыдущем примере, записи бы считались потребленными после их возвращения пользователю в poll. Это потенциально может привести к сбоям после сбора записей в один набор, но до отправки их в БД.

Чтобы избежать этого, нужно вручную зафиксировать смещения только после отправки соответствующих записей в БД. Это позволит четко проконтролировать, когда запись стала считаться потребленной. Однако это может привести к появлению противоположной возможности: сбой может произойти в промежутке, когда данные уже были отправлены в БД, но смещения еще не были зафиксированы (даже если такой промежуток составил всего несколько миллисекунд, возможность возникновения сбоев сохраняется). В этом случае потребляющий процесс будет потреблять сообщения из последнего зафиксированного смещения и повторять вставку последнего набора данных.

В таком режиме Corax предоставляет гарантии доставки сообщений at-least-once, что означает, что сообщение, скорее всего, будет доставлено, но при ошибке возможно его задвоение.

Внимание

Использование автоматической фиксации смещений также может предоставлять гарантии доставки сообщений at-least-once, но только при потреблении всех данных, возвращенных от каждого вызова методу poll(Duration) до начала любого последующего вызова, либо до закрытия (close) потребителя. Если какое-то из этих условий не будет выполнено, то зафиксированное смещение вполне может опередить позицию потребленных сообщений, что может привести к потере записей.

Преимуществом ручной фиксации является непосредственный контроль над тем, когда запись может считаться «потребленной».

Приведенный выше пример использует метод commitSync для маркировки всех полученных записей как «зафиксированных». В некоторых случаях может понадобиться еще более четкий контроль того, какие записи фиксируются. Этого можно добиться, явно установив значение смещения. Пример ниже показывает, как зафиксировать смещение после окончания обработки записей в каждой партиции:

try {
         while(running) {
             ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(Long.MAX_VALUE));
             for (TopicPartition partition : records.partitions()) {
                 List<ConsumerRecord<String, String>> partitionRecords = records.records(partition);
                 for (ConsumerRecord<String, String> record : partitionRecords) {
                     System.out.println(record.offset() + ": " + record.value());
                 }
                 long lastOffset = partitionRecords.get(partitionRecords.size() - 1).offset();
                 consumer.commitSync(Collections.singletonMap(partition, new OffsetAndMetadata(lastOffset + 1)));
             }
         }
     } finally {
       consumer.close();
     }

Внимание

Зафиксировано всегда должно быть смещение следующего сообщения, которое будет прочитано вашим приложением. То есть при вызове метода commitSync(offsets) необходимо добавить к смещению последнего обработанного сообщения 1.

Ручное назначение партиций#

В примерах выше выполнена подписка на интересующие топики и позволено Corax динамически назначить каждому топику соответствующие партиции, основываясь на количестве активных потребителей в группе. Однако в некоторых случаях может понадобиться более четкий контроль над определенными партициями, если, например:

  • процесс обслуживает какое-то локальное состояние, ассоциируемое с этой партицией (например, локальное дисковое хранилище пар «ключ-значение»), и тогда он должен получать только записи для партиции диска, которую он обслуживает;

  • процесс сам по себе полностью доступен и при сбое восстанавливается (например, при помощи фреймворков управления кластером — YARN, Mesos или AWS либо как часть фреймворка обработки потока). В этом случае Corax не нужно обнаруживать сбой и перераспределять партиции, поскольку процесс-потребитель будет перезапущен на другой машине.

Для использования данного режима вместо подписки на топик с использованием метода subscribe можно просто вызвать метод assign(Collection) с полным списком партиций, которые необходимо потребить:

String topic = "foo";
     TopicPartition partition0 = new TopicPartition(topic, 0);
     TopicPartition partition1 = new TopicPartition(topic, 1);
     consumer.assign(Arrays.asList(partition0, partition1));

После назначения партиций можно сделать петлю с вызовом метода poll, как было описано в примерах выше. Группа, которую определяет потребитель, все еще используется для фиксации смещений, но теперь набор партиций будет изменяться только при следующем обращении к методу assign. Ручное назначение партиций не использует координацию групп, поэтому сбои потребителя не приведут к перебалансировке назначенных партиций. Каждый потребитель действует независимо, даже если имеет общий groupId с другим потребителем. Для избежания конфликтов фиксации смещений необходимо всегда проверять уникальность groupId для каждого экземпляра потребителя.

Внимание

Комбинировать использование ручного назначения партиций (например, с использованием метода assign) с динамическим (например, с использованием метода subscribe) невозможно.

Сохранение смещений вне Corax#

Приложение-потребитель не обязательно должно использовать встроенное хранилище смещений, оно может сохранять смещения в любом другом месте. Основным сценарием использования для этого является использование приложения для хранения смещения и результатов потребления в одной и той же системе таким образом, чтобы и результаты потребления, и смещения сохранялись атомарно. Это не всегда возможно, но это делает потребление более атомарным и позволяет получить семантику доставки сообщений exactly once, которая сильнее семантики at-least-once, которую пользователь получает, благодаря функциональности фиксирования смещений.

Ниже приведены примерные сценарии такого использования:

  • Если результаты потребления сохраняются в реляционной БД, сохранение в ней смещения может дать возможность фиксировать результаты потребления и смещение в одной транзакции. Следовательно, либо транзакция окажется успешной и смещение будет обновлено, основываясь на том, какие записи были потреблены, либо результат не будет сохранен и смещение не будет обновлено.

  • Если результаты сохраняются в локальное хранилище, то возможно будет сохраняться в нем и смещение. Например, индекс поиска может быть настроен при подписке на определенную партицию и одновременном сохранении и смещения, и индексированных данных. При атомарном сохранении часто возникает такая возможность, и даже при возникновении сбоев с потерей несохраненных данных, все, что осталось после сбоя, будет иметь сохраненное смещение. Это означает, что в этом случае процесс индексации, вернувшийся с потерей последних обновлений, просто продолжает индексирование с того момента, где обновления не были потеряны.

Каждая запись имеет собственное смещение, поэтому для управления смещением необходимо сделать следующее:

  • сконфигурировать параметр следующим образом: enable.auto.commit=false;

  • использовать смещение, предоставляемое каждым экземпляром класса ConsumerRecord для сохранения текущей позиции;

  • восстановить позицию потребителя, используя метод seek(TopicPartition, long).

Такой тип использования является наиболее простым при ручном назначении партиций. Если назначение партиций происходит автоматически, необходимо уделить дополнительное внимание перебалансировке партиций. Это можно сделать, передав экземпляр слушателя ConsumerRebalanceListener в методы subscribe(Collection, ConsumerRebalanceListener) и subscribe(Pattern, ConsumerRebalanceListener). Например, если партиции берутся из потребителя, он будет пытаться зафиксировать свое смещение для этих партиций, применяя метод ConsumerRebalanceListener.onPartitionsRevoked(Collection). При назначении партиций на потребителя он будет искать смещение для этих новых партиций и корректным образом инициализировать потребителя на эту позицию при помощи метода ConsumerRebalanceListener.onPartitionsAssigned(Collection).

Другим распространенным случаем применения слушателя ConsumerRebalanceListener является освобождение обслуживаемых приложением кешей для перемещенных куда-либо партиций.

Многопоточная обработка#

Внимание

Потребитель Corax не имеет средств безопасности для потоков. Весь сетевой I/O происходит в потоке вызывающего приложения. Обеспечение корректной синхронизации нескольких потоков является ответственностью пользователя. Несинхронизированный доступ приведет к появлению исключения ConcurrentModificationException.

Единственным исключением для вышеописанного правила является метод wakeup(), который может использоваться из внешнего потока для безопасного прерывания активной операции. В этом случае поток, блокирующий операцию, выдаст исключение WakeupException. Такой способ может использоваться для отключения потребителя из других потоков. Код ниже отражает типичный случай применения:

 public class KafkaConsumerRunner implements Runnable {
     private final AtomicBoolean closed = new AtomicBoolean(false);
     private final KafkaConsumer consumer;

     public KafkaConsumerRunner(KafkaConsumer consumer) {
       this.consumer = consumer;
     }

     @Override
     public void run() {
         try {
             consumer.subscribe(Arrays.asList("topic"));
             while (!closed.get()) {
                 ConsumerRecords records = consumer.poll(Duration.ofMillis(10000));
                 // Handle new records
             }
         } catch (WakeupException e) {
             // Ignore exception if closing
             if (!closed.get()) throw e;
         } finally {
             consumer.close();
         }
     }

     // Shutdown hook which can be called from a separate thread
     public void shutdown() {
         closed.set(true);
         consumer.wakeup();
     }
 }

Затем, в отдельном потоке, потребитель может быть отключен с помощью добавления атрибута closed и вызова метода wakeup() для потребителя:

closed.set(true);
consumer.wakeup();

Примечание

Поскольку существует возможность использовать прерывания потока вместо использования метода wakeup() для прерывания блокирующей операции (в этом случае появится исключение InterruptException), использовать их не рекомендуется, так как это может привести к прерыванию «чистого» отключения потребителя. Прерывания можно использовать для случаев, в которых использование метода wakeup() не представляется возможным, например, когда поток потребителя управляется кодом, который не «знает» о наличии клиента Corax.

Такая модель поточности для обработки не была внедрена специально, так как это оставляет несколько возможностей для применения многопоточной обработки записей:

  • Один потребитель — один поток.

    Это простая возможность давать каждому потоку свой экземпляр потребителя.

    Преимущества:

    • Легкость реализации.

    • Часто такой подход является наиболее быстрым и не требует внутрипоточной координации.

    • Такой подход облегчает реализацию порядковой обработки на основе порядка партиций (каждый поток просто обрабатывает сообщения по порядку их получения).

    Недостатки:

    • Большее количество потребителей означает большее количество TCP-подключений к кластеру (по одному на поток). Однако Corax эффективно управляет подключениями, поэтому это не сильно повлияет на работу.

    • Большее количество потребителей означает большее количество запросов, отправляемых на сервер, и немного меньшее количество пакетов с данными, что может привести к падению пропускной способности I/O.

    • Общее количество потоков по всем процессам будет ограничено общим количеством партиций.

  • Разделение потребления и обработки записей.

    Другой альтернативой является наличие одного или нескольких потоков-потребителей, которые потребляют все данные и передают экземпляры ConsumerRecords в блокирующую очередь, потребленную пулом потоков-обработчиков, которые и выполняют обработку записей.

    Преимущество: позволяет масштабировать количество потребителей и обработчиков независимо друг от друга. Это дает возможность иметь один потребитель, который обслуживает несколько потоков-обработчиков, что снимает ограничение по партициям.

    Недостатки:

    • Гарантии порядка обработчиков требуют определенного внимания, так как потоки будут работать независимо друг от друга, и более ранний пакет данных может быть обработан после более позднего просто из-за недостатка времени работы потока.

      Примечание

      Для обработки без порядковых требований это не считается недостатком.

    • Ручное фиксирование позиции усложнится, поскольку будет требовать координации всех потоков для обеспечения успешного завершения обработки партиции.

Такой подход имеет множество вариаций. Например, каждый поток может иметь собственную очередь, а потоки-потребители могут хешироваться в эти очереди, используя метод TopicPartition для обеспечения порядкового потребления и облегчения процесса фиксирования.