Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 04.06 · 30 мин
Продвинутый
KIP-932Share GroupsKafkaShareConsumerqueue semanticsAcknowledgeTypeshare coordinator

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 продвигается одной сущностью. Но как только нагрузка на конкретное сообщение становится тяжёлой и неравномерной — модель ломается.

Проблемы consumer groups для queue-like нагрузок

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 на сбоях» — это лишние оковы.

NOTE

До 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 сообщений.

Жизненный цикл записи в share group

Available

Available: запись лежит в логе, готова к доставке. Все записи между SPSO (share-partition start offset) и SPEO (end offset) — потенциально доступны. Записи после SPEO ещё не вошли в окно.
poll()

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’ы не персистятся — только эпохи, чтобы детектировать «зомби» координаторов.

INFO

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Транзиентная ошибка → обратно в Availabletimeout, retryable downstream failure
REJECTНевосстановимая ошибка → Archivedpoison message, schema corruption

В Kafka 4.2 (KIP-1222) добавлен четвёртый — RENEW, продление acquisition lock без изменения состояния записи. Полезно для долгих обработок, где 30-секундного дефолта не хватает, но и поднимать share.record.lock.duration.ms глобально не хочется.

WARNING

Каждый 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 vs share group

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март 2025early access (только для тестирования)
4.1сентябрь 2025preview (включается флагом, не для prod)
4.2февраль 2026GA, production-ready
WARNING

Кластеры с включёнными 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 на каждую запись.
TIP

В implicit mode код проще: убираешь share.acknowledgement.mode=explicit, и каждый следующий poll() автоматически акает все записи предыдущего как ACCEPT. Если бросил исключение — записи освободятся по lock timeout. Подходит для простых сценариев, где ты ок терять детали ошибки.


Когда выбирать share groups

Дерево решений: share group или consumer group

Нужен порядок?

Нужен ли строгий порядок обработки записей по ключу или 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 — мощный инструмент, но у него есть цена.

Ограничения 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, а не сокращать.
WARNING

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’ов. Сравнение:

Что нужноRabbitMQKafka до 4.2Kafka 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)
Orderingper-queue только при 1 consumerper-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
Проверка знанийKnowledge check
Топик `payments-events` имеет 4 партиции, по нему читает share group из 16 consumer'ов. Один из consumer'ов вызывает release() на запись из-за timeout downstream-сервиса. Что произойдёт дальше с этой записью?
ОтветAnswer
Запись вернётся из состояния Acquired в состояние Available, broker инкрементирует её delivery_count на единицу. Запись снова станет доступной для доставки — её может получить ЛЮБОЙ из 16 consumer'ов share group (включая того же, который её release'нул, но с большой вероятностью — другой, так как координатор балансирует распределение). Если delivery_count в итоге достигнет group.share.delivery.attempt.limit (по умолчанию 5), broker автоматически переведёт запись в состояние Archived и больше не будет её доставлять — это встроенная защита от poison message, не требующая ручной DLQ-логики.
Проверка знанийKnowledge check
Команда хочет мигрировать обработку webhook-доставки с RabbitMQ на Kafka share groups. У них 8 партиций в топике webhook-deliveries и они планируют запустить 200 worker'ов. Какие три ключевых момента нужно учесть при таком переходе?
ОтветAnswer
Во-первых, share group даёт M:N распределение — 200 worker'ов смогут читать с 8 партиций одновременно (несколько consumer'ов на partition), параллелизм больше не привязан к числу партиций. Во-вторых, нет ordering guarantees: если для конкретного эндпойнта нужен порядок доставки webhook'ов — share group не подходит, придётся либо использовать consumer group, либо обеспечить идемпотентность на стороне получателя. В-третьих, default lock duration 30 секунд задаёт минимальный delay повторной доставки при падении worker'а с in-flight записью; если webhook-вызов обычно занимает 2-5 секунд, при сбое worker'а очередная попытка пойдёт только через 30 сек — нужно либо настроить share.record.lock.duration.ms ниже, либо использовать explicit acknowledgement с быстрым RELEASE при ошибках.
Проверка знанийKnowledge check
В чём принципиальное отличие хранения delivery state в share groups по сравнению с consumer groups, и как это связано с топиком __share_group_state?
ОтветAnswer
В consumer groups state в основном клиентский: клиенты сами отслеживают свой position, broker хранит только commit offset в __consumer_offsets — одно число на (group, topic, partition). В share groups broker отслеживает состояние КАЖДОЙ in-flight записи (Available, Acquired, Acknowledged, Archived), её delivery_count и acquisition lock — это per-record state. Этот state персистится в новый внутренний топик __share_group_state в виде записей ShareSnapshot (полный снапшот состояния share-partition) и ShareUpdate (инкрементальные изменения). Share Coordinator — broker, который ведёт партицию __share_group_state для конкретной share-partition — при перезапуске проигрывает свою партицию и восстанавливает in-memory state. Назначения consumer'ов на partition'ы при этом НЕ персистятся, только эпохи — это оптимизация, чтобы уменьшить объём записываемых данных.

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

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

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

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