Skip to content
Learning Platform

Ментальная модель: партиция — это единица параллелизма

Прежде чем разбирать ребалансировку, зафиксируем одну идею, из которой вытекает всё остальное. В Apache Kafka единицей параллелизма потребления является не топик и не сообщение, а партиция. Топик разбит на N партиций, и внутри одной consumer group каждая партиция в любой момент времени назначена ровно одному консьюмеру. Это инвариант, который Kafka защищает любой ценой — именно ради него и существует весь механизм ребалансировки.

Из этого инварианта сразу следуют два практических вывода. Первый: больше консьюмеров в группе, чем партиций в топике, не дадут прироста — лишние консьюмеры будут простаивать без назначения. Второй: порядок сообщений Kafka гарантирует только внутри партиции, а раз партицию читает один консьюмер, то один консьюмер видит строго упорядоченный поток по каждой своей партиции. Масштабирование потребления — это игра в распределение партиций между членами группы, и rebalancing — это и есть тот самый акт распределения.

Эта статья — бесплатный самостоятельный разбор. Она пересекается с модулем про консьюмеров в нашем платном курсе Apache Kafka, но читается отдельно: всё объясняется здесь с нуля до уровня протокольных запросов.

Group coordinator: кто раздаёт партиции

За жизнь каждой группы отвечает один брокер — group coordinator. Какой именно брокер станет координатором, определяется детерминированно, без выборов и переговоров: берётся внутренний топик __consumer_offsets (по умолчанию 50 партиций), вычисляется hash(group.id) % num_partitions(__consumer_offsets), и лидер получившейся партиции и есть координатор группы. Поэтому group.id — это не просто метка: он жёстко привязывает группу к конкретному брокеру.

Координатор делает три вещи. Он ведёт список живых членов группы через heartbeat’ы. Он запускает ребалансировку, когда состав или метаданные меняются. И он хранит закоммиченные offset’ы группы — те самые «докуда дочитали» — в партиции __consumer_offsets, за которую отвечает. Важная деталь: в классическом протоколе координатор сам не считает, кому какая партиция достанется. Он лишь выбирает одного из консьюмеров group leader’ом и пересылает ему метаданные; вычисление назначения происходит на стороне клиента. Это объясняет, почему стратегия назначения — параметр консьюмера, а не брокера.

Что запускает ребалансировку

Координатор инициирует ребаланс при любом из событий: новый консьюмер шлёт JoinGroup; существующий уходит штатно (вызвал close(), отправив LeaveGroup) или признан мёртвым по таймауту; консьюмер сменил подписку через subscribe(); администратор добавил партиции в топик. Во всех случаях нарушается или может быть улучшен инвариант «одна партиция — один владелец», и группа пересобирает назначение.

Stop-the-world: почему классический ребаланс больно бьёт

До эпохи кооперативных протоколов работала так называемая eager (жадная) ребалансировка, и понять её необходимо, чтобы оценить, что именно чинили потом.

Жадный ребаланс — это синхронный барьер. Когда координатор объявляет ребаланс, каждый член группы сначала отзывает все свои партиции и прекращает обработку. Затем все шлют JoinGroup и ждут, пока координатор соберёт ответы от всех. Координатор назначает group leader’а, тот вычисляет распределение и рассылает его через SyncGroup. Только после этого консьюмеры возобновляют чтение. Ключевое слово — все. Добавление одного консьюмера в группу из ста останавливает всю сотню, и пауза тянется столько, сколько занимает самый медленный JoinGroup. В больших группах это десятки секунд полного простоя на каждое изменение состава. Отсюда и термин stop-the-world.

Корень боли в том, что жадный протокол на каждом ребалансе отбирает все партиции у всех, даже если 99% назначений по сути не меняются. Партиция P0, которую как читал консьюмер C1, так и продолжит читать после, всё равно проходит цикл «отозвали — переназначили» с холодным стартом.

Cooperative: ребаланс без полной остановки

Кооперативный инкрементальный протокол смягчает stop-the-world ровно за счёт устранения лишних отзывов. Идея: не отбирать всё подряд, а отбирать только те партиции, которые реально должны сменить владельца, и делать это в две фазы.

В первой фазе все консьюмеры снова шлют JoinGroup, но ничего не отзывают — продолжают читать свои текущие партиции. Лидер вычисляет целевое назначение и сравнивает его с текущим: партиции, которые остаются у прежних владельцев, не трогаются вовсе. Те, что должны переехать, помечаются на отзыв — и только их владельцы их освобождают. Это вызывает вторую короткую ребалансировку, в которой освобождённые партиции отдаются новым хозяевам. Консьюмер, чьи партиции никуда не переезжают, не прерывается ни на миллисекунду. cooperative rebalancing превращает глобальную остановку в локальную — паузу видят лишь те разделы, что меняют владельца.

