Share groups: Kafka как очередь

KIP-932, GA в Kafka 4.2. Когда consumer group не подходит и что вы отдаёте взамен эластичности.

Consumer group упирается в число партиций: партицию читает ровно один участник, и обработчиков не может быть больше, чем партиций. Для потоковой обработки это правильно — так держится гарантия порядка. Для задач типа «разобрать очередь заданий» это ограничение мешает.

Share groups (KIP-932) решают именно это. В share group партиции ни за кем не закреплены: любой участник берёт любую доступную запись. Участников может быть сколько угодно, независимо от числа партиций. Функциональность стала общедоступной в Apache Kafka 4.2.

Механика отличается от consumer group. Взятая запись блокируется на время обработки (acquisition lock, по умолчанию 30 секунд): если участник не успел её подтвердить, лок истекает и запись возвращается в общий доступ. У каждой записи есть счётчик попыток доставки; по умолчанию после пяти неудач она уходит в archived. Подтверждение персональное: acknowledge, release или reject для каждой записи отдельно.

Плата за эластичность — порядок. Обработка занимает разное время, и запись, взятая позже, вполне может завершиться раньше. Никакой гарантии последовательности внутри партиции для share group нет.

Практический критерий выбора: нужен порядок по ключу — consumer group. Нужна эластичность, конкурентная обработка и подтверждение по каждому сообщению — share group. Обе модели работают на одном и том же топике одновременно и не мешают друг другу.

Разобрать это в симуляторе Интерактивный сценарий: настраиваешь кластер, ломаешь его и смотришь, что происходит

Частые вопросы

Заменяют ли share groups RabbitMQ?

Закрывают значительную часть сценариев очереди, оставаясь в экосистеме Kafka. Но привычных для брокеров очередей возможностей — приоритетов, TTL на сообщение, сложной маршрутизации — в них нет.

Можно ли читать один топик и consumer group, и share group?

Да, это независимые механизмы, работающие на одном топике параллельно.

Что происходит при истечении лока?

Запись возвращается в доступные и достаётся другому участнику, а счётчик попыток доставки увеличивается.

Дальше по теме

Что такое партиция в Kafka Consumer lag: что это и как с ним работать acks в Kafka: 0, 1 и all ISR и min.insync.replicas Ребалансировка consumer group