학습 목표

시나리오

결제 정산 컨슈머가 새벽에 멈췄습니다. 아침에 발견했을 때 lag은 210만 건이었고 복구에 4시간이 걸렸습니다.

원인 조사에서 드러난 것은 알림이 있었지만 울리지 않았다는 사실입니다. 알림 조건이 consumer_lag > 100000이었고, 컨슈머 프로세스가 죽으면서 클라이언트 메트릭 자체가 사라져 조건이 평가되지 않았습니다. "값이 크다"는 조건은 값이 없어지는 장애를 잡지 못합니다.

요구사항은 세 가지입니다. (1) 컨슈머가 사라지면 알린다. (2) lag이 크지 않아도 줄지 않으면 알린다. (3) 파티션별 편차를 본다(한 파티션만 막히는 패턴을 잡기 위해).

lag 측정 3가지 방법

측정 주체와 함정
방법 측정 주체 얻는 것 함정
kafka-consumer-groups.sh --describe AdminClient (외부) 파티션별 CURRENT-OFFSET / LOG-END-OFFSET / LAG 커밋된 오프셋 기준입니다. 처리는 진행 중인데 커밋 주기가 길면 lag이 실제보다 커 보입니다. 그리고 폴링이라 부하가 있습니다
컨슈머 JMX records-lag-max 컨슈머 클라이언트 이 컨슈머가 본 파티션 중 최대 lag 컨슈머가 죽으면 메트릭도 사라집니다. 그리고 "최대"이므로 어느 파티션인지 알 수 없고, 할당받지 않은 파티션은 보이지 않습니다
외부 exporter (AdminClient 기반) 전용 프로세스 그룹 × 토픽 × 파티션 lag을 시계열로 컨슈머가 죽어도 계속 측정됩니다(가장 중요한 성질). 대신 컴포넌트가 하나 늘고, 커밋된 오프셋 기준이라는 한계는 같습니다

아키텍처

핵심 메트릭 대시보드 — URP, UnderMinIsr, OfflinePartitions, ActiveController, 요청 처리 유휴율, 컨슈머 lag 여섯 개의 지표를 정상 기준과 함께 배치했습니다. UnderReplicatedPartitions 는 kafka.server type ReplicaManager 에 있고 정상값은 0 이며 0 이 아니면 ISR 이 전체 레플리카보다 적어 복제가 밀리는 상태입니다. UnderMinIsrPartitionCount 도 같은 빈에 있고 정상값 0 이며 min.insync.replicas 미달이므로 쓰기가 거부되기 시작합니다. OfflinePartitionsCount 는 kafka.controller type KafkaController 에 있고 정상값 0 이며 리더가 없는 파티션이 있으면 읽기와 쓰기가 모두 불가능합니다. ActiveControllerCount 는 클러스터 전체 합이 정확히 1 이어야 하고 0 이면 컨트롤러가 없고 2 이상이면 분열 상태입니다. RequestHandlerAvgIdlePercent 는 0 과 1 사이 값으로 0.3 이상을 권하며 0 에 가까우면 요청 처리 스레드가 포화된 것입니다. records-lag-max 는 컨슈머 쪽 지표로 얼마나 뒤처졌는지를 나타냅니다. IsrShrinksPerSec 는 UnderReplicated 보다 먼저 움직이는 조기 신호입니다. 핵심 메트릭 대시보드 — 이 여섯 개만 보면 클러스터 상태를 압니다 값은 정상 기준입니다. 하나라도 벗어나면 그 줄의 "의미"부터 확인하세요. UnderReplicatedPartitions kafka.server:type=ReplicaManager 정상: 0 ISR < 전체 레플리카 — 복제가 밀린다 UnderMinIsrPartitionCount kafka.server:type=ReplicaManager 정상: 0 ISR < min.insync.replicas — 쓰기 거부 시작 OfflinePartitionsCount kafka.controller:type=KafkaController 정상: 0 리더 없는 파티션 — 읽기·쓰기 모두 불가 ActiveControllerCount kafka.controller:type=KafkaController 정상: 클러스터 합 1 0 이면 컨트롤러 없음 · 2 이상이면 split brain RequestHandlerAvgIdlePercent kafka.server:type=KafkaRequestHandlerPool 정상: 0.3 이상 0 에 가까우면 요청 처리 스레드 포화 records-lag-max consumer-fetch-manager-metrics (클라이언트) 정상: 업무 기준 이내 컨슈머가 얼마나 뒤처졌는지 (클라이언트 지표) IsrShrinksPerSec 가 오르면 브로커 한 대가 뒤처지는 신호입니다 — UnderReplicated 보다 먼저 움직입니다. lag 은 브로커가 아니라 컨슈머 쪽 지표입니다 — CLI 로는 kafka-consumer-groups.sh --describe 로 봅니다. records-lag-max 의 전체 이름은 kafka.consumer:type=consumer-fetch-manager-metrics 입니다.
핵심 메트릭 대시보드 — URP / OfflinePartitions / ActiveController / RequestHandlerIdle 배치 예시

