📨 Розділ 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): зміщення (offset) наступного повідомлення, яке буде записане в партицію лідером.
  • 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 повідомлень у топіку важких аналітичних PDF/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. Рішення: Моніторити лаг необхідно у суворій зв’язці зі статусом членства групи (Consumer Group Membership):
    • 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. Консьюмер, закріплений за цією партицією, захлинається, тоді як інші 15 інстансів простоюють без роботи.
  2. Poison Pill у партиції 4: У партицію 4 надійшло пошкоджене повідомлення, консьюмер крутиться в нескінченному циклі повторних спроб (retry loop) і не може просунути офсет уперед.
  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 або Java 21 Virtual Threads):

  1. Processing Offset: офсет повідомлень, які фізично перебувають у пам’яті воркерів і обробляються прямо зараз.
  2. Committed Offset: офсет, який був фактично зафіксований на брокері в топіку __consumer_offsets.
  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 Keys) або зависання на “отруйних пігулках” (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; якщо нових повідомлень немає, лаг залишиться нульовим.
  • ❌ «Достатньо моніторити лише сумарний лаг топіка» — Сумарне значення маскує критичний перекіс, коли одна партиція намертво заблокована, а решта 15 мають нульовий лаг.
  • ❌ «При зростанні лагу треба просто масштабувати консьюмери до 100 подів» — Кількість активних консьюмерів у групі фізично обмежена кількістю партицій топіка. Усі поди понад цю кількість будуть простоювати.
  • ❌ «Lag завжди дорівнює LEO мінус committed offset» — При режимі read_committed консьюмер не може читати далі Last Stable Offset (LSO), тому незакрита транзакція продюсера блокує споживача незалежно від LEO.

Пов’язані теми