Под капотом Kafka: путь сообщения от send() до commit offset. Часть 2

Привет, Хабр! Меня зовут Максим Шуматбаев, я инженер технической поддержки в Arenadata. Это вторая часть разбора, что на самом деле происходит при работе клиентских приложений с Kafka-кластером. В первой части мы проследили запись от получения метаданных до транзакций и exactly-once. Теперь очередь второй половины пути: как потребитель получает данные, делит партиции с другими участниками группы, фиксирует прогресс и почему иногда начинает отставать.

1. Как потребитель читает
У потребителя свои механизмы, свои таймауты и свои тонкости настройки, которые влияют на эффективность чтения. Начнем с принципиального выбора, на котором построено все чтение в Kafka.
Pull вместо push
Push-уведомлений в Kafka нет, брокеры никак не могут самостоятельно доставлять сообщения до клиентов. Потребитель сам запрашивает данные по модели pull. Чтение — это цикл запросов Fetch: потребитель просит сообщения с offset X, брокер отдает батч из лога с этого offset, потребитель обрабатывает и при следующем poll() просит следующую порцию.
Сделано так не случайно, так как pull-модель отдает потребителю контроль над темпом чтения. В push-системах быстрый продюсер может затопить медленного потребителя новыми сообщениями, заваливая его данными быстрее, чем тот успевает их обработать. В Kafka потребитель вытягивает данные тогда, когда он сам хочет.
Цикл pull по шагам:
Потребитель шлет
Fetch(offset = N), просит данные с offset N.Брокер возвращает батч записей с offset N.
Потребитель обрабатывает и повторяет цикл со следующим offset.
В высокоуровневом API (Java KafkaConsumer) все это спрятано за poll(). Но есть деталь: некоторые реализации позволяют кешировать прочитанные сообщения и отдавать их по частям, тем самым снижая нагрузку на брокеры из-за меньшего числа запросов.
Long polling и очередь ожидания (purgatory)
Обычный Fetch вернет пустоту, если новых данных нет. Чтобы не нагружать кластер пустыми запросами, есть long polling: брокер придерживает ответ, пока не накопятся данные. Этим управляют два параметра потребителя: fetch.min.bytes (минимальный объем) и fetch.max.wait.ms (максимальное ожидание на сервере). Данных меньше порога, брокер не отвечает сразу, а ждет, пока не наберется fetch.min.bytes либо не пройдет fetch.max.wait.ms.
Внутри брокера такой отложенный запрос уходит в очередь ожидания, которую исторически зовут purgatory, чистилище. Запрос застревает там между «выполнен» и «отклонен», ожидая условия, отсюда и название. Набрался нужный объем байт или вышло время — запрос покидает purgatory, формируется FetchResponse. Если же данные были сразу, запрос минует очередь и отправляет данные мгновенно.
Long polling резко режет число пустых ответов ценой небольшой искусственной задержки. При низкой нагрузке новые события будто запаздывают, пока брокер ждет, пока они накопятся. Если же fetch.min.bytes слишком большой, то потребитель долго сидит без данных, даже когда отдельные сообщения уже пришли. Куда же без компромисса между задержками и эффективностью.
Consumer Lag: что это и почему растет
Главная метрика эффективного чтения — это лаг потребителя (consumer lag), отставание от конца лога. Формально разница между последним offset в партиции (log end offset) и текущим зафиксированным offset группы. По-простому, сколько сообщений еще не обработано относительно самого свежего в топике.
Причины роста почти всегда сводятся к одному: поток входящих данных обгоняет способность их обрабатывать. Основные:
Медленная обработка. Потребитель долго возится с каждым сообщением, продюсер продолжает писать, разрыв растет
Всплеск трафика. Объем резко вырос (суточный пик, массовое событие, аномальный продюсер), потребители физически не успевают
Перекос партиций (skew). Неудачный ключ, одна перегретая партиция копит в себе сильно больше сообщений, чем остальные, ее потребитель отстает сильнее прочих. Сюда же слишком малое число партиций — больше потребителей, чем партиций, не повысит пропускную способность, так как на одну партицию приходится лишь один потребитель
Ребаланс группы. На время ребаланса чтение встает, лаг растет по всем затронутым партициям
Инфраструктура (I/O). Уперлись в диск или сеть брокера, даже быстрые потребители ограничены скоростью отдачи.
Не всякий лаг — это плохо. Небольшой стабильный лаг — это норма для здорового кластера, данные идут непрерывно, и между записью и чтением всегда короткий промежуток. Более того, такой буфер полезен, он сглаживает пики, потребитель разбирает всплеск чуть позже, но ровным темпом. Здесь важна именно сама динамика лага — неуклонно растущий или резко скакнувший лаг означает, что читающая сторона не справляется.
Лимиты на чтении
Симметрично записи, у потребителя свои лимиты на объем за один запрос: fetch.max.bytes (на весь fetch) и max.partition.fetch.bytes (на партицию). Тут стоит напомнить о необходимости корректной настройки лимитов между записью и чтением, о которой мы говорили в первой части в разделе про лимит размера сообщения. Например, продюсер настроен на большие сообщения, а у потребителя fetch.max.bytes слишком маленький, и чтение встает. Поэтому лимиты всегда согласуют на обоих концах тракта.
2. Consumer group: как устроена группа потребителей
Одиночный потребитель понятен, но очень быстро его станет не хватать. Несколько экземпляров потребителя делят между собой партиции топика и читают их параллельно, в случае падения одного, его партиции подхватывают остальные.
Два параллельных процесса
Потребитель делает одновременно две вещи: читает и обрабатывает данные, и участвует в групповом протоколе, поддерживая членство. Процессы изолированы и идут параллельно. В современных клиентах помимо основного цикла poll() есть фоновый поток heartbeat. Основной шлет Fetch и обрабатывает записи, фоновый периодически шлет координатору Heartbeat, мол, жив и держу партиции. Так, долгая обработка не мешает кластеру видеть, что потребитель активен.

