Эта статья подготовлена в рамках курса «Java‑разработчик. Продвинутый уровень».

Привет, Хабр!

Три часа ночи, алерт: лаг группы перевалил за миллион и продолжает расти. Дежурный делает то, что сделал бы любой — поднимает число экземпляров с четырёх до двенадцати. Лаг после этого растёт быстрее.

Ещё полчаса уходит на то, чтобы понять, почему. Группа не обрабатывает ничего вообще: она непрерывно перебалансируется, каждый цикл съедает секунд сорок, а между циклами полезной работы почти не остаётся. Каждый новый экземпляр запускал очередную перебалансировку, так что дежурный своими руками добивал то, что пытался спасти.

Обидно тут не то, что решение оказалось неверным, а то, что по симптому его выбрать было нельзя. Растущий лаг дают минимум пять разных причин, и лечатся они по‑разному: где‑то экземпляры надо добавить, где‑то убрать, а где‑то вообще не трогать группу и идти чинить базу.

Ниже — как за десять минут понять, с чем именно вы имеете дело, и почему смотреть надо не на лаг, а на то, какие часы у консьюмера истекли.

Три таймера, которые решают судьбу консьюмера

Прежде чем смотреть на метрики, стоит понять, за чем следит брокер. За консьюмером наблюдают три независимых часов, и от того, какие из них истекли, зависит вообще всё.

  • Сессионные часы отмеряют сигналы жизни: клиентская библиотека шлёт их брокеру из отдельного потока, каждые heartbeat.interval.ms. Если брокер не получил ни одного за session.timeout.ms (45 секунд в современных версиях), консьюмер считается мёртвым и выбывает из группы. При этом поток сигналов отдельный, и обработка сообщений его не блокирует.

  • Опросные часы следят за другим — за тем, чтобы между двумя вызовами poll() проходило не больше max.poll.interval.ms, по умолчанию 5 минут. Здесь наоборот: обработка выполняется в том же потоке, поэтому долгая пачка напрямую растягивает интервал. Превысили — брокер решает, что консьюмер завис, и тоже выкидывает его из группы.

Разница между двумя таймерами и есть первый разделяющий вопрос. Сигналы жизни идут, а опрос задерживается — виновата обработка. Сигналы не идут — виноват процесс целиком: пауза сборщика мусора, зависание на сети, нехватка процессора на узле.

  • Третьи часы таймером не являются вовсе, за ними стоит факт: осознанная смена состава группы. Деплой, масштабирование, перезапуск пода. Тут перебалансировка законна, вопрос только в её цене.

Первые две команды

Диагностика начинается не с дашборда, а с состояния группы.

kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --describe --group orders-processor
GROUP             TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG      CONSUMER-ID       HOST
orders-processor  orders  0          4812003         4812150         147      consumer-1-a8f2   /10.0.2.14
orders-processor  orders  1          4790112         5904881         1114769  consumer-2-b3c1   /10.0.2.19
orders-processor  orders  2          4811840         4812150         310      consumer-3-d7e4   /10.0.2.22
orders-processor  orders  3          4811995         4812150         155      consumer-4-f1a9   /10.0.2.27

Лаг сосредоточен на одной партиции из четырёх, остальные держатся у нуля. Значит, группа в целом справляется, а конкретная партиция — нет.

Повторите ту же команду через полминуты и посмотрите, меняется ли столбец с идентификатором консьюмера. Если идентификаторы каждый раз новые, а иногда вместо них прочерк — группа перебалансируется, и всё остальное можно не смотреть, пока это не починено.

Состояние группы видно отдельно:

kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --describe --group orders-processor --state
GROUP             COORDINATOR       ASSIGNMENT-STRATEGY  STATE                #MEMBERS
orders-processor  broker-2:9092     range                PreparingRebalance   12

Состояние PreparingRebalance вместо Stable при повторных вызовах — та самая история из начала статьи.

Одна партиция отстаёт, остальные нет

Разберём случай из вывода выше. Отставание на одной партиции означает одно из двух, и различить их несложно.

Либо ключ партиционирования перекошен: один клиент, один магазин, один идентификатор устройства генерирует непропорционально много событий, и всё это летит в одну партицию. Проверяется подсчётом ключей в отстающей партиции, хоть консольным консьюмером с ограничением по количеству.

kafka-console-consumer.sh --bootstrap-server broker:9092 \
  --topic orders --partition 1 --offset 4790112 \
  --max-messages 1000 --property print.key=true \
  | awk -F'\t' '{print $1}' | sort | uniq -c | sort -rn | head
   847 tenant-4471
    39 tenant-1203
    22 tenant-8890

