Share Groups (KIP-932): queue semantics в Kafka
Десятилетие Kafka жила с одним правилом: одна партиция — один consumer внутри группы. Это правило давало упорядоченность, простую модель offset и линейное масштабирование, но закрывало целый класс use-case’ов — task queues, work distribution, медленный per-message processing. Команды, которым нужны были очереди, ставили рядом RabbitMQ или SQS. KIP-932 «Queues for Kafka» закрывает этот пробел: вводит новый тип группы — share group — где несколько consumer’ов кооперативно читают из одной партиции, а acknowledgment работает на уровне отдельной записи. Preview появился в Kafka 4.0 (март 2025), GA — в Kafka 4.2 (февраль 2026).
Зачем это нужно: ограничения consumer groups
Классическая модель consumer groups устроена так: координатор группы назначает каждую партицию ровно одному consumer’у. Это даёт упорядоченное чтение в пределах партиции и гарантирует, что один offset продвигается одной сущностью. Но как только нагрузка на конкретное сообщение становится тяжёлой и неравномерной — модель ломается.
Head-of-line blocking
Slow consumer head-of-line blocking: один consumer обрабатывает запись 10 секунд (ML inference, внешний API). Все следующие записи в его партиции ждут — даже если другие consumer'ы простаивают.Параллелизм ограничен числом партиций
Жёсткий потолок параллелизма: количество одновременно работающих consumer'ов не может превышать число партиций. Если партиций 8, а ты хочешь 50 worker'ов — никак. Нужен реpartitioning топика, что меняет порядок и всю downstream-логику.Дисбаланс «горячих» партиций
Дисбаланс по партициям: если producer использует hash-партиционирование по key, одна горячая partition получает 80% трафика. Другие consumer'ы простаивают, тот, кому досталась горячая partition, захлёбывается.Нет per-record acknowledgment
Нет per-record retry: при exception consumer должен либо коммитить и терять сообщение, либо не коммитить и блокировать весь partition. Custom DLQ требует ручной реализации со всеми её багами.Для классического event streaming эти ограничения — фича, а не баг: они дают порядок и предсказуемость. Но для task queue, где порядок не важен, а важно «выгрести очередь как можно быстрее с retry на сбоях» — это лишние оковы.
До KIP-932 распространённый workaround выглядел так: создать топик с большим количеством партиций (256, 512), а потом надеяться, что хеш ключа ляжет равномерно. Это работает плохо: partition skew, перебалансировки на каждом scale-out, неудобство при изменении масштаба.
Архитектура: broker-side delivery state
Главное архитектурное отличие share groups от consumer groups — где живёт состояние delivery. В consumer groups координация в основном клиентская: лидер группы вычисляет назначения, клиенты сами трекают offset, broker лишь хранит коммиты. В share groups broker берёт на себя гораздо больше: он отслеживает delivery state каждой отдельной записи в скользящем окне in-flight сообщений.
Available
Available: запись лежит в логе, готова к доставке. Все записи между SPSO (share-partition start offset) и SPEO (end offset) — потенциально доступны. Записи после SPEO ещё не вошли в окно.Acquired (locked, 30s)
Acquired: запись доставлена консьюмеру и заблокирована для других на время acquisition lock (по умолчанию share.record.lock.duration.ms=30000, 30 секунд). Другие consumer'ы того же share group не получат эту запись, пока lock жив.ACCEPT → Acknowledged
ACCEPT: успешная обработка. Запись переходит в Acknowledged, broker увеличивает share-partition start offset (SPSO), запись больше никогда не будет доставлена.RELEASE → Available (retry)
RELEASE: транзиентная ошибка (timeout, недоступность downstream). Запись возвращается в Available, increment delivery_count. Может быть доставлена другому consumer'у того же share group.REJECT → Archived
REJECT: невосстановимая ошибка (poison message, schema corruption). Запись сразу архивируется, не пытаемся доставлять снова. Аналог переноса в DLQ.Lock истёк → Available + delivery_count++
Lock expiration: consumer не успел вызвать acknowledge/release/reject за 30 секунд (умер, завис, GC pause). Broker автоматически возвращает запись в Available с increment delivery_count. Защита от потерянных in-flight записей.delivery_count ≥ limit → Archived
Достигнут лимит delivery_attempts (по умолчанию group.share.delivery.attempt.limit=5). Запись архивируется автоматически — защита от poison message, который никто не может обработать.Скользящее окно in-flight записей живёт между двумя offset’ами:
- SPSO (Share-Partition Start Offset) — всё до этого offset уже Acknowledged или Archived;
- SPEO (Share-Partition End Offset) — верхняя граница окна; записи выше ещё не вошли в делёжку.
Окно ограничено по размеру (по умолчанию 200 записей на share-partition, регулируется group.share.record.lock.partition.limit), чтобы broker не хранил неограниченный delivery state в памяти.
Share Coordinator и __share_group_state
Состояние share group персистится в новом внутреннем топике — __share_group_state. Записи в нём ключуются по (group, topicId, partition). Share Coordinator — broker, который ведёт партицию состояния для конкретной share-partition; он же обслуживает RPC от share consumer’ов.
Coordinator пишет два типа записей:
- ShareSnapshot — полный самодостаточный снапшот состояния share-partition (SPSO, in-flight записи, их состояния и delivery counts);
- ShareUpdate — инкрементальное изменение: «такая-то запись перешла в Acknowledged», «вот это — в Archived».
При перезапуске broker’а coordinator проигрывает партицию __share_group_state от последнего snapshot и восстанавливает in-memory state. Назначения consumer’ов на partition’ы не персистятся — только эпохи, чтобы детектировать «зомби» координаторов.
Share Coordinator — это отдельная сущность от Group Coordinator (который остался обслуживать consumer groups). В Kafka 4.2 broker может одновременно быть Group Coordinator для одних групп и Share Coordinator для других — share groups и consumer groups сосуществуют на одном кластере и могут читать одни и те же топики.
Acknowledgment-модель: ACCEPT, RELEASE, REJECT, RENEW
KafkaShareConsumer работает в двух режимах: implicit (записи автоматически акают на следующем poll()) и explicit (приложение само вызывает acknowledge() для каждой записи). Explicit режим даёт три типа подтверждения:
| AcknowledgeType | Семантика | Применение |
|---|---|---|
ACCEPT | Успешно обработано → Acknowledged | штатный happy path |
RELEASE | Транзиентная ошибка → обратно в Available | timeout, retryable downstream failure |
REJECT | Невосстановимая ошибка → Archived | poison message, schema corruption |
В Kafka 4.2 (KIP-1222) добавлен четвёртый — RENEW, продление acquisition lock без изменения состояния записи. Полезно для долгих обработок, где 30-секундного дефолта не хватает, но и поднимать share.record.lock.duration.ms глобально не хочется.
Каждый RELEASE инкрементирует delivery_count на стороне broker’а. После достижения group.share.delivery.attempt.limit (по умолчанию 5) broker сам архивирует запись — даже если consumer продолжает звать release(). Это встроенная защита от poison message: бесконечный retry-loop невозможен в принципе.
Partition assignment: несколько consumer’ов на одну partition
В consumer groups действует жёсткое правило: одна partition — один consumer. В share groups оно отменено. Координатор может назначить одну и ту же partition нескольким consumer’ам того же share group — они будут разбирать её записи кооперативно, через acquisition lock.
Consumer group: P0 → C0, P1 → C1, P2 → C2, P3 → C3 (1:1)
Consumer group: 4 партиции, 4 consumer'а — каждый получает по одной. Если consumer'ов 8 — половина простаивает. Параллелизм жёстко ограничен числом партиций.Share group: P0 → C0,C1; P1 → C2,C3; P2 → C4,C5; P3 → C6,C7 (M:N)
Share group: 4 партиции, 8 consumer'ов. Broker распределяет partition'ы между ВСЕМИ — несколько consumer'ов могут читать одну partition. Acquisition lock на стороне broker гарантирует, что одна запись не будет обработана дважды одновременно.Scale-out: новый consumer мгновенно подключается к пулу
Scale-out без rebalancing'а в классическом смысле: новый consumer присоединяется и сразу начинает забирать записи из общего пула. Никакой stop-the-world паузы — координатор просто перераспределяет назначения. Это отличается от KIP-848 cooperative rebalancing — там всё ещё есть пересдача partition'ов.Это меняет правила игры для масштабирования: количество worker’ов больше не привязано к количеству партиций. У тебя топик с 8 партициями и 100 worker’ов — все 100 будут активны, разбирая записи через broker-managed lock.
Версии и конфигурация
| Версия Kafka | Дата | Статус share groups |
|---|---|---|
| 4.0 | март 2025 | early access (только для тестирования) |
| 4.1 | сентябрь 2025 | preview (включается флагом, не для prod) |
| 4.2 | февраль 2026 | GA, production-ready |
Кластеры с включёнными share groups в 4.0 нельзя обновить до 4.1 и выше из-за изменений формата __share_group_state. Конфиг 4.0 был объявлен как «не upgrade-safe» — его явно использовали только для evaluation.
Минимальная конфигурация broker’а (Kafka 4.2):
# server.properties
group.coordinator.rebalance.protocols=classic,consumer,share
share.coordinator.state.topic.replication.factor=3
share.coordinator.state.topic.min.isr=2
Группа создаётся как share при подписке через KafkaShareConsumer — никакого ручного создания не нужно, broker сам распознаёт тип по протоколу запроса.
Ключевые tuning-параметры:
share.record.lock.duration.ms— длительность acquisition lock (по умолчанию 30000, диапазон 1000–60000);group.share.record.lock.partition.limit— лимит in-flight записей на share-partition (по умолчанию 200);group.share.delivery.attempt.limit— максимум попыток до архивирования (по умолчанию 5);share.auto.offset.reset— стартовая позиция для новой share group (latest/earliest).
Code example: KafkaShareConsumer
import org.apache.kafka.clients.consumer.AcknowledgeType;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaShareConsumer;
import java.time.Duration;
import java.util.List;
import java.util.Properties;
public class TaskQueueWorker {
public static void main(String[] args) {
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "image-processing-workers");
props.put("key.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
props.put("value.deserializer",
"org.apache.kafka.common.serialization.StringDeserializer");
// Explicit acknowledgment: per-record control
props.put("share.acknowledgement.mode", "explicit");
try (KafkaShareConsumer<String, String> consumer =
new KafkaShareConsumer<>(props)) {
consumer.subscribe(List.of("image-processing-tasks"));
while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> record : records) {
try {
ImageJob job = parseJob(record.value());
processImage(job); // долгая ML обработка
// Успех: подтверждаем конкретную запись
consumer.acknowledge(record, AcknowledgeType.ACCEPT);
} catch (TransientDownstreamException e) {
// Retry: запись вернётся в Available и достанется кому-то ещё
consumer.acknowledge(record, AcknowledgeType.RELEASE);
} catch (PoisonMessageException e) {
// Невосстановимо: сразу в Archived
consumer.acknowledge(record, AcknowledgeType.REJECT);
}
}
// Коммит acknowledgments на broker одним батчем
consumer.commitSync();
}
}
}
}
Несколько ключевых отличий от обычного KafkaConsumer:
- Нет вызова
assign()илиseek()— share consumer не управляет offset’ами; он только подтверждает записи. acknowledge(record, type)принимает конкретную запись, а не batch — это per-record acknowledgment.commitSync()отправляет все накопленные подтверждения на broker одной операцией, чтобы не делать round-trip на каждую запись.
В implicit mode код проще: убираешь share.acknowledgement.mode=explicit, и каждый следующий poll() автоматически акает все записи предыдущего как ACCEPT. Если бросил исключение — записи освободятся по lock timeout. Подходит для простых сценариев, где ты ок терять детали ошибки.
Когда выбирать share groups
Нужен порядок?
Нужен ли строгий порядок обработки записей по ключу или per-partition? Если да — только consumer group. Share group НЕ даёт ordering: внутри одной partition записи разбираются параллельно несколькими consumer'ами.Consumer group
Consumer group: stream processing, event sourcing, CDC, аудит, любая ситуация, где порядок событий критичен (например, обновления балансов счёта).Параллелизм > partition count?
Параллелизм нужно скейлить выше, чем число partition'ов в топике? Тогда consumer group не подходит — потолок параллелизма = число партиций.Share group
Share group: task queue, work distribution, ML inference, image/video processing — задачи независимы, важна пропускная способность, а не порядок.Нужен per-record retry?
Нужен per-record retry с DLQ? В consumer group приходится строить custom DLQ-логику с ручным offset-управлением. Share group даёт это из коробки через REJECT/lock-timeout/delivery_count.Share group
Share group: каждый record можно ACCEPT/RELEASE/REJECT независимо. Poison message protection через delivery_attempt_limit встроена в broker.Типичные use case’ы под share groups:
- Task queues — обработка независимых job’ов (image resize, PDF generation, ML inference);
- Notification fan-out — рассылка уведомлений с retry на флэйки backend;
- Webhook delivery — outbound HTTP с retry и DLQ через REJECT;
- Background processing — всё, что раньше тащили в RabbitMQ или SQS «потому что Kafka так не умеет».
Limitations
Share groups — мощный инструмент, но у него есть цена.
Нет ordering
Нет ordering guarantees. Внутри одной partition несколько consumer'ов разбирают записи параллельно — порядок обработки не определён. Все use case'ы, где порядок важен (event sourcing, CDC, баланс счёта), под share group НЕ подходят.Offset не общий с consumer groups
Share groups существуют параллельно с consumer groups: они не делят offset'ы. Один и тот же топик может читаться share group и consumer group одновременно — позиции независимы. Это не миграция: переключение нагрузки требует копирования логики.SPSO влияет на retention-семантику
Усложнение retention. Share-partition start offset продвигается ТОЛЬКО когда все записи до этой границы Acknowledged или Archived. Если один consumer завис на середине окна — SPSO не двигается, и retention по этим offset'ам не отрабатывает (хотя время-based retention работает как обычно).Нагрузка на broker (state на каждой записи)
Дополнительная нагрузка на broker. Share Coordinator поддерживает in-memory state in-flight записей, пишет ShareSnapshot/ShareUpdate в __share_group_state. Для топиков с большим количеством share-partition'ов и высоким QPS это заметно.Нет транзакций и Streams
Транзакции и Kafka Streams в share groups не поддерживаются (по крайней мере в 4.2 GA). Если нужен exactly-once или DSL Streams — это всё ещё consumer group territory.Latency повторной доставки = lock duration
Lock duration (30 сек по умолчанию) задаёт нижнюю границу latency повторной доставки при падении consumer'а. Если ваша задача обрабатывается 50ms, а consumer умер с записью в Acquired — следующая попытка только через 30 секунд. RENEW (KIP-1222) помогает только продлевать lock, а не сокращать.Share groups не заменяют consumer groups, а дополняют их. Stream processing, event sourcing, CDC, любые ordering-sensitive нагрузки остаются за consumer groups. Share groups — про задачи, где порядок неважен, а важна elastic пропускная способность с per-record retry.
RabbitMQ-like patterns в Kafka
С share groups Kafka впервые становится разумной заменой RabbitMQ для типичных queue use case’ов. Сравнение:
| Что нужно | RabbitMQ | Kafka до 4.2 | Kafka 4.2+ share groups |
|---|---|---|---|
| Per-message ack/nack | да (basic.ack/nack) | нет (только offset commit) | да (ACCEPT/RELEASE/REJECT) |
| Параллелизм > потребителей-на-очередь | да | нет (1 partition = 1 consumer) | да (M:N) |
| Poison message protection | через DLQ exchange | ручная реализация | встроенная (delivery_attempt_limit) |
| Retention сообщений после ack | нет (удаляются) | да (offset commit, лог остаётся) | да (запись архивируется логически, но физически в логе) |
| Replay сообщений | сложно | штатно | штатно (новая share group с earliest) |
| Ordering | per-queue только при 1 consumer | per-partition | нет |
Для команд, где Kafka уже стоит как event backbone, share groups убирают главный аргумент в пользу RabbitMQ — «нам нужны очереди». Не нужно тащить второй брокер с отдельной operational нагрузкой.
Итог
Share groups (KIP-932) — крупнейшее изменение consumer-модели Kafka за десятилетие. Broker берёт на себя отслеживание delivery state каждой записи в скользящем окне; партиции назначаются нескольким consumer’ам кооперативно; per-record acknowledgment даёт ACCEPT/RELEASE/REJECT с встроенной poison message protection. Это закрывает целый класс queue use-case’ов — task queues, work distribution, retry с DLQ — для которых раньше приходилось ставить рядом RabbitMQ. Цена — отказ от ordering, дополнительная нагрузка на broker и независимость offset-пространства от обычных consumer groups. В Kafka 4.2 (февраль 2026) фича объявлена production-ready.
Spark transformWithState: альтернативный stateful pattern