Мониторинг 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_totalkafka_fetch_bytes_totalkafka_produce_records_totalkafka_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. Диск переполняется за часы при высокой нагрузке записи.
Смотри также
- Observability - Метрики с Prometheus - как устроен Prometheus scraping
- Message Brokers - Consumer - как считать lag программно
- Docker - docker-compose - kafka-exporter и kafka-ui в compose
Мини-задание
- Добавь
kafka-exporterиkafka-uiв docker-compose, проверь метрики на/metrics - Отправь 1000 событий через producer, запусти медленный consumer (sleep 10ms per record) - наблюдай как растёт lag в kafka-ui
- Подключи
kpromк своему producer, добавь/metricsendpoint, проверь метрики в Prometheus - Напиши alert rule на lag > 1000 и проверь что Prometheus его поднимает при имитации отставания
- Найди через
kafka-consumer-groups --describeкакой partition имеет наибольший lag