Стоит отдельно сказать про следующее поколение — серверный протокол KIP-848 (group.protocol=consumer, по умолчанию в Kafka 4.0). В нём вычисление назначения переехало с клиента-лидера на брокер-координатор, а тяжёлый JoinGroup/SyncGroup заменён на лёгкий асинхронный ConsumerGroupHeartbeat. Ребаланс перестал быть синхронным барьером в принципе: координатор инкрементально сообщает каждому консьюмеру его новое назначение в ответах на heartbeat. Кооперативность здесь встроена в сам протокол.

Стратегии назначения партиций

Кто именно из членов получит какую партицию, определяет partition.assignment.strategy. Четыре стратегии, которые встречаются на практике:

СтратегияБалансировкаРебалансSticky (сохраняет назначения)Когда применять
RangeAssignorПо топикам: партиции одного индекса у одного консьюмераEager (stop-the-world)НетCo-partition join по нескольким топикам; дефолт-легаси
RoundRobinAssignorРавномерная по всем партициям всех топиковEager (stop-the-world)НетКогда важна ровная нагрузка, а co-partitioning не нужен
StickyAssignorРавномернаяEager, но старается сохранить прежние назначенияДаДорогой ребаланс, но клиент на старом протоколе
CooperativeStickyAssignorРавномернаяCooperative (без полной остановки)ДаДефолтный выбор для классического протокола сегодня

RangeAssignor раскладывает партиции потопично: для каждого топика партиции делятся диапазонами между консьюмерами. Побочный эффект — если топиков мало, а консьюмеров много, нагрузка перекашивается, зато партиции с одинаковым ключом в нескольких топиках попадают одному консьюмеру (полезно для join’ов). RoundRobinAssignor раздаёт все партиции всех топиков по кругу, давая ровную нагрузку, но не заботясь о ребалансе. sticky assignment добавляет ключевое свойство: при ребалансе сохранить максимум прежних назначений, чтобы не терять прогретые кэши и не переигрывать состояние. CooperativeStickyAssignor объединяет липкость с кооперативным протоколом и сегодня является разумным дефолтом для всех, кто ещё на классическом протоколе. Нельзя смешивать eager- и cooperative-стратегии в одной группе во время роллинг-апгрейда — это отдельная процедура в два деплоя.

Offset commit: где живёт «докуда дочитали»

Назначение партиций решает кто читает; offset решает откуда продолжить после рестарта или переезда партиции. Закоммиченный offset — это указатель, который консьюмер сохраняет координатору в __consumer_offsets. При ребалансе новый владелец партиции стартует именно с закоммиченного offset’а, поэтому стратегия коммита напрямую определяет гарантии доставки.

from kafka import KafkaConsumer

consumer = KafkaConsumer(
    "events",
    group_id="analytics",
    enable_auto_commit=False,          # ручной коммит
    partition_assignment_strategy=[
        # cooperative-sticky: без stop-the-world и с сохранением назначений
        "org.apache.kafka.clients.consumer.CooperativeStickyAssignor",
    ],
    max_poll_interval_ms=300000,       # бюджет на обработку одного poll-батча
    session_timeout_ms=10000,          # окно тишины heartbeat до признания мёртвым
)

for msg in consumer:
    process(msg)                       # сначала обрабатываем
    consumer.commit()                  # потом коммитим -> at-least-once

Развилка гарантий сводится к одному вопросу: коммитим до или после обработки.

  • At-most-once: коммит offset’а до обработки. Если консьюмер упал между коммитом и завершением работы, сообщение считается прочитанным, но не обработано. Потери возможны, дубликатов нет.
  • At-least-once: обработка, затем коммит (как в коде выше). Падение между обработкой и коммитом приведёт к повторной выдаче сообщения новому владельцу партиции. Дубликаты возможны, потерь нет. Это самый частый выбор.
  • Exactly-once: достижим не одним коммитом offset’а, а транзакцией. Запись результата и коммит входного offset’а упаковываются в одну Kafka-транзакцию (паттерн read-process-write через sendOffsetsToTransaction), либо обеспечивается идемпотентность приёмника. Без транзакций «exactly-once на словах» — это на деле at-least-once плюс дедупликация на стороне получателя.

auto-commit (enable.auto.commit=true) удобен, но обманчив: он коммитит offset’ы, которые poll() уже выдал приложению, независимо от того, обработаны ли они. Это даёт at-least-once в лучшем случае и тихие потери в худшем — если вы обрабатываете асинхронно, а авто-коммит успел сохранить offset до завершения обработки. Для контроля над гарантиями ручной коммит почти всегда правильнее.

