Мониторинг Kafka: consumer lag, throughput, Prometheus

Мониторинг Kafka: consumer lag, throughput, Prometheus

Kafka без мониторинга - это взрывное устройство с неизвестным таймером. Consumer отстал на миллион сообщений, а ты узнаёшь об этом от пользователей.

Главная метрика - consumer lag: разница между latest offset и committed offset. Растущий lag = consumer не успевает.

Consumer Lag через CLI

# Текущий lag всех consumer groups
docker exec kafka kafka-consumer-groups \
  --bootstrap-server localhost:9092 \
  --list

# Детали конкретной группы
docker exec kafka kafka-consumer-groups \
  --bootstrap-server localhost:9092 \
  --describe \
  --group analytics

# Вывод:
# GROUP     TOPIC          PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CONSUMER-ID
# analytics lesson-events  0          1000            1050            50   analytics-1@host
# analytics lesson-events  1          800             800             0    analytics-2@host
# analytics lesson-events  2          600             650             50   analytics-2@host

Lag 0 - consumer в реальном времени. Lag растёт - нужно масштабировать или оптимизировать.

franz-go Metrics Hook

franz-go предоставляет kgo.Hook для кастомных метрик:

type kafkaMetrics struct {
    producerErrors   prometheus.Counter
    producerBytes    prometheus.Counter
    consumerLag      *prometheus.GaugeVec
    fetchRate        prometheus.Gauge
}

func newKafkaMetrics(reg prometheus.Registerer) *kafkaMetrics {
    m := &kafkaMetrics{
        producerErrors: prometheus.NewCounter(prometheus.CounterOpts{
            Name: "kafka_producer_errors_total",
            Help: "Total Kafka producer errors",
        }),
        producerBytes: prometheus.NewCounter(prometheus.CounterOpts{
            Name: "kafka_producer_bytes_total",
        }),
        consumerLag: prometheus.NewGaugeVec(prometheus.GaugeOpts{
            Name: "kafka_consumer_lag",
            Help: "Consumer lag per partition",
        }, []string{"topic", "partition"}),
        fetchRate: prometheus.NewGauge(prometheus.GaugeOpts{
            Name: "kafka_consumer_fetch_rate",
        }),
    }
    reg.MustRegister(m.producerErrors, m.producerBytes, m.consumerLag, m.fetchRate)
    return m
}

// Реализуем kgo.HookProduceRecordUnbuffered
func (m *kafkaMetrics) OnProduceRecordUnbuffered(r *kgo.Record, err error) {
    if err != nil {
        m.producerErrors.Inc()
    } else {
        m.producerBytes.Add(float64(len(r.Value)))
    }
}

// Реализуем kgo.HookFetchRecordUnbuffered
func (m *kafkaMetrics) OnFetchRecordUnbuffered(r *kgo.Record, _ bool) {
    m.fetchRate.Inc()
}
metrics := newKafkaMetrics(prometheus.DefaultRegisterer)

client, _ := kgo.NewClient(
    kgo.SeedBrokers("localhost:9092"),
    kgo.WithHooks(metrics), // подключаем hook
)

Готовый плагин: kprom

import "github.com/twmb/franz-go/plugin/kprom"

metrics := kprom.NewMetrics("kafka", kprom.GoMetrics())

client, _ := kgo.NewClient(
    kgo.SeedBrokers("localhost:9092"),
    kgo.WithHooks(metrics),
)

kprom автоматически регистрирует стандартный набор метрик с префиксом kafka_:

  • kafka_produce_bytes_total
  • kafka_fetch_bytes_total
  • kafka_produce_records_total
  • kafka_fetch_records_total
  • Latency гистограммы

Consumer Lag в Prometheus

Kafka не экспортирует consumer lag напрямую в Prometheus. Нужен экспортер.

kafka-exporter (простой вариант):

# docker-compose.yml
  kafka-exporter:
    image: danielqsj/kafka-exporter:latest
    command: --kafka.server=kafka:29092
    ports:
      - "9308:9308"
    depends_on:
      - kafka

Метрики доступны на http://localhost:9308/metrics:

# Consumer group lag
kafka_consumergroup_lag{consumergroup="analytics",partition="0",topic="lesson-events"} 50
kafka_consumergroup_lag{consumergroup="analytics",partition="1",topic="lesson-events"} 0

# Топик размер
kafka_topic_partition_current_offset{partition="0",topic="lesson-events"} 1050

Prometheus конфиг:

scrape_configs:
  - job_name: 'kafka-exporter'
    static_configs:
      - targets: ['kafka-exporter:9308']
  - job_name: 'app'
    static_configs:
      - targets: ['app:8080']  # /metrics с kprom

Alerting Rules

# prometheus/alerts.yml
groups:
  - name: kafka
    rules:
      # Lag растёт больше 5 минут подряд
      - alert: KafkaConsumerLagGrowing
        expr: |
          rate(kafka_consumergroup_lag[5m]) > 0
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Consumer lag growing for {{ $labels.consumergroup }}"

      # Lag превысил абсолютный порог (SLA = обработать за 60 секунд)
      - alert: KafkaConsumerLagHigh
        expr: kafka_consumergroup_lag > 10000
        for: 2m
        labels:
          severity: critical
        annotations:
          summary: "High consumer lag: {{ $value }} messages behind"

      # Producer ошибки
      - alert: KafkaProducerErrors
        expr: rate(kafka_producer_errors_total[5m]) > 0
        labels:
          severity: warning

kafka-ui: web-интерфейс

kafka-ui (provectuslabs) - удобная альтернатива CLI для ежедневного мониторинга:

  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    ports:
      - "8090:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092
      KAFKA_CLUSTERS_0_KAFKACONNECT_0_NAME: debezium
      KAFKA_CLUSTERS_0_KAFKACONNECT_0_ADDRESS: http://connect:8083

В kafka-ui видно:

  • Топики: размер, партиции, retention
  • Consumer Groups: lag per partition, last committed offset
  • Messages: просмотр сообщений, поиск по ключу
  • Kafka Connect: статус коннекторов

Метрики JMX (для производственного кластера)

Kafka нативно экспортирует метрики через JMX. JMX Exporter конвертирует их в Prometheus:

  kafka:
    image: confluentinc/cp-kafka:7.6.0
    environment:
      KAFKA_JMX_PORT: 9999
      KAFKA_JMX_HOSTNAME: kafka
      EXTRA_ARGS: -javaagent:/opt/jmx-exporter.jar=7071:/etc/jmx-exporter.yml

JMX даёт детальные метрики брокера: ISR shrink rate, under-replicated partitions, request rate per type.

Типичная ошибка

Не мониторить lag - consumer отстаёт неделями, а команда узнаёт когда диск у Kafka заканчивается (retention не чистит, потому что slow consumer держит старые сегменты). Алерт на lag обязателен.

Не мониторить replication slot lag (CDC) - если используешь Debezium и slot лагает, PostgreSQL не может вычистить WAL. Диск переполняется за часы при высокой нагрузке записи.

Смотри также

Мини-задание

  • Добавь kafka-exporter и kafka-ui в docker-compose, проверь метрики на /metrics
  • Отправь 1000 событий через producer, запусти медленный consumer (sleep 10ms per record) - наблюдай как растёт lag в kafka-ui
  • Подключи kprom к своему producer, добавь /metrics endpoint, проверь метрики в Prometheus
  • Напиши alert rule на lag > 1000 и проверь что Prometheus его поднимает при имитации отставания
  • Найди через kafka-consumer-groups --describe какой partition имеет наибольший lag

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.