Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 04.03 · 25 мин
Средний
Consumer Groupsgroup.idPartition AssignmentGroup Coordinator

Consumer Groups

Consumer groups — это механизм горизонтального масштабирования потребления данных из Kafka. Идея проста: несколько consumer объединяются в группу с общим group.id и совместно читают партиции топика — каждая партиция назначается ровно одному consumer в группе. Добавляете consumer — нагрузка распределяется. Убираете consumer — его партиции переходят к остальным.


group.id: логическое имя группы

group.id — это строковой идентификатор, объединяющий consumer в одну группу. Все consumer с одним group.id кооперируют друг с другом через group coordinator — специальный брокер, ответственный за управление группой.

Роль group.id

Consumer 1 (group_id=‘analytics’)

Consumer 1: group_id='analytics'. Один из трёх экземпляров группы. Получит назначение на некоторые партиции топика 'events'.

Consumer 2 (group_id=‘analytics’)

Consumer 2: тот же group_id='analytics'. Является частью той же группы. Breker назначит ему непересекающийся набор партиций.

Consumer 3 (group_id=‘analytics’)

Consumer 3: тот же group_id='analytics'. Третий участник группы. Все три consumer вместе покрывают весь топик.
координация

Group Coordinator (broker)

Group Coordinator: брокер, ответственный за управление этой consumer group. Определяется по формуле: hash(group.id) % num_partitions(__consumer_offsets). Хранит committed offsets группы в __consumer_offsets.
назначение партиций

P0 → C1

Partition 0 → Consumer 1. Только Consumer 1 читает эту партицию в данной группе.

P1 → C1

Partition 1 → Consumer 1. Consumer 1 получил две партиции (топик 6 партиций, 3 consumer → по 2 каждому).

P2 → C2

Partition 2 → Consumer 2. Consumer 2 читает свои две партиции независимо от C1 и C3.

P3 → C2

Partition 3 → Consumer 2. Назначение равномерное: 6 партиций / 3 consumer = 2 на каждого.

P4 → C3

Partition 4 → Consumer 3. Consumer 3 обрабатывает свои два потока данных независимо.

P5 → C3

Partition 5 → Consumer 3. Полное покрытие топика: все 6 партиций читаются параллельно тремя consumer.

Назначение партиций

Назначение партиций происходит через алгоритм partition assignment strategy. Kafka поддерживает несколько стратегий:

  • RangeAssignor (по умолчанию в legacy protocol): последовательные партиции назначаются consumer по диапазонам
  • RoundRobinAssignor: партиции распределяются по очереди, более равномерно при нескольких топиках
  • StickyAssignor: минимизирует перемещение партиций при ребалансировке
  • CooperativeStickyAssignor: cooperative incremental ребалансировка с минимальным прерыванием
  • Server-side assignment (KIP-848): брокер сам управляет назначением — consumer только сообщают о желаемых топиках
NOTE

В KIP-848 (Kafka 4.0 с group.protocol=consumer) назначение партиций полностью перешло на сторону сервера. Consumer больше не участвует в синхронизации — брокер управляет assignment асинхронно через ConsumerGroupHeartbeat API. Подробнее в уроке 04 этого модуля.


ConsumerGroupDiagram: визуализация группы