session.timeout, max.poll.interval и эпидемия ребалансов

Два таймаута определяют, когда координатор сочтёт консьюмера мёртвым, и путать их — классический источник «группа ребалансируется каждые 30 секунд».

session.timeout.ms — это окно, в течение которого координатор должен услышать heartbeat. Heartbeat’ы в классическом протоколе шлёт фоновый поток консьюмера независимо от обработки, с интервалом heartbeat.interval.ms (правило: примерно треть от session.timeout.ms). Если за окно сессии ни одного heartbeat не пришло — координатор считает консьюмера упавшим и запускает ребаланс. Слишком маленький session.timeout.ms ловит ложные смерти при GC-паузах или сетевых микрозадержках; слишком большой замедляет обнаружение реальных падений.

max.poll.interval.ms про другое — это максимальный бюджет времени между двумя вызовами poll(). Heartbeat говорит «я жив на уровне сети», но если ваш цикл обработки одного батча затягивается дольше max.poll.interval.ms (тяжёлая трансформация, медленный внешний вызов на каждое сообщение), консьюмер не успевает вернуться к poll() вовремя. Координатор считает его зависшим, выкидывает из группы и запускает ребаланс — а консьюмер, закончив обработку, обнаруживает себя исключённым и пытается вернуться, провоцируя ещё один ребаланс. Это самопровоцирующийся цикл. Лечится либо уменьшением max.poll.records (меньше работы на батч), либо увеличением max.poll.interval.ms, либо выносом тяжёлой обработки из poll-потока. Различие принципиально: heartbeat подтверждает живость, а poll() подтверждает прогресс обработки — это разные оси.

ConsumerRebalanceListener: окно для аккуратного коммита

Между «партицию отозвали» и «партицию выдали новому владельцу» есть момент, который приложение обязано использовать правильно, иначе at-least-once превращается в потери. Этот момент даёт интерфейс ConsumerRebalanceListener с двумя коллбэками: onPartitionsRevoked (вызывается перед тем, как партиции уйдут) и onPartitionsAssigned (после того, как пришли новые). Правило простое: в onPartitionsRevoked нужно синхронно закоммитить offset’ы по отзываемым партициям, пока вы ещё их владелец. Если этого не сделать, новый владелец стартует с последнего закоммиченного offset’а и переиграет всё, что вы обработали, но не успели зафиксировать.

В кооперативном протоколе появляется третий коллбэк — onPartitionsLost. Он срабатывает, когда консьюмер уже потерял партиции «грязно» (например, вылетел из группы по таймауту), и коммитить в нём бессмысленно — партиции уже у других. Различение «отозвали штатно» и «потеряли внезапно» как раз и позволяет писать корректный код фиксации состояния при ребалансе вместо слепого коммита всего подряд.

Static membership: ребаланс, которого можно избежать

Часть ребалансов вообще не нужна. Роллинг-рестарт деплоя, под на Kubernetes, который пересоздаётся за 5 секунд, временный сетевой блип — каждый такой эпизод по умолчанию вызывает два ребаланса (на уход и на возврат), хотя по сути это тот же самый консьюмер с теми же партициями.

static membership решает это. Задав каждому экземпляру стабильный group.instance.id, вы говорите координатору: это не новый член, это вернувшийся старый. Если экземпляр перезапустился и вернулся в пределах session.timeout.ms, координатор отдаёт ему ровно те же партиции без какого-либо ребалансировки. На практике для управляемого роллинг-рестарта session.timeout.ms поднимают так, чтобы он покрывал окно перезапуска пода — тогда плановые деплои проходят вообще без перетасовки партиций. Цена — отложенное обнаружение настоящих падений (на величину увеличенного таймаута), и это осознанный компромисс.

Что забрать с собой

Партиция — единица параллелизма, и весь механизм групп служит инварианту «одна партиция — один владелец». Координатор детерминированно привязан к group.id, хранит offset’ы в __consumer_offsets и запускает ребалансы. Жадный протокол останавливает всю группу; кооперативный и серверный KIP-848 отзывают только переезжающие партиции. Стратегия назначения и стратегия коммита вместе задают баланс нагрузки и гарантии доставки. А session.timeout.ms, max.poll.interval.ms и static membership — это рычаги, которыми вы превращаете шумную, постоянно ребалансирующуюся группу в стабильную.

Если хочется пройти этот путь системно — от протокольных запросов ConsumerGroupHeartbeat до production-настройки групп с лабораторными работами, — посмотрите полный курс Apache Kafka (вводные уроки бесплатны) и другие материалы направления Data Engineering.

Ещё в направлении · Data Engineering

Все материалы направления →