Kafka: consumer group и ребалансировка партиций
Зачем нужна consumer group
Partition — это единица параллелизма в Kafka. Одну партицию в один момент времени читает не более одного консьюмера из группы. Consumer group — это механизм распределения партиций топика между несколькими инстансами приложения так, чтобы каждое сообщение из партиции обработал ровно один консьюмер группы.
Если партиций 6 и консьюмеров 3 — каждому достанется по 2 партиции. Если консьюмеров 7 — один останется без работы, Kafka не дробит партицию между несколькими потребителями. Это ключевое ограничение: горизонтальное масштабирование consumer group ограничено числом партиций.
Что такое ребалансировка
Ребалансировка (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.