Consumer Groups
Consumer groups — это механизм горизонтального масштабирования потребления данных из Kafka. Идея проста: несколько consumer объединяются в группу с общим group.id и совместно читают партиции топика — каждая партиция назначается ровно одному consumer в группе. Добавляете consumer — нагрузка распределяется. Убираете consumer — его партиции переходят к остальным.
group.id: логическое имя группы
group.id — это строковой идентификатор, объединяющий consumer в одну группу. Все consumer с одним group.id кооперируют друг с другом через group coordinator — специальный брокер, ответственный за управление группой.
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 только сообщают о желаемых топиках
В KIP-848 (Kafka 4.0 с group.protocol=consumer) назначение партиций полностью перешло на сторону сервера. Consumer больше не участвует в синхронизации — брокер управляет assignment асинхронно через ConsumerGroupHeartbeat API. Подробнее в уроке 04 этого модуля.
ConsumerGroupDiagram: визуализация группы
Диаграмма выше показывает типичное распределение для группы из 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:
Хранение 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 ограничен числом партиций топика. Это фундаментальное ограничение архитектуры.
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.
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'])
Это именно то, что отличает 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 читают независимо. Лаг — ключевая метрика здоровья группы.