Consumer Group: назначение партиций
P0Партиция P0: упорядоченный, неизменяемый лог записей. Гарантирует порядок только внутри партиции — между партициями порядок не гарантирован. Ключ записи определяет, в какую партицию она попадает (hash(key) % numPartitions). Назначена консьюмеру C0.
P1Партиция P1: вторая партиция топика. Каждая партиция может обрабатываться только одним консьюмером в группе одновременно — это обеспечивает ровно-однократную обработку записей в рамках одного consumer group. Назначена консьюмеру C0.
P2Партиция P2: третья партиция. Если количество консьюмеров больше числа партиций, лишние консьюмеры простаивают. Если меньше — консьюмер получает несколько партиций. Назначена консьюмеру C1.
P3Партиция P3: четвёртая партиция топика. Смещение (offset) читается и коммитится независимо для каждой партиции. Консьюмер коммитит offset после обработки записей. При перебалансировке новый консьюмер продолжит с последнего закоммиченного offset. Назначена консьюмеру C1.
P0, P1
Group CoordinatorКоординатор группы (Group Coordinator): специальный брокер, ответственный за данную consumer group. Выбирается по формуле: abs(groupId.hashCode()) % __consumer_offsets.numPartitions. Управляет heartbeat-сессиями (session.timeout.ms), инициирует перебалансировку при подключении/отключении консьюмеров. В KIP-848 заменяется на серверное управление назначением (Server-Side Rebalance Protocol).
P2, P3
Consumer C0Консьюмер C0: экземпляр консьюмера в группе analytics-group. Обрабатывает партиции P0 и P1. Выполняет poll-цикл: consumer.poll(Duration.ofMillis(100)) возвращает записи из назначенных партиций. После обработки вызывает commitSync() или commitAsync() для сохранения offset в __consumer_offsets.
P0Партиция P0 назначена C0: консьюмер C0 читает записи из P0 начиная с закоммиченного offset. Скорость чтения ограничена max.poll.records (по умолчанию 500). Если обработка занимает больше max.poll.interval.ms (5 мин), координатор считает консьюмера мёртвым и запускает rebalance.
P1Партиция P1 назначена C0: обрабатывается последовательно с P0 в рамках poll-цикла. При cooperative rebalance (KIP-429) C0 может временно сохранить P1 пока координатор перераспределяет только изменившиеся партиции, минимизируя stop-the-world паузы.
Consumer C1Консьюмер C1: второй экземпляр в группе. Обрабатывает партиции P2 и P3. Отправляет heartbeat координатору каждые heartbeat.interval.ms (3с). Если координатор не получает heartbeat в течение session.timeout.ms (45с), инициирует rebalance. C1 может быть отдельным процессом или потоком.
P2Партиция P2 назначена C1: C1 читает записи из P2. Стратегия назначения партиций определяется partition.assignment.strategy: RangeAssignor, RoundRobinAssignor, StickyAssignor или CooperativeStickyAssignor (рекомендуется для минимизации rebalance-пауз).
P3Партиция P3 назначена C1: последовательно обрабатывается C1. Consumer lag = LEO(P3) - committed_offset(C1, P3) — показывает, насколько консьюмер отстаёт от продюсера. Мониторируется через kafka-consumer-groups.sh --describe или JMX-метрику records-lag-max.
commit offsets
__consumer_offsets (committed offsets)Топик __consumer_offsets: внутренний топик Kafka с 50 партициями (по умолчанию), где хранятся закоммиченные смещения всех consumer groups. Ключ записи = (group.id, topic, partition), значение = committed_offset. Позволяет консьюмеру возобновить чтение с корректного места после перезапуска или rebalance.

Диаграмма выше показывает типичное распределение для группы из 2 consumer и 4 партиций: Consumer 0 получает P0 и P1, Consumer 1 получает P2 и P3. Group Coordinator хранит offsets в __consumer_offsets.


Group Coordinator: выбор и роль

Group Coordinator — это брокер, ответственный за жизненный цикл конкретной consumer group. Он определяется детерминированно:

coordinator_partition = hash(group.id) % num_partitions(__consumer_offsets)
coordinator_broker = leader(partition coordinator_partition of __consumer_offsets)

По умолчанию __consumer_offsets имеет 50 партиций. Для group.id='analytics': если hash('analytics') % 50 = 12, то coordinator — это брокер, который является лидером партиции 12 топика __consumer_offsets.

Обязанности Group Coordinator:

Group Coordinator — функции

Хранение committed offsets

Хранение offsets: coordinator хранит текущие закоммиченные offsets всех consumer group в топике __consumer_offsets. При рестарте consumer получает последний закоммиченный offset именно от coordinator.

Управление составом группы

Управление membership: coordinator отслеживает, какие consumer активны в группе. При потере heartbeat от consumer — инициирует ребалансировку.

Координация ребалансировки

Координация ребалансировки: при изменении состава группы coordinator организует протокол JoinGroup/SyncGroup (legacy) или управляет сервер-driven assignment (KIP-848).

Обработка heartbeat

Heartbeat обработка: consumer регулярно отправляют HeartbeatRequest (legacy) или ConsumerGroupHeartbeat (KIP-848). Coordinator следит за session timeout для каждого члена.

Приём commit offset запросов

Commit offsets: координирует получение OffsetCommitRequest от consumer и запись в __consumer_offsets. Синхронный commit блокирует consumer до получения ack.

Масштабирование: consumer vs партиции

Максимальный параллелизм consumer group ограничен числом партиций топика. Это фундаментальное ограничение архитектуры.

Consumer vs Партиции: правило масштабирования

4 партиции = 4 consumer (оптимально)

Оптимальный случай: 4 партиции, 4 consumer. Каждый consumer получает ровно 1 партицию. Максимальный параллелизм при полном использовании ресурсов.