브로커 JVM에 JMX Exporter를 javaagent로 붙여 JMX MBean을 Prometheus 형식으로 노출합니다. 예제 1의 compose에 환경변수 두 개를 추가하면 됩니다.

구성 요소
컴포넌트버전역할
JMX Exporter1.0.1브로커 JVM에 javaagent로 붙어 :9404/metrics로 노출합니다
Prometheusv3.13.1스크레이프와 알림 룰 평가
Grafana13.1.1대시보드
Kafka4.3.1예제 1의 3노드 클러스터

전체 코드

디렉터리 구조
kafka-monitoring/
├── docker-compose.monitoring.yml
├── jmx/
│   ├── download-agent.sh        # JMX Exporter jar 내려받기
│   └── kafka-broker.yml         # 노출할 MBean 화이트리스트
├── prometheus/
│   ├── prometheus.yml
│   └── rules/kafka-alerts.yml
└── grafana/
    └── provisioning/datasources/prometheus.yml

JMX Exporter 준비

kafka-monitoring/jmx/download-agent.sh
#!/usr/bin/env bash
# JMX Exporter javaagent 를 내려받습니다. 브로커 컨테이너에 마운트해 씁니다.
set -euo pipefail

VERSION="${JMX_EXPORTER_VERSION:-1.0.1}"
DEST="$(cd "$(dirname "$0")" && pwd)"
JAR="${DEST}/jmx_prometheus_javaagent-${VERSION}.jar"

if [[ -f "$JAR" ]]; then
  echo "이미 있습니다: $JAR"
  exit 0
fi

curl -fsSL -o "$JAR" \
  "https://repo1.maven.org/maven2/io/prometheus/jmx/jmx_prometheus_javaagent/${VERSION}/jmx_prometheus_javaagent-${VERSION}.jar"

echo "내려받음: $JAR"
kafka-monitoring/jmx/kafka-broker.yml — 노출할 MBean만 고릅니다
# JMX Exporter 설정.
#
# whitelistObjectNames 로 스캔 대상을 좁히는 것이 중요합니다.
# Kafka 브로커는 "토픽 × 파티션" 단위 MBean 을 수천 개 만들 수 있고,
# 전부 스캔하면 스크레이프 한 번이 수 초씩 걸려 브로커에 부하가 갑니다.
lowercaseOutputName: true
lowercaseOutputLabelNames: true

whitelistObjectNames:
  - 'kafka.controller:type=KafkaController,name=*'
  - 'kafka.server:type=ReplicaManager,name=*'
  - 'kafka.server:type=BrokerTopicMetrics,name=*'
  - 'kafka.server:type=KafkaRequestHandlerPool,name=*'
  - 'kafka.network:type=SocketServer,name=*'
  - 'kafka.network:type=RequestMetrics,name=*,request=*'
  - 'kafka.log:type=LogFlushStats,name=*'
  - 'java.lang:type=Memory'
  - 'java.lang:type=GarbageCollector,name=*'

