Алертинг Kafka в Kubernetes: consumer lag, партиции и диск

Published: 2026-02-18

Для Kafka в продакшене нужны четыре категории алертов: доступность брокеров, состояние репликации, consumer lag и насыщение ресурсов. У каждой категории своя срочность и свой источник метрик. Вот набор PrometheusRule, который мы используем в Kubernetes вместе с kafka-exporter.


Деплой kafka-exporter

kafka-exporter собирает метрики через Kafka API (AdminClient) и отдаёт их в формате Prometheus. JMX он не использует, поэтому легче других экспортёров Kafka. Разворачиваем Helm-чартом prometheus-kafka-exporter:

yamlapiVersion: helm.toolkit.fluxcd.io/v2
kind: HelmRelease
metadata:
  name: kafka-exporter
spec:
  chart:
    spec:
      chart: prometheus-kafka-exporter
      version: "3.0.1"
      sourceRef:
        kind: HelmRepository
        name: prometheus-community
        namespace: flux-system
  values:
    kafkaServer:
      - kafka.kafka.svc.cluster.local:9092
    prometheus:
      serviceMonitor:
        enabled: true
        additionalLabels:
          release: kube-prom-stack
        interval: 30s
        relabelings:
          - targetLabel: job
            replacement: kafka
    resources:
      requests:
        cpu: 25m
        memory: 64Mi
      limits:
        cpu: 200m
        memory: 128Mi

Блок relabelings переопределяет метку job на kafka — все наши выражения алертов используют job="kafka". Без этого меткой job по умолчанию станет имя HelmRelease или ServiceMonitor, а оно в разных средах разное.


Доступность брокеров

yaml- alert: KafkaDown
  expr: up{job="kafka"} < 1
  for: 5m
  labels:
    severity: disaster
    service: kafka
    groupp: admin
    url: 'https://grafana.example.com/d/kafka-overview/kafka-overview'
    event: 'Kafka Down |{{$labels.instance}}'
  annotations:
    description: >-
      Kafka broker {{$labels.instance}} недоступен более 5 мин.
      Timestamp: {{ with query "time()+10800" }}{{ . | first | value | humanizeTimestamp }}{{ end }}
yaml- alert: KafkaExporterMissing
  expr: absent(kafka_brokers{job="kafka"})
  for: 5m
  labels:
    severity: high
    service: kafka
    event: 'Kafka Exporter Missing'
  annotations:
    description: Серия kafka_brokers отсутствует — экспортёр может быть недоступен или не может подключиться к Kafka.

Алерт с absent() ловит случай, когда сам экспортёр не может собрать метрики с Kafka. Без него KafkaDown никогда не сработает: вместо значения 0 у Prometheus просто не будет данных.

Мониторинг количества брокеров

Для кластеров с несколькими брокерами:

yaml- alert: KafkaBrokerCountChanged
  expr: kafka_brokers{job="kafka"} < 3
  for: 5m
  labels:
    severity: high
    event: 'Kafka Broker Count Low |{{$value}} brokers'
  annotations:
    description: Ожидается 3 Kafka брокера, видно только {{$value}}. Возможно, брокер упал.

Недореплицированные партиции

yaml- alert: KafkaUnderReplicatedPartitions
  expr: sum(kafka_topic_partition_under_replicated_partition{job="kafka"}) > 0
  for: 10m
  labels:
    severity: high
    service: kafka
    event: 'Kafka Under-Replicated Partitions'
  annotations:
    description: >-
      Kafka имеет {{$value}} под-реплицированных партиций более 10 мин.
      Реплика отстаёт или брокер замедлился. Это ведущий индикатор риска потери данных.

Недореплицированные партиции — опережающий индикатор: они появляются раньше, чем брокер реально упадёт. Если брокер откажет, пока партиции недореплицированы, у части сообщений останется единственная копия — без избыточности.

Почему 10m, а не 5m? Кратковременная недорепликация случается при rolling update и выборах лидера. for: 10m отфильтровывает этот штатный шум.

Offline-партиции

Серьёзнее, чем недореплицированные:

yaml- alert: KafkaOfflinePartitions
  expr: sum(kafka_topic_partition_leader{job="kafka"} < 0) > 0
  for: 1m
  labels:
    severity: disaster
    event: 'Kafka Offline Partitions |{{$value}} partitions'
  annotations:
    description: >-
      {{$value}} партиций Kafka не имеют лидера. Producers и consumers для этих
      партиций будут немедленно падать с LEADER_NOT_AVAILABLE.

Consumer lag

yaml- alert: KafkaConsumerLagCritical
  expr: >
    sum(kafka_consumergroup_lag_sum{job="kafka"}) by (consumergroup) > 100000
  for: 5m
  labels:
    severity: average
    service: kafka
    event: 'Kafka Consumer Lag |{{$labels.consumergroup}}'
  annotations:
    description: >-
      Consumer group {{$labels.consumergroup}}: лаг {{$value}} сообщений
      (выше порога 100 000) более 5 мин.

Порог 100000 зависит от нагрузки:

  • консьюмер обрабатывает 10k сообщений/сек → лаг 100k = отставание на 10 секунд (приемлемо)
  • консьюмер обрабатывает 100 сообщений/сек → тот же лаг = 1000 секунд (критично)

Для отдельных consumer groups задавайте свои пороги:

yaml- alert: KafkaHighPriorityConsumerLag
  expr: >
    kafka_consumergroup_lag_sum{job="kafka",
      consumergroup="payment-processor"} > 500
  for: 2m
  labels:
    severity: high
    event: 'Payment Consumer Lag |{{$labels.consumergroup}}'

Consumer group перестала читать

yaml- alert: KafkaConsumerGroupStopped
  expr: >
    delta(kafka_consumergroup_lag_sum{job="kafka"}[5m]) > 1000
    and
    kafka_consumergroup_lag_sum{job="kafka"} > 1000
  for: 5m
  labels:
    severity: high
    event: 'Kafka Consumer Stopped |{{$labels.consumergroup}}'
  annotations:
    description: >-
      Consumer group {{$labels.consumergroup}}: лаг вырос более чем на 1000
      сообщений за 5 минут и не уменьшается.

Использование диска (Kafka PVC)

Метрик диска у kafka-exporter нет — они берутся из kubelet:

yaml- alert: KafkaHighDiskUsage
  expr: >
    sum(kubelet_volume_stats_used_bytes{namespace="kafka"})
      by (namespace, persistentvolumeclaim)
    /
    sum(kubelet_volume_stats_capacity_bytes{namespace="kafka"})
      by (namespace, persistentvolumeclaim)
    > 0.9
  for: 15m
  labels:
    severity: warning
    service: kafka
    event: 'Kafka High Disk |{{$labels.persistentvolumeclaim}}'
  annotations:
    description: >-
      Kafka PVC {{$labels.persistentvolumeclaim}} заполнен более чем на 90%.
      Рассмотрите уменьшение retention или добавление хранилища.

Быстрое решение — уменьшить retention.ms на проблемных топиках:

bashkafka-configs.sh --bootstrap-server kafka:9092 \
  --entity-type topics \
  --entity-name my-topic \
  --alter --add-config retention.ms=3600000  # 1 час

Использование памяти

yaml- alert: KafkaHighMemoryUsage
  expr: process_resident_memory_bytes{job="kafka"} / 1024 / 1024 / 1024 > 4
  for: 15m
  labels:
    severity: warning
    event: 'Kafka High Memory |{{$labels.instance}}'
  annotations:
    description: Kafka broker {{$labels.instance}} использует более 4 ГБ RSS более 15 мин.

Kafka работает на JVM. RSS выше лимита heap обычно означает page cache, и это не обязательно проблема. Но превышение лимита памяти контейнера приводит к OOM kill, поэтому задавайте размер heap явно:

yamlvalues:
  extraEnvVars:
    - name: KAFKA_HEAP_OPTS
      value: "-Xmx2g -Xms2g"
  resources:
    limits:
      memory: 4Gi  # heap (2g) + off-heap + накладные расходы