poll() для получения и обработки записей, тогда как фоновый поток отправляет координатору группы сообщения о том, что consumer активен и получает подтверждения. Благодаря этому consumer не исключается из группы во время длительной обработки, поскольку координатор продолжает получать heartbeats от членов группыУ heartbeat две роли: сообщать координатору, что потребитель жив, и быть обратным каналом связи. В ответ на heartbeat координатор возвращает код вроде REBALANCE_IN_PROGRESS, так он говорит потребителю «пора переприсоединяться». То есть потребитель сам регулярно спрашивает «все в порядке?» и через ответ узнает о начале ребаланса.
Координатор и жизненный цикл группы
У каждой группы свой координатор, конкретный брокер. Выбирается следующим образом: Kafka хеширует group.id, берет остаток от деления на число партиций внутреннего топика __consumer_offsets, и лидер этой партиции становится координатором. На старте потребитель шлет FindCoordinator со своим group.id любому брокеру, узнает адрес координатора и начинает присоединение.
Жизненный цикл в классическом протоколе:
Новый потребитель шлет координатору
JoinGroupсо своимgroup.id.Координатор запускает ребаланс и ждет
JoinGroupот всех активных участников, включая уже состоящих, они тоже переподтверждают участие.Координатор назначает одного потребителя лидером группы (обычно первого присоединившегося). Лидер получает полный список членов и их метаданные.
Остальные (follower) получают пустой список.
Лидер раскладывает партиции по выбранной стратегии и формирует assignment, кто что читает.
Каждый шлет координатору
SyncGroup, лидер вкладывает распределение, follower шлют пустое.Координатор раздает каждому его список партиций в ответах на
SyncGroup.Группа переходит в Stable-состояние, все читают свои партиции.
Каждый цикл поднимает счетчик generation.id. Запрос OffsetCommit несет текущее поколение, и если потребитель шлет commit с устаревшим поколением, координатор отклонит его с ILLEGAL_GENERATION. Так, участник прошлого поколения не перезапишет offset после потери членства.
Статическое членство против мигающих рестартов. По умолчанию потребители динамические: вышел, координатор тут же удаляет тебя, вернулся, считаешься новым членом, запускается ребаланс. При частых перезапусках это превращается в шторм ненужных ребалансов. Помочь решить это может статическое членство (KIP-345): задайте каждому экземпляру стабильный group.instance.id, и координатор запомнит его. При кратковременном исчезновении он не станет сразу останавливать всю группу, а подождет возвращения того же ID в пределах увеличенного session.timeout.ms. Успел переподключиться — получил свои партиции обратно без ребаланса.
Heartbeat и session timeout
У каждого члена группы с координатором сессия длиной session.timeout.ms. Чтобы держать ее живой, потребитель шлет heartbeat не реже heartbeat.interval.ms (по умолчанию интервал 3 с, таймаут сессии 45 с, интервал обязан быть меньше таймаута). Пока heartbeat приходят, участник жив. Пропал на все окно session.timeout.ms - координатор помечает его Dead, забирает партиции и запускает ребаланс.

