📨 Раздел 15 · Вопрос #29

Можно ли читать сообщения из определённой партиции

В Apache Kafka консьюмер может читать сообщения напрямую из конкретно указанных партиций топика. Для этого вместо механизма динамической подписки consumer.subscribe() использует...


🟢 Junior Level

Прямой ответ: Да, можно!

В Apache Kafka консьюмер может читать сообщения напрямую из конкретно указанных партиций топика. Для этого вместо механизма динамической подписки consumer.subscribe() используется метод явного ручного назначения — consumer.assign().

Динамическая подписка (subscribe):
Consumer ──subscribe("orders")──> Group Coordinator на брокере сам решает,
                                  какие партиции выделить (например, партиции 0 и 1).

Ручное назначение (assign):
Consumer ──assign([orders-0])───> Читает СТРОГО и ТОЛЬКО партицию 0.
                                  Брокер не участвует в распределении.

Базовый пример на Java

Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);

// 1. Формируем список целевых партиций (например, партиция 0 топика "orders")
TopicPartition partition0 = new TopicPartition("orders", 0);
consumer.assign(List.of(partition0));

// 2. Читаем сообщения в стандартном цикле poll
while (running) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        System.out.printf("Partition: %d, Offset: %d, Key: %s, Value: %s%n",
            record.partition(), record.offset(), record.key(), record.value());
    }
}

Назначение в Spring Kafka

Во фреймворке Spring Kafka прямое чтение конкретных партиций настраивается декларативно в аннотации @KafkaListener:

@KafkaListener(
    topicPartitions = @org.springframework.kafka.annotation.TopicPartition(
        topic = "orders", 
        partitions = { "0", "2" } // Слушаем только 0 и 2 партиции
    )
)
public void listenSpecificPartitions(ConsumerRecord<String, String> record) {
    log.info("Получено сообщение из партиции {}: {}", record.partition(), record.value());
}

🟡 Middle Level

Фундаментальные отличия: subscribe() против assign()

Методы subscribe() и assign() являются взаимоисключающими (mutually exclusive). Попытка вызвать оба метода на одном экземпляре KafkaConsumer приведет к выбросу исключения IllegalStateException.

Характеристика consumer.subscribe() consumer.assign()
Координатор группы Активно используется (GroupCoordinator) Полностью игнорируется
Ребалансировка (Rebalance) Автоматическая при падении или масштабировании Отсутствует. Консьюмер автономен (Standalone)
Отказоустойчивость Высокая: партиция упавшего пода отойдет соседу Нулевая: при падении партиция перестанет читаться
Наличие group.id Обязательно Опционально (нужно только для хранения смещений)
Добавление новых партиций Обнаруживаются автоматически Игнорируются (нужен ручной перевызов assign())
Область применения Стандартная продакшн-обработка микросервисов Реплей, отладка, аудит, StatefulSet шардирование

Управление смещением (Seek API)

Главная сила assign() раскрывается в сочетании с произвольным позиционированием чтения:

TopicPartition tp = new TopicPartition("orders", 0);
consumer.assign(List.of(tp));

// 1. Чтение с самого начала партиции (с наименьшего доступного LogStartOffset)
consumer.seekToBeginning(List.of(tp));

// 2. Чтение только новых сообщений (с хвоста лога LogEndOffset)
consumer.seekToEnd(List.of(tp));

// 3. Чтение с точного номера смещения
consumer.seek(tp, 450200L);

// 4. Поиск смещения по временной метке (Time-based Seek)
long oneHourAgo = System.currentTimeMillis() - TimeUnit.HOURS.toMillis(1);
Map<TopicPartition, Long> timestampsToSearch = Map.of(tp, oneHourAgo);

// Запрашиваем у брокера первое смещение, записанное после указанного времени
Map<TopicPartition, OffsetAndTimestamp> foundOffsets = consumer.offsetsForTimes(timestampsToSearch);
OffsetAndTimestamp target = foundOffsets.get(tp);

if (target != null) {
    consumer.seek(tp, target.offset());
}

🔴 Senior Level

Работа с group.id и __consumer_offsets при ручном назначении

Распространенное заблуждение: «При использовании assign() параметр group.id не нужен и коммитить офсеты нельзя».

Это не так. Механика работы зависит от конфигурации:

  1. group.id не указан:
    • Вызовы consumer.commitSync() или consumer.commitAsync() выбросят InvalidGroupIdException: To use the group management or offset commit APIs, you must provide a valid group.id in the consumer configuration.
    • Консьюмер не сохраняет состояние на брокере. При каждом старте он будет читать с позиции, определяемой auto.offset.reset (если не был вызван seek()).
  2. group.id указан:
    • Консьюмер НЕ вступает в протокол балансировки группы (не шлет JoinGroup / SyncGroup).
    • Однако при вызове commitSync() координатор брокера сохраняет смещения в системном топике __consumer_offsets под указанным group.id!
    • При перезапуске приложения consumer.assign() автоматически прочитает последний закоммиченный офсет из __consumer_offsets и продолжит обработку без дублей.

Production-архитектура: Паттерн StatefulSet + Partition Pinning

В некоторых сценариях динамическая ребалансировка subscribe() является неприемлемой (Stop-the-World паузы при ребалансе, сброс локальных кэшей). В распределенных базах данных и аналитических движках (Apache Flink, Kafka Streams, локальные кэши на RocksDB) используется паттерн Deterministic Partition Pinning:

