Consumer group — механизм, которым Kafka масштабирует чтение: партиции распределяются между участниками группы, и каждую читает ровно один. Звучит просто, пока не начинается ребалансировка — перераспределение партиций при входе/выходе участника, которое при неудачной настройке превращается в «rebalance storm», когда группа больше перетасовывает партиции, чем обрабатывает сообщения. А коммит offset определяет, что происходит при сбое: потеря или дубль.
Это вторая статья серии «Kafka: глубокое погружение». Опирается на модель из первой статьи про лог, топики и партиции: партиция — единица параллелизма чтения, а consumer group — механизм, который эту единицу раздаёт участникам без пересечений. Все числа и логи ниже — из живого трёхброкерного KRaft-кластера, не из документации. Код обоих стендов открыт в репозитории digital-cookbook: kafka/go/consumer-groups (franz-go) и kafka/java/consumer-groups (kafka-clients).
В статье
- Группа и координатор
- Стратегии назначения партиций
- Eager против incremental cooperative: живой лог ребаланса
- Static membership: пережить перезапуск без ребаланса
- Rebalance storms
- Коммит offset: auto vs manual
Группа и координатор
Consumer group — это набор consumer’ов с общим group.id, между которыми Kafka распределяет партиции топика так, чтобы каждую партицию в группе читал ровно один consumer одновременно. За распределение отвечает не сами consumer’ы, а выделенный брокер — group coordinator (один из брокеров кластера, определяется хешем group.id). Coordinator ведёт членство группы (кто сейчас в ней состоит), принимает heartbeat’ы участников и запускает ребаланс, когда состав группы меняется: кто-то присоединился, вышел по таймауту, упал или явно покинул группу.
Жизненный цикл участника держится на двух таймаутах. session.timeout.ms — сколько coordinator ждёт heartbeat, прежде чем счесть участника мёртвым и исключить его из группы (в живом стенде выставлен в 12 секунд для сценария static membership — подробности ниже). max.poll.interval.ms — отдельный таймер: сколько времени разрешено между двумя вызовами poll(); если consumer завис в обработке дольше этого интервала, coordinator считает его выбывшим независимо от heartbeat’а, потому что современный протокол шлёт heartbeat из отдельного фонового потока, а не синхронно с poll(). Смешивать эти два таймаута — частый источник путаницы: session.timeout про «клиент перестал отвечать вообще», max.poll.interval — про «клиент завис в бизнес-логике».
Партиций в группе не может быть меньше, чем нужно на всех активных consumer’ов с пользой: если consumer’ов больше, чем партиций, лишние получат пустое назначение и будут простаивать — масштабирование чтения одного топика упирается в число его партиций (см. первую статью про то, как это число выбирается при создании топика).
Стратегии назначения партиций
Как именно партиции раскладываются между участниками — решает стратегия назначения (assignor / balancer), и это отдельная от протокола ребаланса ось: стратегия отвечает на вопрос «кому какая партиция», протокол — на вопрос «как именно это назначение доехало до участников» (см. следующий раздел). У обоих клиентов живого стенда — franz-go и kafka-clients — доступны те же четыре стратегии:
- range — партиции каждого топика делятся на смежные диапазоны между consumer’ами по алфавитному порядку
client.id/member.id; при нескольких топиках один и тот же consumer систематически получает младшие диапазоны везде — известная слабость range при совместном потреблении нескольких топиков. - round-robin — все партиции всех подписанных топиков раскладываются по кругу между consumer’ами; ровнее range, но так же перетасовывает всё назначение заново при каждом изменении состава.
- sticky — тот же результат, что round-robin по балансу, но минимизирует движение партиций между последовательными назначениями: старается сохранить прежнее закрепление там, где это не мешает равномерности.
- cooperative-sticky — тот же принцип «минимум движения», что и sticky, но реализован через другой протокол ребаланса (incremental cooperative вместо eager) — из-за этого различие sticky и cooperative-sticky не в том, что назначается, а в том, как назначение доставляется участникам, и это ключевая тема следующего раздела.
Дефолты у клиентов расходятся, и это стоит держать в голове при переносе конфигурации между стеками: franz-go (pkg/kgo, версия v1.21.5 в стенде) по умолчанию использует только CooperativeStickyBalancer — других вариантов «из коробки» нет, их нужно явно перечислить через kgo.Balancers(...). У kafka-clients дефолт классического (classic) протокола — список [RangeAssignor, CooperativeStickyAssignor] в актуальных версиях, но фактическое поведение зависит от согласования с coordinator; в живом стенде именно поведение classic-протокола kafka-clients по умолчанию наблюдалось как eager (подробности — в разделе про static membership ниже).
specs := []strategySpec{
{"range", kgo.RangeBalancer()},
{"roundrobin", kgo.RoundRobinBalancer()},
{"sticky", kgo.StickyBalancer()},
{"cooperative-sticky", kgo.CooperativeStickyBalancer()},
}
// ...
opt := kgo.Balancers(spec.balancer)
cl, err := kgo.NewClient(kgo.SeedBrokers(seeds...), kgo.ConsumerGroup(groupID), opt, ...)private static final List<Spec> SPECS = List.of(
new Spec("range", RangeAssignor.class.getName()),
new Spec("roundrobin", RoundRobinAssignor.class.getName()),
new Spec("sticky", StickyAssignor.class.getName()),
new Spec("cooperative-sticky", CooperativeStickyAssignor.class.getName())
);
// ...
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, List.of(spec.assignorClass()));Eager против incremental cooperative: живой лог ребаланса
Протокол ребаланса — это то, как назначение партиций доезжает до участников группы, и здесь разница между range/roundrobin/sticky с одной стороны и cooperative-sticky с другой — не декларативная, а измеримая. На живом стенде (топик demo-groups, 6 партиций, RF=3) три consumer’а стабилизируют назначение, затем подключается четвёртый — и логируется, что именно каждая стратегия отзывает у уже работающих участников.
При eager-протоколе (range, roundrobin, sticky) вступление нового участника — это stop-the-world для всех остальных: coordinator сначала отзывает у КАЖДОГО прежнего участника ВСЁ его текущее назначение, и только потом раздаёт партиции заново с нуля. Реальный лог стенда:
--- range (IsCooperative=false) ---
[3.213s] consumer-a revoked при вступлении consumer-d: [0 1] — ВСЁ (stop-the-world)
[3.213s] consumer-b revoked: [2 3] — ВСЁ
[3.213s] consumer-c revoked: [4 5] — ВСЁroundrobin и sticky в eager-режиме ведут себя идентично по механике — у всех троих прежних участников полный revoke; различается только то, что именно будет назначено заново, а не то, как отзывается старое. При cooperative-sticky картина другая — ни один прежний участник не теряет ВСЁ своё назначение целиком, ребаланс сходится за несколько раундов, и отзывается только та партиция, что реально переезжает к новому участнику:
--- cooperative-sticky (IsCooperative=true) ---
[12.118s] consumer-a/b/c REVOKED [] (пусто — первый раунд ничего не забрал)
[12.121s] consumer-a revoked при вступлении consumer-d: [] — ПУСТО
[12.121s] consumer-a REVOKED [4] текущее назначение: [0] <- второй раунд: только реально переезжающая партиция
[12.626s] consumer-d ASSIGNED [4] <- досталась новому члену только теперьПрограммные ассерты стенда закрепляют это как проверяемый факт, а не наблюдение по логу: для eager-стратегий у каждого прежнего участника обязан найтись хотя бы один revoke с полным совпадением размера с назначением ДО события; для cooperative-sticky — ни одного такого «полного» revoke. Оба ассерта зелёные и на franz-go, и на kafka-clients. Стоит отдельно отметить тонкость классификации: «полный revoke» — это не то же самое, что «непустой revoke». В одном из прогонов у участника с назначением [0 4] встретился REVOKED [4] — партиция ушла, но [0] осталась у того же владельца; это инкрементальный, а не stop-the-world revoke, хотя список отозванных партиций и не пуст. Различать эти случаи нужно сравнением размера отозванного множества с размером назначения непосредственно перед событием, а не проверкой «пусто/не пусто».
sequenceDiagram
participant C1 as consumer-1
participant C2 as consumer-2
participant GC as group coordinator
participant C3 as consumer-3 (новый)
rect rgb(238, 232, 216)
Note over C1,C3: EAGER (range / roundrobin / sticky) — stop-the-world
C3->>GC: JoinGroup
GC-->>C1: REVOKED [0 1 2] — ВСЁ текущее назначение
GC-->>C2: REVOKED [3 4 5] — ВСЁ текущее назначение
Note over C1,C2: чтение партиций встало у обоих, пока идёт SyncGroup
GC-->>C1: ASSIGNED [0 1]
GC-->>C2: ASSIGNED [2 3]
GC-->>C3: ASSIGNED [4 5]
end
rect rgb(201, 228, 197)
Note over C1,C3: COOPERATIVE-STICKY — инкрементально
C3->>GC: JoinGroup
GC-->>C1: REVOKED [] — партиции не тронуты
GC-->>C2: REVOKED [] — партиции не тронуты
Note over C1,C2: чтение продолжается без остановки
GC-->>C3: ASSIGNED [] — первый раунд ничего не отдал
Note over GC: второй раунд: только реально переезжающая партиция
GC-->>C1: REVOKED [4]
GC-->>C3: ASSIGNED [4]
end
Одна оговорка про таймкоды в логах выше: секундные метки — общие на весь сценарий, который последовательно прогоняет все четыре стратегии в одном процессе (range идёт первой, cooperative-sticky — последней), поэтому 3.213s у range и 12.118s у cooperative-sticky нельзя сравнивать напрямую как «время ребаланса» — во второй входит и выполнение трёх предыдущих сценариев. Что метки действительно показывают про cooperative — это разрыв внутри него самого: 12.118s (первый раунд, пустой revoke) → 12.626s (второй раунд, реальный переезд партиции), то есть ребаланс сходится не за один round-trip JoinGroup/SyncGroup, а за несколько.
Практический вывод: cooperative-sticky не «быстрее» по абсолютному времени сходимости — за счёт нескольких раундов протокола он вполне может занять больше wall-clock, — но он не останавливает чтение у участников, которых изменение состава не касается напрямую. При большой группе и частой ротации участников это и есть разница между «на секунду просела вся группа» и «просела только одна партиция».
Отдельно стоит развести это с автоматическим переносом партиций между брокерами при administrative reassignment (kafka-reassign-partitions) или авто-ребалансировкой лидеров (auto.leader.rebalance.enable) — это два разных механизма на разных уровнях. Ребаланс consumer group, описанный здесь, — про то, какой consumer какую партицию читает. Reassignment — про то, на каком брокере физически лежит реплика партиции; consumer group при этом вообще не меняет состав. Подробно об эксплуатационной стороне reassignment — в статье про эксплуатацию, KRaft и тюнинг.
Static membership: пережить перезапуск без ребаланса
Обычный (динамический) участник группы при graceful-остановке шлёт LeaveGroupRequest, и coordinator немедленно запускает ребаланс среди оставшихся — даже если процесс через секунду перезапустится и вернётся в группу. Для деплоя с частыми рестартами (rolling update, автоскейлинг, просто нестабильная сеть) это означает лишний ребаланс на каждый такой рестарт — ощутимая цена, если consumer’ов и партиций много.
Static membership (KIP-345) решает это явным идентификатором участника — group.instance.id в kafka-clients, kgo.InstanceID(...) в franz-go. Пока новый процесс переподключается с ТЕМ ЖЕ instance-id в пределах session.timeout.ms, coordinator трактует это не как выход и вход нового участника, а как временное отсутствие того же самого: партиции остаются закреплены за instance-id, и при graceful-остановке LeaveGroupRequest не отправляется вовсе — сравните с обычным закрытием, которое шлёт этот запрос всегда.
На живом стенде (топик demo-groups, session.timeout.ms=12000) static-a и static-b стабилизируют назначение, затем static-a останавливается и перезапускается с тем же group.instance.id. Рестарт занял 166–512 мс (тайминги host-зависимы, но с большим запасом укладываются на порядки меньше session.timeout в 12 секунд) — партиции вернулись той же статике, а static-b за всё время не получил ни одного revoke или assign: группа не ребалансировалась вовсе.
a := newMember("static-a", groupID, nil,
kgo.InstanceID("static-a"), kgo.SessionTimeout(sessionTimeout))
// ...
a.close() // graceful close статического члена НЕ шлёт LeaveGroupRequest
a2 := newMember("static-a-restarted", groupID, nil,
kgo.InstanceID("static-a"), kgo.SessionTimeout(sessionTimeout))
// в пределах session.timeout партиции возвращаются той же статикеprops.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, SESSION_TIMEOUT_MS);
props.put(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, instanceId); // "static-a"
// graceful close() с заданным group.instance.id — LeaveGroupRequest
// НЕ отправляется (KIP-345), coordinator ждёт session.timeout перед тем,
// как считать участника выбывшим окончательноКонтраст с динамическим членством в том же стенде показателен вдвойне: без group.instance.id перезапуск dynamic-c вызывает немедленный revoke/assign у dynamic-d на обоих клиентах — но по-разному в деталях. У franz-go (дефолт — CooperativeStickyBalancer) это REVOKED []/ASSIGNED [] — событие засчитано, но инкрементально пусто. У kafka-clients (дефолт классического протокола в этом сценарии де-факто ведёт себя как eager) — полный REVOKED [3 4 5], назначение пересобирается с нуля. Сам факт «ребаланс произошёл» одинаков, но его цена — нет: это ещё одна иллюстрация разницы eager/cooperative из предыдущего раздела, только на другом триггере (не «новый участник», а «временное исчезновение без static membership»).
Rebalance storms
Rebalance storm — ситуация, когда группа проводит в перетасовке партиций больше времени, чем в их фактическом чтении: один ребаланс триггерит условия для следующего, и группа не успевает стабилизироваться. Три типичных источника, которые усиливают друг друга:
max.poll.interval.msменьше реального времени обработки батча. Если бизнес-логика внутри цикла обработки (не heartbeat, а именно синхронная работа между вызовамиpoll()) регулярно превышает интервал, coordinator периодически считает участника зависшим, исключает его, тот переподключается — и особенно больно это бьёт при eager-протоколе, где каждое такое исключение — stop-the-world для всей группы, а не только для одного участника.- Частые рестарты без static membership. Каждый перезапуск процесса без
group.instance.id— этоLeaveGroupRequestпри остановке иJoinGroupпри старте, то есть минимум один гарантированный ребаланс на цикл рестарта; при автоскейлинге или нестабильной сети это может повторяться десятками раз за короткое окно. - Слишком короткий
session.timeout.msотносительно реальной сетевой задержки или пауз GC. Участник, у которого heartbeat изредка опаздывает из-за сборки мусора или сетевого джиттера, периодически выпадает из группы и возвращается — с тем же эффектом, что и первый пункт, но по другой причине.
Смягчение — по тем же трём осям: держать max.poll.interval.ms с запасом над худшим реалистичным временем обработки батча (а тяжёлую работу — выносить из цикла poll в отдельный воркер-пул, как это, кстати, и делает демонстрация auto-commit ниже — с обратной стороны это же решение создаёт риск потери данных, если не увязать его с ручным коммитом); использовать static membership там, где рестарты предсказуемо частые; не занижать session.timeout.ms ради «быстрее обнаружить падение» без запаса на реальные паузы среды. cooperative-sticky не устраняет причины шторма, но ограничивает его цену: даже при частых пересборках назначения простаивают только реально переезжающие партиции, а не вся группа разом.
Коммит offset: auto vs manual
Позиция чтения (offset), которую consumer группы закоммитил, — это то, откуда начнёт читать группа при следующем подключении: сбой процесса ДО коммита откатывает чтение назад (заново обработать — риск дубля), сбой ПОСЛЕ коммита, но до завершения обработки, продвигает чтение вперёд без реального завершения работы (риск потери). Выбор стратегии коммита — это выбор между at-least-once (готовы к дублям, не готовы к потерям) и at-most-once (готовы к потерям, не готовы к дублям); Kafka сама по себе не даёт exactly-once на этом уровне — она даёт инструменты для одной из двух гарантий, а какую выбрать, решает конфигурация коммита и код обработчика.
Auto-commit коммитит offset по таймеру (auto.commit.interval.ms, по умолчанию 5000 мс) независимо от того, завершилась ли реальная обработка полученных записей. На живом стенде это воспроизведено намеренно укороченными таймингами (интервал коммита 100 мс против «обработки» в фоновом воркере на 400 мс — не прод-дефолты, а демонстрация гонки, реальный дефолт Kafka — 5000 мс) — consumer раздаёт полученные записи в фоновые воркеры и сразу идёт за следующей пачкой, не дожидаясь их завершения. Позиция чтения (а с ней auto-commit) уходит вперёд быстрее, чем реально завершается обработка:
[auto] отправлено 20 сообщений
[auto] 'краш': раздиспатчено=20, реально доделано ДО краша=0
[auto] итог: отправлено=20, реально доделано до краша=0, дочитано новым консьюмером=0, суммарно=0
[assert] OK (auto-commit): потеря продемонстрирована — отправлено=20, суммарно доделано+дочитано=0 (< 20) => at-most-onceНа Go все 20 из 20 сообщений потеряны безвозвратно: новый consumer группы продолжает с уже ушедшей вперёд закоммиченной позиции и просто не видит то, что было раздиспатчено, но не доделано. На Java в зафиксированном прогоне картина та же (0 из 20 доделано до краша), хотя в отдельных прогонах JVM-планировщик иногда успевал доделать до трёх воркеров раньше «краша» — конкретное число плавает от прогона к прогону (это гонка потоков, а не разное поведение клиентов), но сам факт потерь — at-most-once — подтверждён на обоих.
Ручной коммит (commitSync-аналог: CommitRecords в franz-go, commitSync в kafka-clients) убирает эту гонку, если коммитить строго ПОСЛЕ завершения обработки, а не до. На стенде «упали» на середине пачки: обработано синхронно 12 из 20 записей, закоммичены только предыдущие батчи (3 из них), последняя пачка из 9 обработанных — не закоммичена, потому что «крах» случился до вызова коммита:
[manual] отправлено 20 сообщений
[manual] 'краш' после 12 обработанных записей (закоммичено раньше, батчами, без последней незакоммиченной пачки)
[assert] OK (manual-commit): все 20 сообщений покрыты (run1=12, run2=17, дублей=9) — потерь нет, at-least-once подтверждёнНовый consumer передоставил незакоммиченный хвост: 9 дублей (уже обработанных в первом прогоне, но не закоммиченных до краха) плюс 8 новых записей — итого 17 в повторном прогоне, суммарное покрытие 20 из 20 без потерь. Идентификация записей велась по значению (manual-msg-N), а не по Kafka-offset — топик demo-groups общий для всех сценариев стенда, и у шести его партиций независимые, несовпадающие последовательности offset’ов, так что offset сам по себе не идентифицирует конкретное сообщение однозначно между прогонами.
process := func(m *member, fetches kgo.Fetches) {
var batchRecs []*kgo.Record
fetches.EachRecord(func(r *kgo.Record) {
// ... обработка записи ...
batchRecs = append(batchRecs, r)
})
if len(batchRecs) == 0 {
return
}
// commitSync-аналог: коммитим ПОСЛЕ обработки всей пачки, не до
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
err := m.cl.CommitRecords(ctx, batchRecs...)
cancel()
}
c1 := newMember("manual-1", groupID, process, kgo.DisableAutoCommit())Member.Processor process = (m, records) -> {
List<ConsumerRecord<String, String>> batchRecs = new ArrayList<>();
for (ConsumerRecord<String, String> r : records) {
// ... обработка записи ...
batchRecs.add(r);
}
if (batchRecs.isEmpty()) return;
// commitSync-аналог: коммитим ПОСЛЕ обработки всей пачки, не до
m.rawConsumer().commitSync(batchOffsets(batchRecs));
};
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");Есть и цена у ручного коммита — латентность round-trip’а до coordinator’а. На стенде 10 последовательных commitSync-вызовов по одной записи заняли 11.7–13.4 мс (порядка 1.2–1.3 мс на коммит), а один батч-коммит на те же 10 записей — около 1.1 мс: на порядок быстрее, потому что это один round-trip до coordinator группы вместо десяти. Абсолютные миллисекунды здесь host-зависимы (сеть и диск брокера в конкретном прогоне), но сам факт — batching коммитов на порядок дешевле поштучных — воспроизводится структурно, а не случайно.
Итог по всем 19 ассертам стенда целиком (четыре сценария rebalance/strategies/static/commits вместе) — зелёный на обоих клиентах, 0 упавших проверок, штатное завершение процесса.
Дальше в серии: как партиция ведёт себя при отказе брокера и что означает acks на практике — в статье «Репликация и надёжность Kafka». Эксплуатационная сторона — метрики consumer lag, ручной и автоматический reassignment партиций между брокерами (не путать с ребалансом группы, описанным здесь) — в статье про KRaft, эксплуатацию и тюнинг. Если консьюмер живёт на JVM — сборка мусора, паузы и их влияние на heartbeat и session.timeout разобраны в «JVM и messaging: Kafka на практике». А для навигации по всей теме messaging сразу — карта messaging-landscape-map.
Комментарии