Слишком короткий session.timeout.ms дает ложные срабатывания: пауза GC или мелкий сетевой сбой, из-за чего потребитель пропустил пару heartbeat и запустился ненужный ребаланс. Слишком длинный — система медленно реагирует на реальный сбой, дольше держит партиции упавшего узла. Для чувствительных случаев увеличивают интервалы или берут статическое членство. Состояние можно отслеживать по метрикам heartbeat-rate, last-heartbeat-seconds-ago и прочим.
Таймаут на стороне клиент: max.poll.interval.ms
Есть еще таймер, и его постоянно путают с session timeout, хотя они про разное. max.poll.interval.ms (по умолчанию 5 минут) ограничивает интервал между вызовами poll(). Он нужен, когда приложение живое (heartbeat идут из фонового потока), но застряло в обработке и перестало вызывать poll(). Превысил max.poll.interval.ms, клиент сам помечает себя зависшим и покидает группу.
Разница принципиальна:
session.timeout.msследит за связью с кластером. Пропали heartbeat, координатор считает процесс мертвым;max.poll.interval.msследит за активностью обработки. Heartbeat идут, но данные не вычитываются, значит, обработка по факту встала, хотя формально по heartbeat все хорошо.
Пример из практики. Читаете крупный батч и обрабатываете его, скажем, 7 минут. Heartbeat летят из фонового потока, координатор доволен, а poll() не вызывался дольше 5 минут. Клиент решает, что потребитель завис, и запускает ребаланс, партиции едут к другому экземпляру, который читает тот же батч и снова застревает на 7 минут. Группа уходит в бесконечный цикл ребалансов и не двигается. Лечится ускорением обработки, уменьшением max.poll.records (читать меньше за раз) или увеличением max.poll.interval.ms под реальное время обработки.
Причины и цена ребаланса
Координатор пересобирает группу в нескольких случаях. Пришел новый потребитель с тем же group.id. Кто-то вышел через close() или LeaveGroup, либо пропал. Сменилось число партиций топика. Группа поменяла подписку через subscribe(). Любое из этого запускает новый цикл JoinGroup и SyncGroup, поднимает поколение, раздает новые assignment.
В классической реализации (eager rebalance) на время перераспределения вся группа встает. Каждый потребитель сначала отказывается от партиций, подтверждает готовность, потом получает новый assignment. Ребаланс длится 5 секунд, это 5 секунд простоя всей группы с последующим разбором накопленного. Плюс возможна повторная обработка и сброс локальных кэшей.
Для смягчения подобных ситуаций ввели кооперативный (инкрементальный) ребаланс. При изменениях отзываются не все партиции. Добавился новый потребитель, у существующих заберут лишь пару разделов для него, остальные продолжат работать без остановки. Протокол разбит на несколько этапов, зато большую часть времени группа читает. Цена — это чуть большая сложность и поддержка на стороне клиента (partition.assignment.strategy).
Что нового в Kafka 4.0, KIP-848. В 4.0 доступен новый протокол ребаланса группы. Главная идея, перенести вычисление assignment с клиента-лидера на брокер-координатор, что резко сокращает время ожидания и ускоряет ребалансы. Классический протокол (eager и cooperative) никуда не делся и работает, но для новых высоконагруженных сценариев смотреть стоит на новый.
3. Offset: фиксация прогресса чтения
Потребитель читает, но как кластер помнит, докуда он дочитал, чтобы после рестарта или ребаланса продолжить с нужного места? Для этого нужен offset и его фиксация.
Что такое offset
offset это позиция сообщения внутри партиции. Каждое сообщение получает порядковый номер, начиная с 0. Порядок гарантирован только внутри одной партиции, поэтому offset уникален лишь в пределах своей партиции и не является глобальным ID сообщения. Это координата в логе.
Читая партицию, потребитель использует offset для отслеживания прогресса. Обработал сообщение, может зафиксировать (commit) offset следующего. Обработал offset 25, при коммите сохраняет 26, следующий для чтения. То есть зафиксированный offset всегда указывает на позицию после последнего обработанного.

