MAATRIX / Блог / Отставание потребителя: что эта цифра на самом деле говорит о вашей системе

Отставание потребителя: что эта цифра на самом деле говорит о вашей системе

MAATRIX

Если вы открыли дашборд с очередями сообщений и увидели растущую линию consumer lag, первый порыв — паниковать. Второй, более частый — проигнорировать, потому что «система же работает, ничего не упало». Оба порыва ошибочны. Consumer lag — это не индикатор аварии, это индикатор темпа: он показывает, успевает ли ваша обработка за тем, что производится, причём задолго до того, как это станет заметно пользователям. Разберём, что это за цифра, откуда она берётся и как её читать, чтобы не пропустить момент, когда система ещё жива, но уже не успевает.

Что физически считается под словом «lag»

В основе любой очереди сообщений — будь то Kafka, RabbitMQ с очередями-стримами, NATS JetStream или Redis Streams — лежит одна и та же идея: у потока сообщений есть последовательные позиции. В Kafka это офсет (offset) — целочисленный порядковый номер сообщения внутри партиции топика. Producer дописывает сообщения в конец лога партиции, и офсет последнего записанного сообщения называется log end offset (LEO, или high water mark).

Consumer, читая сообщения, периодически фиксирует (коммитит) offset, до которого он гарантированно обработал данные — committed offset. Разница между этими двумя значениями и есть lag:

lag = log end offset − committed offset

Это разница не во времени, а в количестве непрочитанных (или непрокоммиченных) записей. Именно поэтому lag в 10 000 сообщений может означать секунды отставания для лёгкого потока телеметрии и часы — для потока, где на каждое сообщение уходит тяжёлый запрос к внешнему API.

В Kafka эта пара чисел хранится и вычисляется на брокере и в consumer group coordinator: брокер знает LEO по каждой партиции, а committed offset хранится в служебном топике __consumer_offsets (или, в устаревших схемах, в ZooKeeper — сейчас это редкость после перехода на KRaft). Инструмент kafka-consumer-groups.sh --describe буквально вычитывает оба числа и вычитает одно из другого — никакой магии внутри нет.

Важный нюанс: lag считается на партицию, а не на топик в целом. У топика с 12 партициями будет 12 отдельных значений lag, и агрегированная цифра на дашборде — это обычно их сумма. Если она растёт, стоит сразу смотреть разбивку по партициям: часто выясняется, что 11 партиций в нуле, а одна растёт бесконечно — это уже не вопрос производительности consumer group в целом, а конкретная партиция с «горячим» ключом или зависшим потребителем.

Почему растущий lag — это сигнал, а не приговор

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

  1. Producer резко ускорился. Пришёл всплеск трафика, партнёр включил массовую выгрузку, cron-задача выгрузила батч за час вместо равномерного потока. Consumer работает с той же скоростью, что и всегда, но входящий поток стал шире.
  2. Consumer замедлился или встал. Упал под нагрузкой внешний сервис, в который consumer стучится на каждое сообщение; воркер завис в дедлоке; под capacity consumer'а (CPU, память, сеть) просело из-за соседних процессов на той же ноде.
  3. Consumer недогружен ресурсами относительно партиций. Классика: у топика 24 партиции, а в consumer group — 3 инстанса. Каждый инстанс тянет по 8 партиций последовательно (в рамках одного потока), и суммарной пропускной способности не хватает даже при ровной нагрузке.

Разница между первым и вторым сценарием критична для реакции. В первом случае система обычно самовосстанавливается: всплеск закончился, consumer постепенно «доедает» накопленное, lag снижается сам. Во втором — без вмешательства lag будет расти бесконечно, потому что причина не во внешнем всплеске, а во внутренней деградации.

Отличить их просто по одному графику: смотрите не только на lag, но и на скорость роста LEO (сколько сообщений producer пишет в секунду) и скорость роста committed offset (сколько сообщений consumer реально обрабатывает в секунду) отдельно. Если обе линии растут, но первая быстрее — это сценарий 1 или 3. Если линия committed offset легла в ноль или сильно просела, а LEO продолжает расти как обычно — это сценарий 2, и здесь нужно разбираться, что случилось с потребителем прямо сейчас, а не ждать самовосстановления.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Почему нулевой lag — не всегда хорошая новость

