📨 Раздел 15 · Вопрос #26

Как мониторить lag консьюмера

Lag показывает, насколько потребитель отстает от продюсера в реальном времени.


🟢 Junior Level

Что такое Consumer Lag простыми словами

Consumer Lag (отставание консьюмера) — это числовая разница между последним записанным сообщением в партиции Kafka и последним сообщением, которое консьюмер прочитал и закоммитил.

Lag показывает, насколько потребитель отстает от продюсера в реальном времени.

Лог партиции на брокере:
[0] ... [800: Закоммичено] ─── (Lag: 200 сообщений) ───> [1000: Log End Offset]
              ▲                                                    ▲
              │                                                    │
     Committed Offset (Консьюмер)                         Log End Offset (Продюсер)
\[\text{Lag} = \text{LogEndOffset (LEO)} - \text{CommittedOffset}\]
  • 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)
  1. kafka-exporter (написан на Go): Опрашивает брокеры по легковесному протоколу Kafka и экспортирует метрики в формате Prometheus:
    • kafka_consumergroup_lag{group="...",topic="...",partition="..."}
    • kafka_consumergroup_lag_sum{group="...",topic="..."} (суммарный лаг по всем партициям топика)
  2. 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)
  1. Log End Offset (LEO): офсет следующего физического сообщения на лидере. Включает нереплицированные и незакоммиченные транзакционные данные.
  2. High Watermark (HW): максимальный офсет, подтвержденный всеми репликами из ISR (In-Sync Replicas). Обычный консьюмер с isolation.level=read_uncommitted может читать до HW.
  3. 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.

  1. Если продюсер в ночное время не отправляет сообщений в топик, LEO не растет.
  2. Если консьюмер обработал все дневные сообщения, закоммитил офсет (например, 50 000), а затем процесс упал по OutOfMemoryError или сетевой сокет разорвался, значение CommittedOffset в __consumer_offsets остается равным 50 000.
  3. Формула возвращает: $\text{Lag} = 50000 - 50000 = 0$.
  4. Мониторинг по чистому лагу рапортует, что все в порядке.
  5. Решение: Мониторить лаг необходимо в связке со статусом членства группы:
    • 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) либо зависания отдельного потока:

  1. Hot Key: Продюсер отправляет сообщения со статическим или несбалансированным ключом (например, все события крупного оптового клиента отправляются с key = "VIP_CUSTOMER_1"). По формуле murmur2(key) % numPartitions абсолютно все эти события маршрутизируются строго в партицию 4. Потребитель партиции 4 захлебывается, а остальные 15 потребителей простаивают.
  2. Poison Pill в партиции 4: В партицию 4 попало сообщение с битым телом, консьюмер крутится в бесконечном блокирующем цикле ретраев и не может сдвинуть офсет.
  3. Устранение:
    • Если виноват Hot Key — добавить к ключу случайный суффикс (соление ключа / key salting: "VIP_CUSTOMER_1_salt" + random(1..4)), либо пересмотреть схему ключей.
    • Если виноват Poison Pill — настроить отправку в Dead Letter Queue (DLQ).

3. Чем отличается Committed Offset Lag от реального Processing Lag в многопоточном консьюмере?

Ответ: В архитектурах, где поток poll() только забирает пачки сообщений, а обработка делегируется пулу воркеров (ExecutorService / Virtual Threads):

  1. Processing Offset: офсет сообщений, которые физически находятся в памяти воркеров и обрабатываются в данный момент.
  2. Committed Offset: офсет, который был фактически зафиксирован на брокере.
  3. Внешние инструменты (Prometheus, CLI, Burrow) имеют доступ только к Committed Offset на брокере.
  4. Если приложение коммитит офсеты пачками раз в 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), поэтому открытая транзакция искусственно удерживает консьюмер.

Связанные темы