847 сообщений из 1000 от одного арендатора — вот и перекос. Решается это на стороне продюсера: добавить в ключ ещё одну составляющую, если порядок внутри арендатора не критичен, либо выделить крупных отдельно.

Либо болен конкретный экземпляр, которому досталась эта партиция. Идентификатор консьюмера из первой таблицы даёт адрес узла, и дальше вопрос переезжает туда: паузы сборщика, нехватка процессора, зависший вызов наружу.

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

Лаг растёт равномерно, а процессор свободен

Другая картина: отстают все партиции сразу, группа стабильна, перебалансировок нет. Смотрим на потребление ресурсов и видим загрузку процессора процентов на 15.

Свободный процессор при растущем лаге почти всегда означает, что консьюмер не работает, а ждёт. Ждать он может базу, чужой API, ответ по сети — что угодно за пределами процесса.

Подтверждается это дампом потоков в момент нагрузки:

jcmd $(pgrep -f orders-processor) Thread.print | grep -A5 'kafka-consumer'
"orders-consumer-0" #24 prio=5 tid=0x00007f2a nid=0x1f4c runnable
   java.lang.Thread.State: RUNNABLE
        at java.net.SocketInputStream.socketRead0(Native Method)
        at org.postgresql.core.PGStream.receiveChar(PGStream.java:465)
        at ru.example.orders.OrderRepository.save(OrderRepository.java:88)

Поток в состоянии RUNNABLE, но стоит на чтении из сокета. Формально работает, фактически ждёт ответа базы. Десять таких потоков дадут ту же картину: лаг растёт, процессор простаивает.

Дальше расследование расходится надвое. Либо ускорять внешний вызов — индекс, пул соединений, кеш. Либо перестать делать по вызову на сообщение и писать пачками:

var records = consumer.poll(Duration.ofMillis(500));
var batch = new ArrayList<Order>(records.count());

for (var record : records) {
    batch.add(parse(record));
}
repository.saveAll(batch);      // один запрос вместо пятисот
consumer.commitSync();

Пакетная запись на нагрузке из мелких сообщений обычно даёт кратный выигрыш, потому что убирает круговые походы до базы. Именно круговые походы, а не саму запись, и стоило оптимизировать с самого начала.

Группа не может перестать перебалансироваться

Первое, что надо выяснить, — какие часы истекли. Логи консьюмера:

grep -E 'Member .* sending LeaveGroup|consumer poll timeout|Attempt to heartbeat failed' app.log | tail -5
14:22:07 INFO  Member consumer-2-b3c1 sending LeaveGroup request due to consumer poll timeout has expired.
14:23:11 INFO  Member consumer-4-f1a9 sending LeaveGroup request due to consumer poll timeout has expired.

Формулировка про истёкший опросный таймаут означает, что обработка пачки заняла больше max.poll.interval.ms. Сообщение про неудавшиеся сигналы жизни означало бы другое — что встал весь процесс.

Берёте 99-й перцентиль времени обработки одного сообщения и умножаете на max.poll.records. Двести миллисекунд на сообщение при пятистах записях за опрос дают 100 секунд на пачку, а таймаут по умолчанию 300. Пока укладываетесь. Стоит внешнему сервису притормозить вдвое, и вы за пределами.

Отсюда первое действие — не поднимать таймаут, а уменьшить пачку:

max.poll.records=100

Поднятие max.poll.interval.ms тоже работает, но обработка быстрее не стала, а обнаружение реально зависшего консьюмера отодвинулось на те же минуты.

Второе действие — сменить стратегию распределения партиций. По умолчанию действует та, при которой любое изменение состава отбирает у всех всё и раздаёт заново:

partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

Кооперативная стратегия переносит только те партиции, что реально меняют владельца.

К тому же немного про сам деплой. Перезапуск пода выглядит для группы как уход участника, и каждый под при раскатке запускает свой цикл. Статическое членство это снимает:

group.instance.id=orders-processor-3
session.timeout.ms=60000

Консьюмер с постоянным идентификатором, вернувшийся в пределах сессионного таймаута, получает свои партиции обратно без перебалансировки. Идентификатор должен быть стабильным между перезапусками, в Kubernetes для этого подходит порядковый номер пода из StatefulSet.