Здесь легко попасть в ловушку дашборда: «lag = 0» на панели мониторинга выглядит как идеальное здоровье системы. На практике ноль означает одно из двух:

  • Consumer действительно успевает обрабатывать всё, что производится, с запасом — это то, к чему вы стремитесь.
  • Consumer простаивает, потому что ему нечего обрабатывать (producer сам ничего не пишет), или он подписан на партиции, которые фактически не используются.

Разница не видна на графике lag как такового — нужен второй график, throughput (сообщений в секунду) на стороне producer. Если throughput тоже около нуля, а не только lag — это не «отличная работа consumer», это либо ожидаемое затишье (например, ночью для бизнес-топика с рабочими часами), либо симптом того, что producer сам не пишет данные — а это может значить, что где-то выше по цепочке сломался сбор данных, а не сама очередь.

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

Поэтому правило простое: смотреть на lag в паре с throughput. Ноль lag при ненулевом стабильном throughput — это хорошо. Ноль lag при нулевом throughput — это нейтрально или подозрительно, в зависимости от того, ожидаете вы затишье или нет.

Как lag предупреждает о проблеме раньше отказа

Ценность consumer lag как метрики именно в том, что он реагирует раньше, чем система «упадёт» в привычном смысле — раньше, чем появятся таймауты у пользователей, раньше, чем очередь займёт весь диск, раньше, чем алерты по CPU или памяти сработают на самом consumer.

Разберём типичную деградацию по шагам:

  1. Внешняя зависимость consumer (база данных, сторонний API, другой микросервис) начинает отвечать на 20% медленнее обычного — само по себе это может не триггерить ни один алерт по latency, если порог настроен грубо.
  2. Consumer, который раньше обрабатывал сообщение за 50 мс, теперь тратит 65 мс. Пропускная способность падает примерно на 20%, а входящий поток не изменился.
  3. Lag начинает расти линейно — это первый видимый симптом всей цепочки, и он появляется задолго до того, как что-то ещё «покраснеет» на дашборде.
  4. Если ничего не сделать, накопление продолжается: у Kafka есть retention (по времени или по объёму), и если consumer не успевает вычитать данные до истечения retention, они физически удаляются брокером — теряются безвозвратно, даже если consumer их так и не прочитал.
  5. Дальше — если у producer вообще нет буфера или backpressure, он может начать получать ошибки записи при достижении лимитов диска на брокере, и это уже полноценный инцидент, видимый пользователям.

То есть между «стало на 20% медленнее» и «пользователи видят ошибки» может пройти много времени — и именно это время даёт вам lag, если вы на него смотрите. Это метрика с фазой предупреждения, а не только фиксации факта.

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

Практика: где смотреть lag и как его снимать

Для Kafka самый прямой способ — штатная утилита:

kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
  --describe --group my-consumer-group

Вывод покажет по каждой партиции: CURRENT-OFFSET, LOG-END-OFFSET и уже посчитанный LAG. Это удобно для разовой диагностики, но неудобно для непрерывного мониторинга — команду не будешь дёргать каждую секунду в проде.

Для постоянного сбора метрик практичнее подключить экспортёр, который снимает lag прямо из JMX-метрик брокера или через Consumer API и отдаёт их в формате Prometheus. Дальше — обычная связка: Prometheus собирает метрику kafka_consumergroup_lag (имя зависит от конкретного экспортёра) по каждой паре топик/партиция/группа, Grafana строит график, Alertmanager реагирует на правило. Если у вас ещё не развёрнута такая связка, разворачивание Prometheus и Grafana на VPS — это отдельная, но короткая задача: подробно описано в статье про установку Grafana и Prometheus на VPS.

Для RabbitMQ аналог lag — глубина очереди (messages_ready + messages_unacknowledged) в связке со скоростью consume (ack rate) — идея та же, только единица измерения не offset, а количество сообщений в очереди напрямую, без разделения на партиции.

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

- alert: ConsumerLagGrowing
  expr: deriv(kafka_consumergroup_lag[10m]) > 0
  for: 10m
  labels:
    severity: warning
  annotations:
    summary: "Lag группы {{ $labels.consumergroup }} на топике {{ $labels.topic }} устойчиво растёт"

Такое правило точнее, чем статичный порог lag > 100000, потому что нормальное значение lag сильно зависит от топика: для потока кликов на сайте с высокой партицированностью норма — тысячи, для очереди биллинговых транзакций норма — единицы.