Kubernetes StatefulSet (order-worker):
Pod 0 (HOSTNAME="order-worker-0") ──> assign(TopicPartition("orders", 0)) ──> Локальный RocksDB 0
Pod 1 (HOSTNAME="order-worker-1") ──> assign(TopicPartition("orders", 1)) ──> Локальный RocksDB 1
Pod 2 (HOSTNAME="order-worker-2") ──> assign(TopicPartition("orders", 2)) ──> Локальный RocksDB 2

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

  1. Полное отсутствие пауз на ребалансировку при перезапуске подов (Deployment).
  2. Локальное дисковое состояние (Embedded RocksDB) жестко привязано к партиции и не перекачивается по сети между подами при рестартах.
  3. Простота расследования инцидентов: логирование партиции изолировано в конкретном контейнере.

Недостатки:

Если Pod 1 падает по OOM, Kubernetes обязан перезапустить именно его. Пока под поднимается, данные партиции 1 не обрабатываются никем, так как другие поды не могут автоматически подхватить чужую партицию.


4 Tricky Questions

1. Что произойдет, если в работающем приложении сначала вызвать consumer.subscribe(List.of("topicA")), а затем на том же объекте вызвать consumer.assign(List.of(new TopicPartition("topicB", 0)))?

Ответ: Вызов завершится исключением IllegalStateException: java.lang.IllegalStateException: Subscription to topics, partitions and pattern are mutually exclusive. Внутри KafkaConsumer поле SubscriptionState хранит тип текущего назначения: NONE, AUTO_TOPICS, AUTO_PATTERN или USER_ASSIGNED. Переключение между динамической подпиской и ручным назначением возможно только после явного вызова consumer.unsubscribe(), который очищает текущее состояние и сбрасывает соединение с координатором группы.

2. Администратор увеличил количество партиций в топике с 4 до 8 на работающем кластере. Как отреагируют консьюмеры с subscribe() и консьюмеры с assign()?

Ответ:

  1. Консьюмеры с subscribe(): Фоновый поток метаданных регулярно обновляет информацию о кластере (metadata.max.age.ms, по умолчанию 5 минут). Обнаружив новые партиции, консьюмеры инициируют ребалансировку группы, и новые партиции 4..7 будут автоматически распределены между участниками группы.
  2. Консьюмеры с assign(): Они не увидят новых партиций никогда. Список назначенных партиций зафиксирован в памяти клиента статически. Приложение продолжит читать только исходные партиции 0..3. Чтобы начать читать партиции 4..7, приложение обязано либо перезапуститься, либо периодически программно опрашивать consumer.partitionsFor("topic") и вручную вызывать consumer.assign(updatedPartitions).

3. Будет ли работать ConsumerRebalanceListener при использовании consumer.assign()?

Ответ: Нет, никогда. Интерфейс ConsumerRebalanceListener с методами onPartitionsRevoked() и onPartitionsAssigned() передается только в качестве аргумента метода subscribe(). Метод assign() принимает исключительно коллекцию TopicPartition и не имеет перегрузок для листенеров. Поскольку протокол координации группы не задействован, событий отзыва или назначения партиций со стороны брокера не генерируется.

4. Можно ли организовать Exactly-Once обработку с записью в реляционную БД, используя assign() вместо subscribe()?

Ответ: Да, и это один из самых надежных классических паттернов.

  1. Консьюмер назначает партицию через assign(tp).
  2. Приложение читает сообщения и записывает бизнес-данные в PostgreSQL/Oracle.
  3. В ту же самую транзакцию базы данных (в служебную таблицу kafka_offsets) сохраняется пара (partition_id, processed_offset):
    BEGIN TRANSACTION;
    INSERT INTO orders (id, amount) VALUES (1, 100);
    UPDATE kafka_offsets SET offset = 105 WHERE partition = 0;
    COMMIT;
    
  4. При старте или после сбоя сервис выполняет: long lastOffset = db.query("SELECT offset FROM kafka_offsets WHERE partition = 0"); consumer.seek(tp, lastOffset + 1);
  5. Поскольку смещение коммитится атомарно вместе с бизнес-данными в единой транзакции БД, дубликаты и потери данных полностью исключены (строгая семантика Exactly-Once без необходимости использовать тяжелый Kafka 2PC Transactions API).

🎯 Шпаргалка для интервью

30-секундный ответ (Elevator Pitch)

«Да, читать конкретную партицию можно с помощью метода consumer.assign(Collection<TopicPartition>). В отличие от subscribe(), ручное назначение полностью отключает протокол координации группы и ребалансировку — консьюмер работает как полностью автономный процесс (Standalone Consumer). Это дает абсолютный детерминированный контроль над партициями и позволяет использовать произвольное позиционирование через Seek API (seek(), seekToBeginning(), offsetsForTimes()). Плата за это — отсутствие встроенной отказоустойчивости: если консьюмер с assign() падает, его партиции автоматически никому не передаются. Основные применения: реплей данных, отладка, аудит, а также шардированные архитектуры с локальным стейтом (StatefulSets + RocksDB)».

Ключевые метрики для Production-мониторинга

  • kafka.consumer:type=consumer-fetch-manager-metrics,client-id={id} -> records-lag: лаг конкретной назначенной партиции.
  • Статус процесса подов в Kubernetes: так как автоматического перехвата партиций нет, падение пода требует немедленного алерта.

Красные флаги (Чего говорить НЕЛЬЗЯ)

  • ❌ «subscribe() и assign() можно вызывать одновременно на одном консьюмере» — Это приведет к немедленному IllegalStateException.
  • ❌ «При assign() нельзя коммитить офсеты в Kafka» — Можно, если указан group.id; смещения будут сохраняться в __consumer_offsets.
  • ❌ «Если под с assign() упадет, другой под подхватит его партицию» — Нет, при ручном назначении координатор группы отключен, партиция повиснет до рестарта пода.
  • ❌ «Новые партиции, созданные админом на лету, автоматически появятся в assign()» — Нет, список статичен в памяти, требуется явный перевызов assign().

Связанные темы