rules:
  # --- 컨트롤러 (KRaft) ---------------------------------------------------
  # ActiveControllerCount: 클러스터 전체 합이 정확히 1이어야 합니다.
  # 0이면 컨트롤러가 없고, 2 이상이면 split-brain 입니다.
  - pattern: 'kafka.controller<>Value'
    name: kafka_controller_$1
    type: GAUGE

  # --- 복제 -------------------------------------------------------------
  # UnderReplicatedPartitions: 정상은 0. 0이 아니면 복제가 밀리는 중입니다.
  - pattern: 'kafka.server<>Value'
    name: kafka_server_replicamanager_$1
    type: GAUGE

  # IsrShrinksPerSec / IsrExpandsPerSec: 반복적으로 움직이면 브로커가 불안정합니다.
  - pattern: 'kafka.server<>Count'
    name: kafka_server_replicamanager_$1_total
    type: COUNTER

  # --- 처리량 -----------------------------------------------------------
  # 토픽 라벨이 없는(=브로커 전체) 지표만 받습니다.
  - pattern: 'kafka.server<>Count'
    name: kafka_server_brokertopicmetrics_$1_total
    type: COUNTER

  # --- 요청 처리 여유 ----------------------------------------------------
  # RequestHandlerAvgIdlePercent: 0에 가까우면 io 스레드가 포화 상태입니다.
  # 0.3 미만이 지속되면 num.io.threads 를 검토합니다.
  - pattern: 'kafka.server<>OneMinuteRate'
    name: kafka_server_requesthandler_avg_idle_percent
    type: GAUGE

  # NetworkProcessorAvgIdlePercent: 네트워크 스레드 여유.
  - pattern: 'kafka.network<>Value'
    name: kafka_network_processor_avg_idle_percent
    type: GAUGE

  # --- 요청 지연 (요청 타입별 p99) ----------------------------------------
  - pattern: 'kafka.network<>99thPercentile'
    name: kafka_network_request_total_time_ms_p99
    labels:
      request: '$1'
    type: GAUGE

  # --- JVM ---------------------------------------------------------------
  - pattern: 'java.lang(\w+)'
    name: jvm_memory_heap_$1_bytes
    type: GAUGE
  - pattern: 'java.lang<>CollectionTime'
    name: jvm_gc_collection_time_ms_total
    labels:
      gc: '$1'
    type: COUNTER

브로커에 javaagent 붙이기

예제 1의 docker-compose.yml세 줄을 추가하면 됩니다. 아래는 override 파일로 분리한 형태입니다 — 원본을 고치지 않아도 됩니다.

kafka-lab/docker-compose.override.yml (예제 1 디렉터리에 둡니다)
# 예제 1의 compose 를 고치지 않고 JMX Exporter 만 얹습니다.
# docker compose 는 같은 디렉터리의 docker-compose.override.yml 을 자동으로 병합합니다.
---
x-jmx: &jmx
  volumes:
    # download-agent.sh 로 내려받은 jar 와 설정을 마운트합니다.
    - ../kafka-monitoring/jmx/jmx_prometheus_javaagent-1.0.1.jar:/opt/jmx/agent.jar:ro
    - ../kafka-monitoring/jmx/kafka-broker.yml:/opt/jmx/config.yml:ro
  environment:
    # KAFKA_OPTS 는 브로커 기동 스크립트가 JVM 인자로 넘겨 줍니다.
    # 형식: -javaagent:{jar}={포트}:{설정파일}
    KAFKA_OPTS: '-javaagent:/opt/jmx/agent.jar=9404:/opt/jmx/config.yml'

services:
  kafka-1:
    <<: *jmx
  kafka-2:
    <<: *jmx
  kafka-3:
    <<: *jmx

Prometheus 설정

kafka-monitoring/prometheus/prometheus.yml
global:
  # 15초 간격이면 1분 단위 증가율을 안정적으로 계산할 수 있습니다.
  scrape_interval: 15s
  evaluation_interval: 15s