Если очередь работает поверх собственного брокера сообщений, а не готового Kafka/RabbitMQ, механика хранения и чтения устроена похоже — она разобрана в статье как брокер хранит сообщения на диске, а базовые принципы устройства очереди — в статье как устроена очередь сообщений.

Частые ошибки при интерпретации lag

Ошибка: сравнивать lag разных топиков напрямую. Абсолютное число lag ничего не значит без контекста скорости обработки конкретного топика. Lag в 5000 для топика с одним сообщением в секунду — это почти полтора часа отставания. Lag в 5000 для топика с 50 000 сообщений в секунду — это доли секунды. Правильнее переводить lag в оценочное время через деление на текущий throughput потребителя, а не сравнивать сырые числа между топиками.

Ошибка: ставить один порог алерта на все consumer group сразу. Разные группы обслуживают разные по характеру потоки, и единый глобальный порог либо будет постоянно шуметь на «шумных» топиках, либо пропустит реальную проблему на «тихих». Про то, как вообще выбирать пороги алертов так, чтобы они не превращались в фоновый шум, есть отдельный разбор — алерты, которые не бесят: пороги.

Ошибка: не различать lag по коммиту и lag по факту обработки. Если consumer коммитит offset сразу после чтения сообщения из брокера, а не после того, как реально обработал его (записал в базу, отправил дальше), то lag по коммиту будет заниженным — он покажет, что сообщение «прочитано», хотя бизнес-логика по нему ещё не выполнена. Это разные гарантии: at-most-once по коммиту офсета не то же самое, что «данные обработаны». Для честной картины нужно смотреть либо на коммит после обработки (at-least-once с риском повторной обработки при сбое), либо вести отдельную метрику завершения бизнес-логики параллельно с lag брокера.

Ошибка: игнорировать lag на неактивных партициях. Если partitioning сделан по ключу с неравномерным распределением (например, по ID клиента, где один крупный клиент генерирует 40% трафика), одна партиция может копить lag, пока остальные в порядке. Суммарная метрика по топику это замаскирует — нужен разбор по партициям, не только агрегат.

Нужен сервер под эту задачу?

Разверните VPS MAATRIX за пару минут: NVMe, AMD EPYC, root-доступ, локации UK, США, Франция и РФ. Оплата картой РФ и по СБП.

Арендовать сервер

Нужны сами нейросети для контента?

Генерируйте изображения, видео и озвучку нейросетями на falapi.io — десятки моделей в одном окне. Оплата картой РФ и по СБП.

Частые вопросы

Какой lag считается нормальным?

Универсального числа нет — норма зависит от throughput топика и требований к задержке в бизнес-логике. Практичнее ориентироваться на время (lag ÷ текущая скорость обработки) и на тренд, а не на абсолютное значение: устойчивый рост важнее конкретной цифры.

Может ли lag быть отрицательным?

Формально нет — committed offset не может обогнать log end offset при штатной работе. Если в инструментах мониторинга видно отрицательное значение, чаще всего это артефакт гонки при снятии метрики (offset снялся до записи нового сообщения) или баг в экспортёре, а не реальное состояние данных.

Что делать, если lag растёт, а добавить ресурсы consumer'у некуда?

Сначала проверить, не упирается ли consumer во внешнюю зависимость (БД, API) — добавление CPU/RAM самому consumer'у не поможет, если узкое место не в нём. Если узкое место действительно в потребителе, а партиций больше, чем работающих инстансов, — добавление ещё одного инстанса в consumer group часто даёт кратный эффект, потому что Kafka автоматически перераспределит партиции между инстансами.

Почему lag скачет каждый раз при деплое?

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

Нужно ли мониторить lag, если сообщений в системе немного?

Да, если задержка обработки критична для бизнеса (платежи, уведомления, синхронизация остатков) — даже на низком трафике lag покажет, если consumer систематически не успевает или падает. На чисто фоновых, не чувствительных к задержке потоках (например, архивные логи) можно ограничиться более редкой проверкой без алертов в реальном времени.

Обсудить статью, задать вопрос или начать новую тему

Есть вопрос по этой статье, идея для обсуждения или просто хотите поделиться опытом? Сообщество MAATRIX ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.

Перейти в сообщество →