Такая реализация позволяет компактно хранить прогресс чтения: потребитель хранит одно число на партицию и может в любой момент перемотать offset назад и перечитать. Это сильно проще традиционных брокеров, где сервер пытался бы отслеживать подтверждение каждого отдельного сообщения.
Операция OffsetCommit
Фиксация offset — это отдельная операция. Технически потребитель шлет OffsetCommitRequest с последним обработанным offset плюс 1 по каждой партиции и своим group.id. Запрос идет координатору группы, который сохраняет offset во внутренний топик __consumer_offsets (ключ — это group.id плюс топик плюс партиция, значение это зафиксированный offset). Так, прогресс чтения хранится обычным логом внутри самой Kafka.
Commit успешен только после того, как запись попала в __consumer_offsets и реплицировалась на все реплики этого топика, лишь тогда координатор шлет OffsetCommitResponse. Поэтому при падении координатора или кластера зафиксированные offset не теряются.
Offset отражает прогресс всей группы. Участник упал или присоединился, координатор переназначает партиции, новый читатель партиции спрашивает у координатора (OffsetFetchRequest) последний зафиксированный offset и продолжает с него.
Отдельность операции commit дает гибкость, приложение само решает, когда считать сообщение обработанным. Но фиксировать после каждого сообщения дорого, это нагружает кластер и режет пропускную способность чтения. Поэтому по умолчанию Kafka коммитит автоматически с интервалом: enable.auto.commit=true, auto.commit.interval.ms=5000.
4. Эксплуатационные ограничения: квоты и throttling
Последний механизм влияет сразу на все, что мы разобрали в обеих частях этой статьи, — это квоты. Kafka умеет ограничивать ресурсы: пропускную способность конкретного пользователя или клиента и долю процессорного времени брокера на его запросы. Если клиент превысил квоту, то брокер не отдает ошибку, а искусственно тормозит обработку. В протоколе для этого есть поле ThrottleTimeMs в ответах — сигнал клиентской библиотеке, что ее запросы были искусственно замедлены.
Для приложения это внезапный рост времени операций при нулевом проценте ошибок. Никаких exception, все успешно, просто медленно. На графиках это рост latency при стабильном уровне ошибок.
Где искать, когда все работает, но медленно. Ключевые клиентские метрики — это produce-throttle-time-avg и fetch-throttle-time-avg. Ненулевые значения говорят, что брокеры намеренно задерживают ответы продюсеру или потребителю. Отдельно опасен асинхронный продюсер: при сильном throttling его буфер быстро забивается неотправленными сообщениями, и после исчерпания buffer.memory вызовы send() блокируются, приложение перестает принимать новые данные. Поэтому при отладке производительности всегда держите throttle-метрики на виду, квота — это частая и неочевидная причина замедления.
Заключение
Все, что мы разобрали, это лишь малая часть того, как устроена Kafka. За скобками остались Kafka Streams и Kafka Connect, тонкости репликации и работы контроллера, KRaft, устройство хранения на брокерах, безопасность, мониторинг кластера и многое другое. Каждая из этих тем заслуживает отдельного разговора.
На этом разбор пути сообщения через Kafka завершен, но сам цикл — нет. В следующих материалах продолжим погружаться в устройство Kafka. Если остались вопросы или есть темы, которые хотелось бы разобрать подробнее, пишите в комментариях.
KioskNews shows a cleaned-up reading view extracted from the publisher’s page — the original always lives on their site, not ours.