Как мониторить lag консьюмера
Lag показывает, насколько потребитель отстает от продюсера в реальном времени.
🟢 Junior Level
Что такое Consumer Lag простыми словами
Consumer Lag (отставание консьюмера) — это числовая разница между последним записанным сообщением в партиции Kafka и последним сообщением, которое консьюмер прочитал и закоммитил.
Lag показывает, насколько потребитель отстает от продюсера в реальном времени.
Лог партиции на брокере:
[0] ... [800: Закоммичено] ─── (Lag: 200 сообщений) ───> [1000: Log End Offset]
▲ ▲
│ │
Committed Offset (Консьюмер) Log End Offset (Продюсер)
- Log End Offset (LEO): смещение следующего сообщения, которое будет записано в партицию лидером.
- Committed Offset: смещение, зафиксированное консьюмером в системном топике
__consumer_offsets. - Lag = 0: консьюмер работает в реальном времени, обрабатывая все входящие данные без задержек.
- Lag растет: продюсер генерирует сообщения быстрее, чем сервис успевает их обрабатывать (узкое место в CPU, БД или внешних вызовах).
Базовый просмотр через консольную утилиту Kafka CLI
Стандартная CLI-утилита kafka-consumer-groups.sh позволяет получить моментальный снимок (snapshot) состояния группы:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group order-processing-group
Пример вывода консоли:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST
order-processing-group orders 0 145000 150000 5000 consumer-1-abc /10.0.1.15
order-processing-group orders 1 148000 148050 50 consumer-2-def /10.0.1.16
order-processing-group orders 2 120000 120000 0 consumer-3-ghi /10.0.1.17
🟡 Middle Level
Программное вычисление Lag через Java AdminClient
Вместо сторонних скриптов Java-сервис может самостоятельно рассчитывать свой lag с помощью официального AdminClient:
public class KafkaLagChecker {
public static Map<TopicPartition, Long> getConsumerGroupLag(
String bootstrapServers,
String groupId,
String topic) throws ExecutionException, InterruptedException {
Properties props = new Properties();
props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
try (AdminClient adminClient = AdminClient.create(props)) {
// 1. Получаем закоммиченные смещения группы консьюмеров
Map<TopicPartition, OffsetAndMetadata> committedOffsets = adminClient
.listConsumerGroupOffsets(groupId)
.partitionsToOffsetAndMetadata()
.get();
// 2. Формируем запрос на получение актуального LogEndOffset (Latest Offset)
Map<TopicPartition, OffsetSpec> requestLatestOffsets = committedOffsets.keySet().stream()
.filter(tp -> tp.topic().equals(topic))
.collect(Collectors.toMap(Function.identity(), tp -> OffsetSpec.latest()));
Map<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> endOffsets = adminClient
.listOffsets(requestLatestOffsets)
.all()
.get();
// 3. Вычисляем Lag для каждой партиции: LEO - CommittedOffset
Map<TopicPartition, Long> lagMap = new HashMap<>();
for (Map.Entry<TopicPartition, OffsetSpec> entry : requestLatestOffsets.entrySet()) {
TopicPartition tp = entry.getKey();
long currentOffset = committedOffsets.get(tp).offset();
long logEndOffset = endOffsets.get(tp).offset();
lagMap.put(tp, Math.max(0, logEndOffset - currentOffset));
}
return lagMap;
}
}
}
Архитектура промышленного мониторинга
В production-контурах опрос через CLI или периодический AdminClient неприменим из-за высокой нагрузки на брокер. Используется связка Prometheus + Grafana:
[Kafka Cluster] ──JMX Exporter──> [Prometheus] ──PromQL──> [Grafana Dashboard]
│ │
[kafka-exporter] ───────────────────────┘ ▼
[Alertmanager]
(Slack / PagerDuty / Telegram)
- kafka-exporter (написан на Go):
Опрашивает брокеры по легковесному протоколу Kafka и экспортирует метрики в формате Prometheus:
kafka_consumergroup_lag{group="...",topic="...",partition="..."}kafka_consumergroup_lag_sum{group="...",topic="..."}(суммарный лаг по всем партициям топика)
- JMX-метрики самого консьюмера (
org.apache.kafka.clients.consumer):records-lag-max: максимальный лаг среди всех назначенных партиций данного инстанса консьюмера.records-lead-min: расстояние до самого старого сообщения (показывает риск удаления данных политикой retention до того, как консьюмер успеет их прочитать!).
🔴 Senior Level
Тонкости вычисления Lag: LEO vs HW vs LSO
При поверхностном взгляде формула LEO - CommittedOffset кажется простой. Однако в корпоративных средах существуют три разных верхних указателя смещения:
ПАРТИЦИЯ НА БРОКЕРЕ
Offset: ... 100 ──── 101 ──── 102 ──── 103 ──── 104 ──── 105 ──── 106
▲ ▲ ▲
│ │ │
LSO HW LEO
(Last Stable Offset) (High Watermark) (Log End Offset)
- Log End Offset (LEO): офсет следующего физического сообщения на лидере. Включает нереплицированные и незакоммиченные транзакционные данные.
- High Watermark (HW): максимальный офсет, подтвержденный всеми репликами из ISR (In-Sync Replicas). Обычный консьюмер с
isolation.level=read_uncommittedможет читать до HW. - Last Stable Offset (LSO): офсет первой незавершенной транзакции (Ongoing Transaction). Консьюмер с
isolation.level=read_committedникогда не прочитает сообщения выше LSO, даже если они уже подтверждены репликами!
Инженерная ловушка: Если висящая транзакция продюсера не завершается часами (забытый
commitTransaction()или сбой координатора), LSO застревает на месте. Консьюмер сread_committedблокируется. Мониторинг по LEO покажет растущий лаг, но проблема заключается не в медлительности консьюмера, а в зависшей транзакции продюсера!
Смещение во времени (Time-based Lag vs Message Count)
Мониторинг лага исключительно в штуках сообщений — частый антипаттерн:
- Лаг в 50 000 сообщений при нагрузке 200 000 RPS — это задержка всего в 250 миллисекунд (норма).
- Лаг в 50 сообщений в топике тяжелых XML-отчетов — это отставание на 4 часа (катастрофа).
Поэтому зрелые команды мониторят Time Lag (отставание по времени в секундах): \(\text{Time Lag} = \text{Current Timestamp} - \text{Timestamp последнего закоммиченного сообщения}\)
Инструмент LinkedIn Burrow или Kafka Exporter с включенным тайм-трекингом вычисляют задержку по временным меткам лога (CreateTime / LogAppendTime), предоставляя статус здоровья группы:
OK— лаг стабилен или уменьшается.WARN— лаг не уменьшается на протяжении нескольких интервалов.ERROR— лаг монотонно растет, скорость потребления ниже скорости генерации.STOP— консьюмер перестал коммитить офсеты (поток завис или упал).
Автоматическое масштабирование консьюмеров (KEDA в Kubernetes)
В Cloud-Native архитектуре метрика Lag используется для горизонтального автомасштабирования подов (Horizontal Pod Autoscaler) через KEDA (Kubernetes Event-driven Autoscaling):
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
name: order-consumer-scaler
spec:
scaleTargetRef:
name: order-consumer-deployment
minReplicaCount: 2
maxReplicaCount: 16 # Не может быть больше количества партиций в топике!
triggers:
- type: kafka
metadata:
bootstrapServers: kafka-cluster-kafka-bootstrap:9092
consumerGroup: order-processing-group
topic: orders
lagThreshold: "1000" # На каждые 1000 сообщений лага добавлять +1 под
activationLagThreshold: "100"
4 Tricky Questions
1. Почему утилиты мониторинга могут показывать Lag = 0, в то время как консьюмер на самом деле полностью «мертв» (упал или завис)?
Ответ:
Lag вычисляется как математическая разность между LEO партиции и последним сохраненным CommittedOffset.
- Если продюсер в ночное время не отправляет сообщений в топик, LEO не растет.
- Если консьюмер обработал все дневные сообщения, закоммитил офсет (например, 50 000), а затем процесс упал по
OutOfMemoryErrorили сетевой сокет разорвался, значениеCommittedOffsetв__consumer_offsetsостается равным 50 000. - Формула возвращает: $\text{Lag} = 50000 - 50000 = 0$.
- Мониторинг по чистому лагу рапортует, что все в порядке.
- Решение: Мониторить лаг необходимо в связке со статусом членства группы:
kafka_consumergroup_members(число активных членов группы).- Метрика брокера
last-heartbeat-seconds-ago. Если активных консьюмеров в группе нет, статус группы должен переходить вDEAD/EMPTY.
2. У консьюмера в топике из 16 партиций суммарный лаг составляет 100 000 сообщений. При этом 15 партиций имеют Lag = 0, а партиция № 4 имеет Lag = 100 000. В чем причина такого перекоса (Lag Skew), и как его устранить?
Ответ: Это классический симптом перекоса ключей партиционирования (Hot Key Problem) либо зависания отдельного потока:
- Hot Key: Продюсер отправляет сообщения со статическим или несбалансированным ключом (например, все события крупного оптового клиента отправляются с
key = "VIP_CUSTOMER_1"). По формулеmurmur2(key) % numPartitionsабсолютно все эти события маршрутизируются строго в партицию 4. Потребитель партиции 4 захлебывается, а остальные 15 потребителей простаивают. - Poison Pill в партиции 4: В партицию 4 попало сообщение с битым телом, консьюмер крутится в бесконечном блокирующем цикле ретраев и не может сдвинуть офсет.
- Устранение:
- Если виноват Hot Key — добавить к ключу случайный суффикс (соление ключа / key salting:
"VIP_CUSTOMER_1_salt" + random(1..4)), либо пересмотреть схему ключей. - Если виноват Poison Pill — настроить отправку в Dead Letter Queue (DLQ).
- Если виноват Hot Key — добавить к ключу случайный суффикс (соление ключа / key salting:
3. Чем отличается Committed Offset Lag от реального Processing Lag в многопоточном консьюмере?
Ответ:
В архитектурах, где поток poll() только забирает пачки сообщений, а обработка делегируется пулу воркеров (ExecutorService / Virtual Threads):
- Processing Offset: офсет сообщений, которые физически находятся в памяти воркеров и обрабатываются в данный момент.
- Committed Offset: офсет, который был фактически зафиксирован на брокере.
- Внешние инструменты (Prometheus, CLI, Burrow) имеют доступ только к Committed Offset на брокере.
- Если приложение коммитит офсеты пачками раз в 10 секунд (или асинхронно через
commitAsync()), внешний лаг будет показывать скачки («пилу») в тысячи сообщений, хотя реального отставания бизнес-обработки нет. Для точного мониторинга приложение должно отдавать внутренние Micrometer-метрики фактически обработанных записей.
4. Что такое метрика records-lead-min, и почему она может быть важнее, чем records-lag-max?
Ответ:
records-lag-maxизмеряет расстояние консьюмера вперед до конца лога (сколько осталось догнать до новых сообщений).records-lead-minизмеряет расстояние консьюмера назад до начала лога (Log Start Offset), то есть сколько сообщений отделяет текущую позицию консьюмера от зоны физического удаления старых данных.- Опасность: Если на топике настроен агрессивный retention (например, хранить данные 2 часа), и консьюмер отстал на 1 час 55 минут, его
records-lead-minстремится к нулю. Как только lead упадет до 0, брокер удалит сегмент с офсетом, который консьюмер еще не успел вычитать. При следующем чтении консьюмер выбросит фатальное исключениеOffsetOutOfRangeException, произойдет потеря данных. Поэтому низкийrecords-lead— это индикатор неминуемой потери данных.
🎯 Шпаргалка для интервью
30-секундный ответ (Elevator Pitch)
«Consumer Lag — это разница между последним смещением в партиции (Log End Offset) и закоммиченным смещением группы (Committed Offset). Он показывает объем необработанных сообщений. Мониторить лаг необходимо в разрезе каждой партиции (
Lag per partition), чтобы вовремя выявлять перекосы ключей (Hot Key) и зависшие партиции (Poison Pills). В продакшене мониторинг организуется через связку kafka-exporter + Prometheus + Grafana, отслеживая не только абсолютное число сообщений, но и производную роста лага (Lag Growth Rate) и отставание по времени (Time Lag). Для предотвращения ложноположительных алертов мониторинг лага всегда дополняется проверкой наличия активных консьюмеров в группе».
Ключевые метрики для Production-мониторинга
kafka_consumergroup_lag_sum: общий лаг группы по топику.kafka_consumergroup_lag: лаг конкретной партиции (для поиска перекосов).deriv(kafka_consumergroup_lag[5m]): скорость роста лага (если $> 0$ — потребитель не справляется).kafka_consumergroup_members: число активных подов-консьюмеров.records-lead-min: запас расстояния до Log Start Offset (угрозаOffsetOutOfRangeException).
Красные флаги (Чего говорить НЕЛЬЗЯ)
- ❌ «Если Lag = 0, значит сервис работает идеально» — Сервис может быть отключен или лежать в OOM; если продюсер ничего не пишет, лаг останется нулевым.
- ❌ «Достаточно мониторить только суммарный лаг топика» — Суммарный лаг скрывает перекос, когда 1 партиция намертво заблокирована, а остальные 15 имеют нулевой лаг.
- ❌ «При росте лага нужно просто масштабировать консьюмеры до 100 подов» — Количество консьюмеров в группе ограничено числом партиций топика. Консьюмеры сверх этого лимита будут простаивать.
- ❌ «Lag всегда измеряется как LEO минус committed offset» — При
read_committedконсьюмер не может читать дальше Last Stable Offset (LSO), поэтому открытая транзакция искусственно удерживает консьюмер.