Если такие production‑сценарии для вас привычны, можно заодно проверить себя во вступительном тесте для Java‑разработчиков продвинутого уровня. Он покажет, какие темы уже закрыты, а где ещё есть пробелы.

Отравленное сообщение, которое не даёт группе сдвинуться

Есть сценарий, при котором лаг растёт, группа стабильна, процессор занят на сто процентов, а смещение не двигается ни на единицу.

grep -c 'Failed to process record' app.log
tail -3 app.log
418223
14:41:02 ERROR Failed to process record at offset 4790112: Cannot deserialize value
14:41:02 ERROR Failed to process record at offset 4790112: Cannot deserialize value
14:41:03 ERROR Failed to process record at offset 4790112: Cannot deserialize value

Одно и то же смещение в каждой строке. Консьюмер прочитал битую запись, упал на обработке, перезапустился, снова прочитал ту же запись, и так по кругу.

Отличается это от медленной обработки одним признаком: смещение вообще не растёт, тогда как при медленной оно ползёт, просто недостаточно быстро. Сравнение двух выводов --describe с интервалом в минуту отвечает на этот вопрос.

Фиксится тем, чтобы битая запись не останавливала поток:

for (var record : records) {
    try {
        process(record);
    } catch (Exception e) {
        log.error("не смогли обработать смещение {}", record.offset(), e);
        deadLetterProducer.send(toDeadLetter(record, e));
    }
}
consumer.commitSync();

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

Если чинить надо прямо сейчас, а выкладка новой версии займёт время, смещение сдвигается вручную. Только с предварительным прогоном, показывающим, что именно произойдёт:

kafka-consumer-groups.sh --bootstrap-server broker:9092 \
  --group orders-processor --topic orders:1 \
  --reset-offsets --shift-by 1 --dry-run

Группу для этого придётся остановить: активной группе смещения не сбросить.

Обработка длиннее любого разумного таймаута

Наивное решение — вынести обработку в пул потоков, чтобы цикл опроса не блокировался. Но пока задача выполняется в фоне, цикл уедет вперёд и закоммитит смещения необработанных сообщений. Падение процесса в этот момент означает потерю данных.

Можно приостановить партиции и продолжать опрашивать вхолостую:

var records = consumer.poll(Duration.ofMillis(500));

if (!records.isEmpty()) {
    consumer.pause(consumer.assignment());          // больше не выдавай записи
    executor.submit(() -> {
        process(records);
        pending.set(false);
    });
}

if (!pending.get() && !consumer.paused().isEmpty()) {
    consumer.resume(consumer.paused());
    consumer.commitSync();
}

Вызов poll() продолжает происходить регулярно и сбрасывает опросный таймер, но записей не возвращает — партиции приостановлены. Обработка идёт своим темпом в отдельном потоке, коммит происходит после её завершения.

Что делать в следующий раз

Проще держать в голове четыре вопроса

  1. Первый — жива ли вообще группа. Прогоняем --describe --state два раза с интервалом в минуту: если состояние скачет, а идентификаторы консьюмеров меняются, группа перебалансируется, и дальше можно не копать.

  2. Второй — отстают все партиции или одна. Одна значит перекос ключа либо больной узел, и экземпляры тут бесполезны. Все вместе — не тянет сама обработка.

  3. Третий вопрос про процессор: он свободен, а лаг растёт? Значит консьюмер не работает, а стоит и ждёт кого‑то снаружи, и дамп потоков покажет кого.

  4. Четвёртый добивает остаток: смотришь в логи, какие часы истекли. Опросный таймаут — обработка медленная. Пропали сигналы жизни — встал весь процесс.

А теперь то, ради чего всё это писалось. Добавление консьюмеров помогает только в одном случае из пяти: когда обработка упирается в процессор, а партиций больше, чем экземпляров.

В остальных четырёх новые поды либо просто простаивают, либо, как в той ночной истории, добивают группу очередной перебалансировкой при каждом запуске.

Когда Kafka начинает вести себя нестабильно в проде, важно понимать, что происходит внутри consumer group и как на её поведение влияют настройки клиента. Если хотите увереннее разбираться с такими ситуациями до ночного алерта и подбирать решение по причине, а не по симптому, на бесплатных открытых уроках OTUS можно продолжить тему на практике.

  • 26 августа в 20:00 — «Работа с Kafka через библиотеку Kafka Clients». Записаться

  • 21 сентября в 20:00— «HTTP-сервер на чистой Java за 30 минут». Записаться

Весь список открытых уроков августа собрали в дайджесте.