4 партиции, 2 consumer — нормально

Недостаток consumer: 4 партиции, 2 consumer. Каждый consumer получает 2 партиции. Throughput ниже, но корректно. Добавление 2 consumer восстановит параллелизм.

4 партиции, 6 consumer — 2 idle (пустые)

Избыток consumer: 4 партиции, 6 consumer. Только 4 consumer получат партиции. 2 consumer будут простаивать — партиций на всех не хватает. Это не ошибка, но 2 слота простаивают зря.
# Правило: добавляем consumer до num_partitions — дальше бессмысленно
# Увеличить параллелизм можно только увеличив число партиций топика

from kafka import KafkaConsumer

# Consumer 1 — в группе 'payments-group'
consumer1 = KafkaConsumer(
    bootstrap_servers=['localhost:9092'],
    group_id='payments-group',
    auto_offset_reset='earliest',
)
consumer1.subscribe(['payments'])

# Consumer 2 — в той же группе 'payments-group'
consumer2 = KafkaConsumer(
    bootstrap_servers=['localhost:9092'],
    group_id='payments-group',
    auto_offset_reset='earliest',
)
consumer2.subscribe(['payments'])
# consumer1 и consumer2 поделят партиции топика 'payments' между собой

Несколько consumer groups: независимое чтение

Различные consumer groups читают один топик полностью независимо. Каждая группа ведёт свои offsets. Это ключевое отличие Kafka от традиционных message broker.

Несколько consumer groups — независимое чтение
Topic: orders (3 партиции)Топик 'orders' содержит все заказы. Каждая из трёх consumer groups читает его независимо, ведя свои offsets.

billing-group (offset=5000)

group_id='billing': читает все заказы для выставления счетов. Offset этой группы — 5000. Другие группы не влияют на этот offset.

analytics-group (offset=3000)

group_id='analytics': читает те же заказы для аналитики. Может быть на offset=3000 — отстаёт, но это не проблема.

notifications-group (offset=4500)

group_id='notifications': читает заказы для отправки уведомлений. Offset=4500. Три независимых потока обработки одних данных.
# Три разные группы — три независимых потребителя одного топика
billing_consumer = KafkaConsumer(
    bootstrap_servers=['localhost:9092'],
    group_id='billing-group',   # отдельный offset
)
analytics_consumer = KafkaConsumer(
    bootstrap_servers=['localhost:9092'],
    group_id='analytics-group', # свой независимый offset
)
# Обе группы читают топик 'orders' — полностью независимо
billing_consumer.subscribe(['orders'])
analytics_consumer.subscribe(['orders'])
TIP

Это именно то, что отличает Kafka от RabbitMQ: в RabbitMQ сообщение удаляется после первого получателя. В Kafka три consumer group получают одинаковые данные без дублирования в топике — каждая ведёт свой offset. Новая group, созданная с auto_offset_reset='earliest', прочитает всю историю с начала, не мешая другим группам.


Мониторинг consumer group lag

Consumer lag — это разница между последним записанным offset (Log End Offset, LEO) и последним закоммиченным offset группы. Лаг показывает, насколько группа отстаёт от продюсеров.

# CLI: мониторинг лага consumer group (Kafka 4.0)
kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --describe \
  --group analytics-group

# Пример вывода:
# GROUP           TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
# analytics-group orders  0          5000            5200            200
# analytics-group orders  1          4800            5100            300
# analytics-group orders  2          5100            5100            0

Суммарный лаг 200 + 300 + 0 = 500 означает, что группа не обработала 500 сообщений. Если лаг растёт — producer пишет быстрее, чем consumer обрабатывает. Решение: добавить consumer (до числа партиций) или увеличить число партиций.


Итог

Consumer groups — фундаментальный механизм масштабирования потребления в Kafka. group.id объединяет consumer. Group coordinator управляет назначением партиций и offsets. Максимальный параллелизм = число партиций. Несколько groups читают независимо. Лаг — ключевая метрика здоровья группы.

Проверка знанийKnowledge check
Топик 'events' имеет 4 партиции. В consumer group 'analytics' работают 6 consumer. Что произойдёт?
ОтветAnswer
Только 4 consumer получат назначение по одной партиции каждый. Оставшиеся 2 consumer будут простаивать — они присоединены к группе, но не получили ни одной партиции, поскольку партиций меньше, чем consumer. Это не ошибка, но 2 consumer расходуют ресурсы впустую. Для использования всех 6 consumer нужно увеличить число партиций топика до 6 (или кратного 6).

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 7