Главная › Блог › Kafka: consumer group и ребалансировка партиций

Kafka: consumer group и ребалансировка партиций

Зачем нужна consumer group

Partition — это единица параллелизма в Kafka. Одну партицию в один момент времени читает не более одного консьюмера из группы. Consumer group — это механизм распределения партиций топика между несколькими инстансами приложения так, чтобы каждое сообщение из партиции обработал ровно один консьюмер группы.

Если партиций 6 и консьюмеров 3 — каждому достанется по 2 партиции. Если консьюмеров 7 — один останется без работы, Kafka не дробит партицию между несколькими потребителями. Это ключевое ограничение: горизонтальное масштабирование consumer group ограничено числом партиций.

Партиции — это кассы в супермаркете, консьюмеры — кассиры. Каждую кассу обслуживает один кассир. Если касс 6, а кассиров 10 — четверо будут курить в подсобке. Ребалансировка — это момент, когда менеджер временно закрывает все кассы и пересаживает кассиров по новой схеме.

Что такое ребалансировка

Ребалансировка (rebalance) — это процесс пересчёта того, какой консьюмер читает какие партиции. Она запускается, когда:

  • новый консьюмер присоединился к группе;
  • консьюмер вышел из группы (упал, завершился, не прислал heartbeat вовремя);
  • изменилось число партиций в топике;
  • истёк session.timeout.ms или consumer не успел обработать batch за max.poll.interval.ms.

За координацию отвечает group coordinator — один из брокеров, назначенный для конкретной группы. Консьюмеры шлют ему heartbeat через отдельный поток (начиная с KIP-62), а присвоение партиций делает один из консьюмеров — лидер группы, выбранный coordinator'ом, используя стратегию партиционирования (RangeAssignor, RoundRobinAssignor, StickyAssignor, CooperativeStickyAssignor).

Stop-the-world: классическая проблема eager rebalance

До KIP-429 ребалансировка была eager: на время пересчёта все консьюмеры группы отзывают все свои партиции (revoke), ждут нового assignment и только потом начинают читать заново. Даже если партиция всё равно остаётся у того же консьюмера — она на секунды выпадает из обработки.

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

Частая ошибка — долгая обработка сообщения внутри poll() цикла без увеличения max.poll.interval.ms. Consumer считается живым пока присылает heartbeat, но если между вызовами poll() проходит больше таймаута — coordinator считает его мёртвым и выкидывает из группы, начинается ребалансировка.

CooperativeStickyAssignor — решение

Начиная с Kafka 2.4 появился incremental cooperative rebalancing. Идея: консьюмер отзывает только те партиции, которые реально должны перейти к другому, а остальные продолжает читать без паузы. Ребалансировка может занять несколько раундов, но throughput не падает до нуля.

props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
    CooperativeStickyAssignor.class.getName());

Для прода это почти всегда правильный выбор вместо дефолтного набора ассайнеров — особенно если группа большая и часто масштабируется.

Коммит offset и ребалансировка — тонкое место

Перед тем как партиция будет отозвана, нужно успеть закоммитить offset обработанных сообщений — иначе после ребалансировки новый владелец партиции начнёт читать с последнего закоммиченного offset и часть сообщений обработается повторно. Для этого есть listener:

consumer.subscribe(topics, new ConsumerRebalanceListener() {
    public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
        consumer.commitSync(currentOffsets);
    }
    public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
        // можно явно seek() на нужный offset
    }
});

Без явного коммита в onPartitionsRevoked приходится рассчитывать только на auto-commit по таймеру, который может не успеть сработать до отзыва партиции — отсюда дубликаты при at-least-once.

Можно ли, чтобы одно сообщение из партиции обработали два консьюмера одной группы параллельно?
Нет, в рамках одной consumer group партицию читает только один консьюмер. Параллельная обработка одной партиции двумя потребителями возможна только если они в разных group.id — тогда это уже два независимых потока чтения, не связанных балансировкой нагрузки.
Почему увеличение числа консьюмеров сверх числа партиций не ускоряет обработку?
Партиция неделима между консьюмерами группы. Лишние консьюмеры просто остаются без назначенных партиций (idle) и ждут, пока освободится партиция при следующей ребалансировке.
Что произойдёт, если consumer завис в обработке сообщения дольше max.poll.interval.ms?
Coordinator посчитает его выбывшим из группы, запустит ребалансировку и переназначит его партиции другим консьюмерам, хотя сам процесс физически жив и продолжает heartbeat-поток.

Ещё разборы