Отставание потребителя: что эта цифра на самом деле говорит о вашей системе
Если вы открыли дашборд с очередями сообщений и увидели растущую линию 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 — это разница скоростей, а не абсолютное состояние. Он растёт в трёх принципиально разных ситуациях, и путать их — частая ошибка:
- Producer резко ускорился. Пришёл всплеск трафика, партнёр включил массовую выгрузку, cron-задача выгрузила батч за час вместо равномерного потока. Consumer работает с той же скоростью, что и всегда, но входящий поток стал шире.
- Consumer замедлился или встал. Упал под нагрузкой внешний сервис, в который consumer стучится на каждое сообщение; воркер завис в дедлоке; под capacity consumer'а (CPU, память, сеть) просело из-за соседних процессов на той же ноде.
- 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.
Разберём типичную деградацию по шагам:
- Внешняя зависимость consumer (база данных, сторонний API, другой микросервис) начинает отвечать на 20% медленнее обычного — само по себе это может не триггерить ни один алерт по latency, если порог настроен грубо.
- Consumer, который раньше обрабатывал сообщение за 50 мс, теперь тратит 65 мс. Пропускная способность падает примерно на 20%, а входящий поток не изменился.
- Lag начинает расти линейно — это первый видимый симптом всей цепочки, и он появляется задолго до того, как что-то ещё «покраснеет» на дашборде.
- Если ничего не сделать, накопление продолжается: у Kafka есть retention (по времени или по объёму), и если consumer не успевает вычитать данные до истечения retention, они физически удаляются брокером — теряются безвозвратно, даже если consumer их так и не прочитал.
- Дальше — если у 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 ждёт. Для общения, пожалуйста, зарегистрируйтесь в нашем личном кабинете.
Перейти в сообщество →