rule_files:
  - /etc/prometheus/rules/*.yml

scrape_configs:
  - job_name: kafka-broker
    static_configs:
      - targets:
          - 'kafka-1:9404'
          - 'kafka-2:9404'
          - 'kafka-3:9404'
    relabel_configs:
      # instance 라벨을 브로커 호스트명으로 정리합니다.
      - source_labels: [__address__]
        regex: '([^:]+):.*'
        target_label: broker
        replacement: '$1'

  # 애플리케이션(컨슈머/프로듀서)이 Micrometer 등으로 노출하는 메트릭.
  # Spring Boot 라면 /actuator/prometheus 입니다(예제 2 · 예제 6).
  - job_name: kafka-clients
    metrics_path: /actuator/prometheus
    static_configs:
      - targets: ['host.docker.internal:8080']
    # 로컬에서 애플리케이션이 안 떠 있으면 스크레이프가 실패합니다. 정상입니다.

알림 룰 — "값이 크다"가 아니라 "줄지 않는다"

kafka-monitoring/prometheus/rules/kafka-alerts.yml
groups:
  # ======================================================================
  # 클러스터 상태 — 여기가 울리면 데이터 유실 위험이 있습니다
  # ======================================================================
  - name: kafka-cluster
    rules:
      # 컨트롤러가 정확히 1개여야 합니다.
      # 0이면 메타데이터 변경이 불가능하고, 2 이상이면 split-brain 입니다.
      - alert: KafkaNoActiveController
        expr: sum(kafka_controller_ActiveControllerCount) != 1
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: '활성 컨트롤러 수가 1이 아닙니다 (현재 {{ $value }})'
          description: 'KRaft 컨트롤러 쿼럼을 확인하세요: kafka-metadata-quorum.sh describe --status'

      # 오프라인 파티션 = 리더가 없는 파티션. 읽기도 쓰기도 불가능합니다.
      - alert: KafkaOfflinePartitions
        expr: sum(kafka_controller_OfflinePartitionsCount) > 0
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: '리더가 없는 파티션이 {{ $value }}개 있습니다'

      # URP > 0 이 지속되면 다음 브로커 장애에서 유실이 발생할 수 있습니다.
      # 롤링 재시작 중에는 일시적으로 올라가므로 for 를 10분으로 둡니다.
      - alert: KafkaUnderReplicatedPartitions
        expr: sum(kafka_server_replicamanager_UnderReplicatedPartitions) > 0
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: '복제가 밀린 파티션이 {{ $value }}개입니다'
          description: 'kafka-topics.sh --describe --under-replicated-partitions 로 확인하세요'

      # min.insync.replicas 에 미달한 파티션. acks=all 쓰기가 실패합니다.
      - alert: KafkaUnderMinIsr
        expr: sum(kafka_server_replicamanager_UnderMinIsrPartitionCount) > 0
        for: 2m
        labels:
          severity: critical
        annotations:
          summary: 'min.insync.replicas 미달 파티션 {{ $value }}개 — acks=all 쓰기가 실패합니다'

      # ISR 축소가 반복되면 브로커나 네트워크가 불안정합니다.
      - alert: KafkaIsrChurn
        expr: sum(rate(kafka_server_replicamanager_IsrShrinksPerSec_total[10m])) > 0
        for: 15m
        labels:
          severity: warning
        annotations:
          summary: 'ISR 축소가 반복되고 있습니다 (GC · 디스크 · 네트워크 확인)'

      # 요청 핸들러 여유가 없으면 처리 지연이 시작됩니다.
      - alert: KafkaRequestHandlerSaturated
        expr: avg by (broker) (kafka_server_requesthandler_avg_idle_percent) < 0.3
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: '{{ $labels.broker }} 의 요청 핸들러 유휴율이 {{ $value | humanizePercentage }}입니다'

  # ======================================================================
  # 컨슈머 lag — 시나리오의 사고를 막는 룰들
  # ======================================================================
  - name: kafka-consumer-lag
    rules:
      # (1) 컨슈머가 사라진 것을 감지합니다.
      #     absent_over_time 은 "메트릭이 없어진 것" 자체를 조건으로 씁니다.
      #     시나리오의 사고가 여기서 잡힙니다.
      - alert: KafkaConsumerDisappeared
        expr: absent_over_time(kafka_consumer_fetch_manager_records_lag_max{client_id!=""}[10m])
        for: 5m
        labels:
          severity: critical
        annotations:
          summary: '컨슈머 메트릭이 10분간 보이지 않습니다 — 프로세스가 죽었을 수 있습니다'

      # (2) lag 이 크지 않아도 "줄지 않으면" 알립니다.
      #     deriv 가 0 이상이면 lag 이 감소하지 않고 있다는 뜻입니다.
      #     절대값 조건만 두면 천천히 쌓이는 장애를 놓칩니다.
      - alert: KafkaConsumerLagNotDecreasing
        expr: |
          kafka_consumer_fetch_manager_records_lag_max > 1000
          and
          deriv(kafka_consumer_fetch_manager_records_lag_max[15m]) >= 0
        for: 15m
        labels:
          severity: warning
        annotations:
          summary: '{{ $labels.client_id }} 의 lag 이 15분간 줄지 않습니다 (현재 {{ $value }})'
          description: '처리 실패 반복 · 외부 시스템 지연 · 리밸런스 루프를 확인하세요'

      # (3) 절대값 조건도 함께 둡니다(급격한 증가를 빠르게 잡기 위해).
      - alert: KafkaConsumerLagHigh
        expr: kafka_consumer_fetch_manager_records_lag_max > 100000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: '{{ $labels.client_id }} 의 lag 이 {{ $value }} 입니다'

      # (4) 소비가 완전히 멈춘 것을 감지합니다.
      #     lag 이 있는데 records-consumed-rate 가 0 이면 처리가 멈춘 것입니다.
      - alert: KafkaConsumerStalled
        expr: |
          kafka_consumer_fetch_manager_records_lag_max > 0
          and
          rate(kafka_consumer_fetch_manager_records_consumed_total[5m]) == 0
        for: 5m
        labels:
          severity: critical
        annotations:
          summary: '{{ $labels.client_id }} 가 lag 을 두고 소비를 멈췄습니다'

  # ======================================================================
  # 프로듀서
  # ======================================================================
  - name: kafka-producer
    rules:
      # 전송 에러는 0이 정상입니다. 1건이라도 나면 유실 가능성을 검토해야 합니다.
      - alert: KafkaProducerErrors
        expr: rate(kafka_producer_record_error_total[5m]) > 0
        for: 2m
        labels:
          severity: critical
        annotations:
          summary: '{{ $labels.client_id }} 에서 전송 실패가 발생했습니다'
          description: 'min.insync.replicas 미달 · 레코드 크기 초과 · 권한 문제를 확인하세요'

      # 재시도가 늘면 브로커가 불안정하거나 ISR 이 부족합니다.
      - alert: KafkaProducerRetries
        expr: rate(kafka_producer_record_retry_total[5m]) > 1
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: '{{ $labels.client_id }} 의 재시도율이 높습니다'

Compose

kafka-monitoring/docker-compose.monitoring.yml
---
name: kafka-lab-monitoring

services:
  prometheus:
    image: prom/prometheus:v3.13.1
    container_name: prometheus
    ports:
      - '9090:9090'
    command:
      - '--config.file=/etc/prometheus/prometheus.yml'
      - '--storage.tsdb.path=/prometheus'
      # 로컬 실습이므로 보관 기간을 짧게 둡니다.
      - '--storage.tsdb.retention.time=7d'
      # 룰 파일을 고친 뒤 재시작 없이 reload 할 수 있게 합니다:
      #   curl -X POST http://localhost:9090/-/reload
      - '--web.enable-lifecycle'
    volumes:
      - ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml:ro
      - ./prometheus/rules:/etc/prometheus/rules:ro
      - prometheus-data:/prometheus
    extra_hosts:
      # 호스트에서 돌리는 애플리케이션을 스크레이프하기 위해 필요합니다(Linux).
      - 'host.docker.internal:host-gateway'
    networks:
      - kafka-lab_default

  grafana:
    image: grafana/grafana:13.1.1
    container_name: grafana
    ports:
      - '3000:3000'
    environment:
      GF_SECURITY_ADMIN_USER: admin
      GF_SECURITY_ADMIN_PASSWORD: admin
      # 로컬 실습 편의. 프로덕션에서는 절대 켜지 마세요.
      GF_AUTH_ANONYMOUS_ENABLED: 'true'
      GF_AUTH_ANONYMOUS_ORG_ROLE: Viewer
    volumes:
      - ./grafana/provisioning:/etc/grafana/provisioning:ro
      - grafana-data:/var/lib/grafana
    depends_on:
      - prometheus
    networks:
      - kafka-lab_default

volumes:
  prometheus-data:
  grafana-data:

networks:
  kafka-lab_default:
    external: true
kafka-monitoring/grafana/provisioning/datasources/prometheus.yml
apiVersion: 1

datasources:
  - name: Prometheus
    type: prometheus
    access: proxy
    # 컨테이너 네트워크 안에서의 주소입니다.
    url: http://prometheus:9090
    isDefault: true
    editable: false

실행 방법

순서대로 실행
# 1. JMX Exporter 내려받기
cd kafka-monitoring
chmod +x jmx/download-agent.sh && ./jmx/download-agent.sh

# 2. 예제 1의 compose 에 javaagent 설정을 넣고 브로커를 재기동합니다.
#    (위 "브로커에 javaagent 붙이기" 절 참고)
cd ../kafka-lab
docker compose up -d --force-recreate

# 3. 브로커가 메트릭을 노출하는지 확인
curl -s http://localhost:29092 >/dev/null 2>&1 || true
docker exec kafka-1 curl -s http://localhost:9404/metrics | grep -c '^kafka_'
# → 0 보다 큰 수가 나와야 합니다.

# 4. Prometheus + Grafana 기동
cd ../kafka-monitoring
docker compose -f docker-compose.monitoring.yml up -d

# 5. 확인
open http://localhost:9090/targets   # 브로커 3개가 UP 이어야 합니다
open http://localhost:9090/alerts    # 룰이 로드되었는지
open http://localhost:3000           # Grafana (admin/admin)

검증 방법

1. 필수 메트릭 7종이 보이는가

브로커에서 반드시 봐야 할 메트릭 (JMX MBean 이름 기준)
MBean정상 범위벗어나면
kafka.controller:type=KafkaController,
name=ActiveControllerCount
클러스터 전체 합계 정확히 1 0이면 메타데이터 변경 불가, 2 이상이면 split-brain
kafka.controller:type=KafkaController,
name=OfflinePartitionsCount
0 리더 없는 파티션 — 읽기·쓰기 모두 불가
kafka.server:type=ReplicaManager,
name=UnderReplicatedPartitions
0 복제가 밀리는 중. 다음 장애에서 유실 위험
kafka.server:type=ReplicaManager,
name=UnderMinIsrPartitionCount
0 acks=all 쓰기가 실패하는 중 (예제 3)
kafka.server:type=ReplicaManager,
name=IsrShrinksPerSec
평상시 0 반복되면 브로커 GC·디스크·네트워크 문제
kafka.server:type=KafkaRequestHandlerPool,
name=RequestHandlerAvgIdlePercent
0.3 이상 io 스레드 포화. num.io.threads(기본 8) 검토
kafka.network:type=SocketServer,
name=NetworkProcessorAvgIdlePercent
0.3 이상 네트워크 스레드 포화. num.network.threads(기본 3) 검토
Prometheus에서 직접 확인
# 컨트롤러가 정확히 1개인지
curl -s 'http://localhost:9090/api/v1/query?query=sum(kafka_controller_ActiveControllerCount)' | jq '.data.result'

# URP 합계
curl -s 'http://localhost:9090/api/v1/query?query=sum(kafka_server_replicamanager_UnderReplicatedPartitions)' | jq '.data.result'

2. lag 3가지 측정값을 비교한다

같은 시점에 세 방법으로 재 봅니다
cd ../kafka-lab

# 컨슈머를 느리게 만들어 lag 을 만듭니다.
# (예제 4의 컨슈머에서 max.poll.records 를 줄이거나 처리에 sleep 을 넣습니다)

echo '--- (A) AdminClient 기준 ---'
./kcli kafka-consumer-groups.sh --describe --group orders-warehouse-loader

echo '--- (B) 컨슈머 JMX 기준 ---'
curl -s 'http://localhost:9090/api/v1/query?query=kafka_consumer_fetch_manager_records_lag_max' | jq -r \
  '.data.result[] | "\(.metric.client_id): \(.value[1])"'

echo '--- (C) 파티션별 log-end-offset 과 커밋 오프셋의 차 ---'
./kcli kafka-get-offsets.sh --topic orders --time -1

(A)와 (B)가 다른 것이 정상입니다. (A)는 커밋된 오프셋 기준이고 (B)는 컨슈머가 실제로 fetch한 위치 기준입니다. 커밋 주기가 길면 (A)가 (B)보다 크게 나옵니다. 둘의 차이가 곧 "커밋되지 않은 처리량"이며, 장애 시 재처리될 구간의 크기입니다.

3. 컨슈머가 사라지는 알림이 실제로 울리는가

시나리오의 사고를 재현
# 컨슈머를 강제 종료합니다.
pkill -9 -f 'com.example.kafka.consumer.Main'

# 10분 뒤 KafkaConsumerDisappeared 가 pending → firing 이 됩니다.
# 대기하지 않고 확인하려면 룰의 [10m] 을 [1m] 로 줄이고 reload 하세요.
sed -i 's/\[10m\])$/[1m])/' prometheus/rules/kafka-alerts.yml
curl -s -X POST http://localhost:9090/-/reload

sleep 90
curl -s http://localhost:9090/api/v1/alerts | jq -r \
  '.data.alerts[] | "\(.labels.alertname) \(.state)"'
기대 출력
KafkaConsumerDisappeared firing

절대값 조건(KafkaConsumerLagHigh)은 울리지 않습니다 — 메트릭이 없어졌기 때문입니다. absent_over_time이 그 빈틈을 메웁니다. 이것이 시나리오의 4시간 장애를 막는 한 줄입니다.

4. 파티션별 편차를 본다

records-lag-max최대값만 주므로 "어느 파티션이 막혔는지"를 알 수 없습니다. 파티션별 lag은 AdminClient 기반으로만 얻을 수 있습니다.

파티션별 lag을 표로
./kcli kafka-consumer-groups.sh --describe --group payment-settlement \
  | awk 'NR==1 || $6 != "0"' 
한 파티션만 막힌 패턴 — 예제 6의 증상
GROUP                TOPIC     PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
payment-settlement   payments  1          40213           331882          291669

파티션 0과 2는 lag 0인데 1만 29만이면 브로커나 컨슈머 용량 문제가 아니라 파티션 1의 특정 레코드에서 처리가 막힌 것입니다. 이 패턴을 알아보는 것이 조사 시간을 크게 줄입니다 (케이스 2).

프로덕션 고려사항

로컬 예제와 프로덕션의 차이
항목이 예제프로덕션
파티션별 lag CLI로 수동 확인 AdminClient 기반 lag exporter를 별도 프로세스로 띄웁니다. 컨슈머가 죽어도 측정이 계속되는 것이 핵심 가치입니다
MBean 스캔 범위 화이트리스트 9개 토픽·파티션 단위 MBean까지 켜면 수만 개 시계열이 생겨 Prometheus와 브로커 양쪽에 부하가 갑니다. 반드시 좁히세요
알림 라우팅 Prometheus 내부에서만 Alertmanager로 심각도별 채널(PagerDuty / Slack)을 분리하고 억제 규칙을 둡니다. URP 알림이 100개 오면 아무도 안 봅니다
보존 기간 7일 장애 사후분석에 필요한 기간(보통 30~90일) + 원격 저장소. 사고 회고에서 "그때 메트릭이 없다"가 가장 흔한 문제입니다
Grafana 인증 anonymous 허용 절대 금지. SSO 연동 + 편집 권한 분리
lag의 단위 건수 건수보다 시간 기반 lag(가장 오래된 미처리 레코드의 나이)이 SLO에 맞습니다. 건수는 메시지 크기와 처리 속도에 따라 의미가 달라집니다
클라이언트 메트릭 Spring actuator 가정 Java 클라이언트는 metrics()로 JMX에 노출하므로 애플리케이션에도 JMX Exporter를 붙이거나 Micrometer 브리지를 씁니다
알림 임계값 고정값 서비스별 SLO에서 역산합니다. "lag 10만"이 아니라 "복구 목표 시간 안에 소화 가능한 양"이 기준입니다

자주 하는 실수

이어서 볼 곳

공식 문서 출처