Step 12 — 관측성과 운영

학습 목표

  • Kafka 애플리케이션에서 반드시 봐야 할 4대 신호(컨슈머 랙 / 처리 지연 / 에러·DLT / 리밸런스)를 구분해 정의한다
  • Boot 자동 설정이 등록하는 MicrometerConsumerListener / MicrometerProducerListener 가 노출하는 지표를 /actuator/metrics 로 직접 확인한다
  • 클라이언트가 보고하는 랙이 컨슈머와 함께 사라지는 것을 재현하고, AdminClient 기반 랙 게이지를 직접 만들어 대체한다
  • observation-enabled 를 프로듀서·컨슈머 양쪽에 켜서 traceparent 헤더로 trace 가 이어지는 것을 로그와 헤더 덤프로 확인한다
  • 모든 리스너 로그에 topic-partition@offset 을 남기는 MDC 규약을 적용하고 before/after 를 비교한다
  • 브로커는 살아 있는데 리스너 컨테이너만 죽은 상태를 재현하고, 이를 잡아내는 커스텀 HealthIndicator 를 만든다

선행 스텝: Step 11 — 테스트 예상 소요: 90분


12-0. 실습 준비

build.gradle 에는 이미 관측성 의존성이 들어 있습니다(실습 프로젝트 셋업 P-4).

implementation 'org.springframework.boot:spring-boot-starter-actuator'
implementation 'io.micrometer:micrometer-registry-prometheus'
implementation 'io.micrometer:micrometer-tracing-bridge-brave'
implementation 'io.zipkin.reporter2:zipkin-reporter-brave'

이 스텝에서만 application.ymlmanagement 블록을 넓힙니다.

management:
  endpoints.web.exposure.include: health,info,metrics,prometheus
  endpoint.health:
    show-details: always            # 커스텀 HealthIndicator 의 details 를 보려면 필수
    group.readiness.include: kafkaListeners   # 12-10 에서 만듭니다
  metrics.tags.application: order-service
  tracing.sampling.probability: 1.0 # 실습이므로 전부 샘플링. 운영은 0.01~0.1

logging.pattern.level: "%5p [${spring.application.name:},%X{traceId:-},%X{spanId:-}]"

컨슈머 그룹은 s12-inventory, s12-notification 을 씁니다. 실행 전 오프셋을 리셋하세요.

kcg --group s12-inventory --topic orders --reset-offsets --to-earliest --execute

💡 브로커 자체의 JMX 지표(UnderReplicatedPartitions, RequestHandlerAvgIdlePercent 등)와 브로커 모니터링 구성은 Kafka 코스 Step 14 에서 다룹니다. 이 스텝은 애플리케이션이 스스로 내보내는 신호만 봅니다.


12-1. 무엇을 봐야 하는가 — 4대 신호

Kafka 애플리케이션에서 대시보드에 올릴 지표는 사실 많지 않습니다. 네 가지면 대부분의 장애를 잡습니다.

신호무엇을 뜻하나나빠지면 생기는 일대표 지표
컨슈머 랙브로커에 쌓인 마지막 오프셋과 그룹이 커밋한 오프셋의 차이처리가 생산을 못 따라감. 방치하면 retention 에 걸려 메시지가 소비되기 전에 삭제kafka.consumer.fetch.manager.records.lag.max
처리 지연(end-to-end)프로듀서가 보낸 시각 ~ 컨슈머가 처리 완료한 시각랙이 0 이어도 한 건 처리에 5초면 사용자 체감은 장애spring.kafka.listener 타이머
에러·DLT 유입률리스너 예외 발생 건수, orders.DLT 로 흘러간 건수조용히 버려지는 주문. 에러 로그가 없어도 DLT 는 쌓임DLT 토픽의 records.consumed.rate
리밸런스 빈도그룹 멤버십이 재구성된 횟수리밸런스 중에는 아무도 소비하지 않음. 잦으면 처리량이 톱니처럼 요동kafka.consumer.coordinator.rebalance.total

이 넷은 서로를 설명합니다. 랙이 튀면 → 처리 지연을 보고 → 에러율을 보고 → 리밸런스를 봅니다. 리밸런스가 잦아서 랙이 튀는 경우와, 처리가 느려서 랙이 튀는 경우는 대응이 완전히 다릅니다.

  랙 급증
   ├─ 처리 지연 정상 + 리밸런스 급증  → 컨슈머가 계속 쫓겨남 (max.poll.interval.ms 초과, Step 07)
   ├─ 처리 지연 급증 + 에러율 0       → 다운스트림(DB·외부 API)이 느려짐
   ├─ 처리 지연 급증 + 에러율 급증    → 블로킹 재시도가 파티션을 잡고 있음 (Step 07)
   └─ 전부 정상인데 랙만 큼           → 생산량이 늘었음. concurrency/파티션 증설 검토

12-2. Micrometer 연동 — 자동으로 나오는 것들

Spring Boot 는 KafkaAutoConfiguration 에서 MicrometerConsumerListenerMicrometerProducerListenerDefaultKafkaConsumerFactory / DefaultKafkaProducerFactory 에 자동으로 붙입니다. 여러분이 아무것도 안 해도 Kafka 클라이언트가 내부적으로 갖고 있는 JMX 지표가 Micrometer 로 넘어옵니다.

curl -s localhost:8080/actuator/metrics | jq -r '.names[]' | grep '^kafka'

결과 (일부)

kafka.consumer.coordinator.assigned.partitions
kafka.consumer.coordinator.rebalance.total
kafka.consumer.fetch.manager.fetch.latency.avg
kafka.consumer.fetch.manager.records.consumed.rate
kafka.consumer.fetch.manager.records.lag
kafka.consumer.fetch.manager.records.lag.max
kafka.producer.record.error.rate
kafka.producer.record.send.rate
kafka.producer.request.latency.avg

핵심만 추리면 이렇습니다.

지표타입의미볼 때 주의
kafka.consumer.fetch.manager.records.lag.maxGauge이 컨슈머가 담당한 파티션 중 가장 큰 랙커밋 기준이 아니라 fetch 위치 기준
kafka.consumer.fetch.manager.records.consumed.rateGauge초당 소비 레코드 수0 이면 처리 중이거나 파티션 미할당
kafka.consumer.coordinator.rebalance.totalCounter이 컨슈머가 겪은 누적 리밸런스 횟수재시작하면 0으로 리셋
kafka.consumer.coordinator.assigned.partitionsGauge현재 할당된 파티션 수0 이면 놀고 있는 스레드(Step 03)
kafka.producer.record.send.rateGauge초당 전송 레코드 수
kafka.producer.record.error.rateGauge초당 전송 실패 수0 이 아니면 즉시 알람

한 지표를 열어 봅니다.

curl -s 'localhost:8080/actuator/metrics/kafka.consumer.fetch.manager.records.lag.max' | jq

결과

{
  "name": "kafka.consumer.fetch.manager.records.lag.max",
  "baseUnit": "records",
  "measurements": [ { "statistic": "VALUE", "value": 47.0 } ],
  "availableTags": [
    { "tag": "spring.id", "values": ["consumerFactory.s12-inventory-0"] },
    { "tag": "client.id", "values": ["consumer-s12-inventory-1", "consumer-s12-inventory-2", "consumer-s12-inventory-3"] }
  ]
}

measurements[0].value47.0 인데, client.id 태그에는 값이 세 개입니다. 태그를 안 지정하면 세 컨슈머 스레드의 값이 합산(sum) 됩니다. ?tag=client.id:consumer-s12-inventory-2 를 붙여 하나만 보면 31.0 입니다.

⚠️ 함정 — lag.max 를 태그 없이 그래프로 그리면 값이 부풀려진다 Micrometer 의 /actuator/metrics 는 같은 이름의 미터를 합산해서 보여 줍니다. concurrency: 3 이면 컨슈머가 3개이므로 "최대 랙"이 세 개 더해진 값이 나옵니다. 위 예에서 실제 최댓값은 31 인데 화면에는 47 이 찍혔습니다. Prometheus 로 긁을 때는 시계열이 client_id 별로 분리되므로 max by (client_id) 로 집계하면 됩니다. Actuator 화면의 숫자를 그대로 임계값으로 쓰지 마세요.


12-3. /actuator/prometheus 실제 출력

micrometer-registry-prometheus 가 클래스패스에 있으면 /actuator/prometheus 가 열립니다. 지표 이름의 ._ 로, 태그는 라벨로 변환됩니다.

curl -s localhost:8080/actuator/prometheus | grep -E '^kafka_consumer_fetch_manager_records_lag'

결과

# TYPE kafka_consumer_fetch_manager_records_lag_max gauge
kafka_consumer_fetch_manager_records_lag_max{application="order-service",client_id="consumer-s12-inventory-1",kafka_version="3.6.1",spring_id="consumerFactory.s12-inventory-0",} 12.0
kafka_consumer_fetch_manager_records_lag_max{application="order-service",client_id="consumer-s12-inventory-2",kafka_version="3.6.1",spring_id="consumerFactory.s12-inventory-1",} 31.0
# TYPE kafka_consumer_fetch_manager_records_lag gauge
kafka_consumer_fetch_manager_records_lag{application="order-service",client_id="consumer-s12-inventory-2",partition="1",topic="orders",} 31.0

records.lag(단수) 쪽에만 topic / partition 라벨이 붙습니다. 파티션별 알람은 이쪽으로만 가능하고, lag.max 는 파티션이 수백 개일 때 시계열 카디널리티를 줄이는 요약용입니다.

# prometheus.yml
scrape_configs:
  - job_name: order-service
    metrics_path: /actuator/prometheus
    scrape_interval: 15s
    static_configs:
      - targets: ['host.docker.internal:8080']
max by (topic) (kafka_consumer_fetch_manager_records_lag)      # 그룹 전체의 최대 랙
increase(kafka_consumer_coordinator_rebalance_total[5m])       # 5분간 리밸런스 증가량
rate(kafka_producer_record_error_total[1m]) > 0                # 프로듀서 전송 실패

12-4. ⚠️ 함정 — 클라이언트가 보고하는 랙은 "그 클라이언트가 아는 랙"이다

여기가 이 스텝에서 가장 중요한 대목입니다.

kafka.consumer.fetch.manager.records.lag컨슈머 객체 안의 카운터입니다. 컨슈머가 브로커에서 fetch 응답을 받을 때 같이 오는 logStartOffset/logEndOffset 으로 계산해 갱신합니다. 즉 그 컨슈머가 살아서 poll 하고 있을 때만 존재하는 값입니다.

앱이 통째로 죽으면 그나마 스크레이프 "타깃 down" 으로 잡힙니다. 진짜 문제는 앱은 살아 있는데 리스너 컨테이너만 죽은 경우입니다. Step 10 의 registry.stop() 이나, Step 07 의 max.poll.interval.ms 초과로 그룹에서 쫓겨난 상태를 재현합니다.

curl -X POST localhost:8080/admin/listeners/stop

결과

INFO 13418 --- [nio-8080-exec-3] o.s.k.l.KafkaMessageListenerContainer    : s12-inventory: Consumer stopped
INFO 13418 --- [ntainer#0-0-C-1] o.a.k.c.c.i.ConsumerCoordinator          : [Consumer clientId=consumer-s12-inventory-1, groupId=s12-inventory] Revoke previously assigned partitions orders-0
INFO 13418 --- [ntainer#0-0-C-1] o.a.k.c.m.Metrics                        : Metrics scheduler closed

이 상태에서 다시 긁으면 grep -c kafka_consumer_fetch_manager_records_lag0 을 돌려줍니다. 지표가 사라졌습니다. 그동안 프로듀서는 계속 발행하므로 실제 랙은 무한히 커지는 중인데, max(kafka_consumer_fetch_manager_records_lag) > 1000 같은 알람 룰은 평가할 시계열이 없어서 발화하지 않습니다.

실제 랙:   0 ──── 120 ──── 480 ──── 2,100 ──── 8,700 ──── ∞
보고된 랙: 0 ──── 120 ──── (지표 소멸) ─────────────────────
알람:      정상 ── 정상 ── 정상 ── 정상 ── 정상 ── 정상  ← 끝까지 안 울림

⚠️ 함정 — 랙이 무한대로 커지는 순간에 랙 지표가 없어진다 컨슈머가 죽으면 그 컨슈머가 만들던 미터도 함께 등록 해제됩니다. 가장 위험한 순간에 관측이 끊깁니다. 증상은 "대시보드의 랙 그래프가 뚝 끊기고 알람은 조용함" 입니다. 그래프가 0 으로 떨어지는 게 아니라 선이 사라집니다. 해결은 두 가지를 함께 걸어야 합니다. ① absent(kafka_consumer_fetch_manager_records_lag{...}) 룰로 지표 부재 자체를 알람한다. ② 브로커에서 직접 조회하는 랙 게이지를 별도로 만든다 → 12-5. ②가 정공법입니다. 컨슈머가 다 죽어도 앱이 살아 있으면 AdminClient 는 답을 줍니다.


12-5. 컨슈머 랙을 직접 노출하기 (핵심 실습)

AdminClient그룹이 커밋한 오프셋을 가져오고, KafkaConsumer.endOffsets토픽의 마지막 오프셋을 가져와 뺍니다. 이 계산은 kafka-consumer-groups.sh --describe 가 하는 일과 정확히 같습니다.

@Component
@Profile("step12-lag")
public class ConsumerGroupLagMetrics {

    private static final String GROUP = "s12-inventory";

    private final AdminClient admin;                // KafkaAdmin.getConfigurationProperties() 로 생성
    private final Consumer<String, Object> probe;   // cf.createConsumer("lag-probe", "-probe")
    private final MultiGauge lagGauge;              // MultiGauge.builder("kafka.consumer.group.lag")...

    @Scheduled(fixedDelay = 10_000, initialDelay = 5_000)
    public void refresh() throws Exception {
        Map<TopicPartition, OffsetAndMetadata> committed =
                admin.listConsumerGroupOffsets(GROUP)
                     .partitionsToOffsetAndMetadata()
                     .get(5, TimeUnit.SECONDS);

        if (committed.isEmpty()) {              // 그룹이 아직 커밋한 적 없음
            lagGauge.register(List.of(), true);
            return;
        }
        Map<TopicPartition, Long> end = probe.endOffsets(committed.keySet());

        List<MultiGauge.Row<?>> rows = new ArrayList<>();
        long total = 0;
        for (var e : committed.entrySet()) {
            TopicPartition tp = e.getKey();
            long lag = Math.max(0, end.getOrDefault(tp, 0L) - e.getValue().offset());
            total += lag;
            rows.add(MultiGauge.Row.of(
                    Tags.of("group", GROUP, "topic", tp.topic(),
                            "partition", String.valueOf(tp.partition())),
                    lag));
        }
        lagGauge.register(rows, true);          // ★ true = 기존 행을 덮어쓴다
        log.info("group={} totalLag={} partitions={}", GROUP, total, rows.size());
    }
}

@Scheduled 를 쓰므로 애플리케이션 클래스에 @EnableScheduling 이 필요합니다. 전체 코드는 Practice.java 에 있습니다.

결과 (10초마다)

INFO 13418 --- [   scheduling-1] c.e.o.s.ConsumerGroupLagMetrics          : group=s12-inventory totalLag=0 partitions=3
INFO 13418 --- [   scheduling-1] c.e.o.s.ConsumerGroupLagMetrics          : group=s12-inventory totalLag=284 partitions=3
INFO 13418 --- [   scheduling-1] c.e.o.s.ConsumerGroupLagMetrics          : group=s12-inventory totalLag=1176 partitions=3
curl -s 'localhost:8080/actuator/metrics/kafka.consumer.group.lag' | jq -c '.measurements, [.availableTags[].tag]'

결과

[{"statistic":"VALUE","value":1176.0}]
["group","topic","partition","application"]
curl -s localhost:8080/actuator/prometheus | grep '^kafka_consumer_group_lag'

결과

# TYPE kafka_consumer_group_lag gauge
kafka_consumer_group_lag{application="order-service",group="s12-inventory",partition="0",topic="orders",} 391.0
kafka_consumer_group_lag{application="order-service",group="s12-inventory",partition="1",topic="orders",} 402.0
kafka_consumer_group_lag{application="order-service",group="s12-inventory",partition="2",topic="orders",} 383.0

이제 리스너 컨테이너를 멈춰도 지표가 살아 있습니다. 12-4 와 같은 실험을 반복합니다.

curl -X POST localhost:8080/admin/listeners/stop
sleep 30
curl -s localhost:8080/actuator/prometheus | grep '^kafka_consumer_group_lag{.*partition="1"'

결과

kafka_consumer_group_lag{application="order-service",group="s12-inventory",partition="1",topic="orders",} 1284.0

클라이언트 지표는 사라졌지만 이쪽은 계속 올라갑니다. 이것이 알람을 걸어야 할 지표입니다.

클라이언트 지표 (records.lag)AdminClient 게이지 (consumer.group.lag)
기준컨슈머의 fetch 위치그룹의 커밋된 오프셋
컨슈머가 죽으면지표 소멸계속 보고됨
갱신 주기fetch 마다(빠름)@Scheduled 주기(10초)
비용0listConsumerGroupOffsets + endOffsets 호출
용도순간 처리 추이알람

💡 실무 팁 — 랙 조회 전용 컨슈머는 절대 poll() 하지 마세요. cf.createConsumer("lag-probe", "-probe") 로 만든 컨슈머는 groupId 를 실제 그룹과 다르게 두고 subscribe 하지 않습니다. 실수로 실제 그룹 ID 로 subscribe 하면 랙 측정용 컨슈머가 그룹에 합류해 리밸런스를 일으키고 파티션을 가져갑니다. endOffsets 는 그룹 멤버십과 무관한 메타데이터 조회라 assign/subscribe 없이 동작합니다.

💡 probe.endOffsets(...) 대신 admin.listOffsets(Map.of(tp, OffsetSpec.latest())) 로도 같은 값을 얻습니다. 컨슈머 객체를 하나 덜 만들어도 되므로 운영 코드에서는 이쪽이 더 깔끔합니다. 이 교재는 두 API 의 대응을 보이려고 endOffsets 를 썼습니다.


12-6. @KafkaListener 처리 시간 측정

ContainerProperties.setMicrometerEnabled(true)기본값이 true 입니다. 리스너 컨테이너는 레코드 하나를 처리할 때마다 spring.kafka.listener 타이머에 기록합니다.

curl -s 'localhost:8080/actuator/metrics/spring.kafka.listener' | jq

결과

{
  "name": "spring.kafka.listener",
  "description": "Kafka Listener Timer",
  "baseUnit": "seconds",
  "measurements": [
    { "statistic": "COUNT",     "value": 1842.0 },
    { "statistic": "TOTAL_TIME","value": 41.196 },
    { "statistic": "MAX",       "value": 1.084 }
  ],
  "availableTags": [
    { "tag": "result",    "values": ["success", "failure"] },
    { "tag": "name",      "values": ["s12-inventory-0", "s12-inventory-1", "s12-inventory-2"] },
    { "tag": "exception", "values": ["none", "IllegalStateException"] }
  ]
}

평균 처리 시간은 TOTAL_TIME / COUNT = 41.196 / 1842 = 22.4ms, 최댓값은 1.084초입니다. 실패만 따로 보려면 ?tag=result:failure 를 붙입니다(위 데이터에서는 COUNT=19, TOTAL_TIME=0.412).

이 타이머는 리스너 메서드 호출 구간만 잽니다. 더 잘게 나누고 싶으면 Timer.Sample 을 직접 씁니다.

@KafkaListener(topics = "orders", groupId = "s12-inventory")
public void onOrder(OrderCreated order) {
    Timer.Sample sample = Timer.start(registry);
    String outcome = "ok";
    try {
        inventoryService.deduct(order);          // 실제 비즈니스 구간
    } catch (RuntimeException ex) {
        outcome = "error";
        throw ex;
    } finally {
        sample.stop(registry.timer("order.inventory.deduct", "outcome", outcome));
    }
}

결과

{ "name": "order.inventory.deduct",
  "measurements": [ { "statistic": "COUNT", "value": 1842.0 },
                    { "statistic": "TOTAL_TIME", "value": 38.940 },
                    { "statistic": "MAX", "value": 1.061 } ] }

spring.kafka.listener 의 41.196초 중 38.940초가 재고 차감입니다. 나머지 2.256초(건당 1.2ms)가 역직렬화·인터셉터·Observation 등 프레임워크 오버헤드입니다. "Kafka 가 느리다"는 신고의 대부분은 이 5% 가 아니라 나머지 95% 쪽입니다.

💡 실무 팁 — @Timed 를 붙였는데 지표가 안 나온다면 TimedAspect 빈이 없는 것입니다. @Timed 는 AOP 기반이라 아래 빈을 직접 등록해야 동작합니다. Boot 는 자동 등록하지 않습니다.

@Bean
public TimedAspect timedAspect(MeterRegistry registry) { return new TimedAspect(registry); }

게다가 @KafkaListener 메서드에 붙일 때는 프록시가 리스너 등록을 가로채면서 어노테이션 탐지가 어긋나는 경우가 있습니다. 리스너 안에서는 Timer.Sample 을 쓰는 편이 확실합니다.

백분위수가 필요하면 management.metrics.distribution.percentiles-histogram.spring.kafka.listener: true 를 켜고 slo 로 버킷 경계를 지정합니다.


12-7. observationEnabled 로 분산 추적

Step 05 에서는 traceId 를 헤더에 직접 넣고 직접 꺼냈습니다. Spring Kafka 3.x 는 이것을 프레임워크가 합니다.

spring:
  kafka:
    template:
      observation-enabled: true
    listener:
      observation-enabled: true

이 두 줄이면 KafkaTemplate.send 와 리스너 실행이 Micrometer Observation 으로 감싸이고, micrometer-tracing-bridge-brave 가 그것을 span 으로 바꿉니다. 프로듀서는 W3C traceparent 헤더를 레코드에 심고, 컨슈머는 그 헤더에서 컨텍스트를 복원합니다.

@Profile("step12-trace")
@RestController
class OrderController {
    @PostMapping("/orders")
    public String create(@RequestParam int seq) {
        OrderCreated order = OrderCreated.of(seq);
        log.info("HTTP 요청 수신 orderId={}", order.orderId());
        kafkaTemplate.send("orders", order.orderId(), order);
        return order.orderId();
    }
}
curl -X POST 'localhost:8080/orders?seq=7'

결과

INFO  [order-service,6f8c2a1b9d4e5f30,3a1b2c4d5e6f7081] 13418 --- [nio-8080-exec-1] c.e.o.step12.OrderController        : HTTP 요청 수신 orderId=ORD-0007
INFO  [order-service,6f8c2a1b9d4e5f30,9c7d1e2f3a4b5c60] 13418 --- [ad | producer-1] c.e.o.step12.TracedProducer         : sent orders-1@318 orderId=ORD-0007
INFO  [order-service,6f8c2a1b9d4e5f30,b2c3d4e5f6a70819] 13418 --- [ntainer#0-1-C-1] c.e.o.step12.TracedListener         : 재고 차감 orderId=ORD-0007 orders-1@318
INFO  [order-service,6f8c2a1b9d4e5f30,c8d9e0f1a2b3c4d5] 13418 --- [ntainer#0-1-C-1] c.e.o.step12.TracedListener         : 알림 발송 orderId=ORD-0007

대괄호 안 가운데 값(6f8c2a1b9d4e5f30)이 네 줄 모두 같습니다. HTTP 요청 → 프로듀서 → 컨슈머가 하나의 trace 입니다. 세 번째 값(spanId)은 각기 다릅니다. 컨슈머 스레드는 HTTP 스레드와 아무 관계가 없는데도 이어졌습니다. 헤더 덕분입니다.

kcc --topic orders --partition 1 --offset 318 --max-messages 1 --property print.headers=true

결과

traceparent:00-6f8c2a1b9d4e5f306f8c2a1b9d4e5f30-9c7d1e2f3a4b5c60-01,__TypeId__:com.example.order.domain.OrderCreated	{"orderId":"ORD-0007","customerId":1007,"sku":"SKU-002","quantity":3,"amount":10000,"createdAt":"2025-01-01T00:07:00Z"}

traceparent 는 W3C 형식 버전-traceId-parentSpanId-플래그 입니다. -01샘플링됨 이라는 뜻입니다. -00 이면 수집기가 버립니다.

항목Step 05 수동 전파observation-enabled
헤더 이름직접 정한 X-Trace-Id표준 traceparent (tracestate, B3 도 선택 가능)
심는 코드ProducerRecord.headers().add(...) 를 매번없음
꺼내는 코드@Header("X-Trace-Id") String + MDC.put없음
MDC 연동직접 put/remove자동 (%X{traceId} 바로 사용)
부모-자식 관계없음. 문자열 한 개span 트리. 지연 구간이 분리됨
Zipkin/OTel 연동불가그대로 전송

수동 전파는 "같은 요청인지"만 알려 주고, Observation 은 "어느 구간에서 느렸는지" 까지 알려 줍니다. 위 trace 는 총 512ms 중 POST /orders 48ms, orders send 11ms, orders receive 464ms 로 쪼개지고, 그 464ms 중 441ms 가 inventory.deduct 였습니다. 문자열 하나만 전파했다면 이 분해는 불가능합니다.


12-8. ⚠️ 함정 — 컨슈머만 켜면 trace 가 끊긴다

observation-enabled 는 프로듀서(template)와 컨슈머(listener)가 별도 설정입니다. 컨슈머만 켜 놓고 "왜 안 이어지지" 하는 경우가 가장 흔합니다.

spring:
  kafka:
    listener:
      observation-enabled: true    # 켬
    # template.observation-enabled 를 빠뜨림  ← 기본 false

결과

INFO  [order-service,6f8c2a1b9d4e5f30,3a1b2c4d5e6f7081] 13418 --- [nio-8080-exec-1] c.e.o.step12.OrderController        : HTTP 요청 수신 orderId=ORD-0007
INFO  [order-service,,] 13418 --- [ad | producer-1] c.e.o.step12.TracedProducer         : sent orders-1@319 orderId=ORD-0007
INFO  [order-service,d41d8cd98f00b204,e9800998ecf8427e] 13418 --- [ntainer#0-1-C-1] c.e.o.step12.TracedListener         : 재고 차감 orderId=ORD-0007 orders-1@319
  • 프로듀서 줄은 [order-service,,] — traceId 가 비어 있습니다.
  • 컨슈머 줄은 traceId 가 있지만 HTTP 요청과 다른 값입니다. 새 trace 를 시작한 것입니다.

헤더를 보면 원인이 명확합니다.

kcc --topic orders --partition 1 --offset 319 --max-messages 1 --property print.headers=true

결과

__TypeId__:com.example.order.domain.OrderCreated	{"orderId":"ORD-0007",...}

traceparent아예 없습니다. 컨슈머 쪽 Observation 은 "부모가 없으니 루트 span 을 만든다"고 판단합니다. 에러도 경고도 없습니다. trace 는 잘 그려지는데 조각조각 끊겨 있을 뿐입니다.

⚠️ 함정 — 이미 쌓여 있던 옛 메시지에는 traceparent 가 없다 설정을 켜고 배포한 뒤에도, 켜기 전에 발행된 메시지는 헤더가 없어 각자 새 trace 로 시작합니다. 랙이 큰 상태에서 배포하면 한동안 "끊긴 trace" 가 대량으로 섞여 보입니다. 장애가 아니라 정상입니다. 마찬가지로 다른 팀 서비스가 발행한 토픽을 구독한다면, 그쪽이 켜기 전까지는 절대 이어지지 않습니다. 확인 방법은 하나뿐입니다. print.headers=true 로 실제 레코드에 traceparent 가 있는지 보는 것.

💡 @RetryableTopic(Step 08)으로 retry 토픽을 거친 메시지는 재발행 시점의 trace 로 이어집니다. 원본 메시지와 재시도 메시지가 같은 trace 로 묶이므로, 재시도 지연이 span 트리에 그대로 드러납니다.


12-9. 로그에 파티션/오프셋 남기기 (실전 필수)

장애 대응은 결국 "그 메시지 하나" 를 찾는 일입니다. topic-partition@offset 이 없으면 찾을 방법이 없습니다.

before

INFO 13418 --- [ntainer#0-1-C-1] c.e.o.step12.InventoryListener          : 재고 차감 실패 orderId=ORD-0311
ERROR 13418 --- [ntainer#0-1-C-1] o.s.k.l.DefaultErrorHandler             : Record in retry, will be re-delivered

이 로그로 할 수 있는 일은 없습니다. ORD-0311 이 어느 파티션 몇 번 오프셋인지 모르면 kcc --offset 으로 꺼내 볼 수도, 그 지점부터 리셋할 수도 없습니다.

RecordInterceptor 로 MDC 에 넣습니다. 리스너 코드를 한 줄도 안 고쳐도 됩니다.

public class MdcRecordInterceptor implements RecordInterceptor<String, Object> {

    @Override
    public ConsumerRecord<String, Object> intercept(ConsumerRecord<String, Object> record,
                                                    Consumer<String, Object> consumer) {
        MDC.put("kafka", record.topic() + "-" + record.partition() + "@" + record.offset());
        MDC.put("kafkaKey", String.valueOf(record.key()));
        return record;
    }

    @Override
    public void afterRecord(ConsumerRecord<String, Object> record, Consumer<String, Object> consumer) {
        MDC.remove("kafka");          // 반드시 지운다. 스레드 재사용
        MDC.remove("kafkaKey");
    }
}

컨테이너 팩토리에 factory.setRecordInterceptor(new MdcRecordInterceptor()) 로 붙이고, 로그 패턴에 [%X{kafka:-}] 를 추가합니다.

logging.pattern.level: "%5p [${spring.application.name:},%X{traceId:-},%X{spanId:-}] [%X{kafka:-}]"

after

INFO  [order-service,6f8c2a1b9d4e5f30,b2c3d4e5f6a70819] [orders-1@311] 13418 --- [ntainer#0-1-C-1] c.e.o.step12.InventoryListener : 재고 차감 실패 orderId=ORD-0311
ERROR [order-service,6f8c2a1b9d4e5f30,b2c3d4e5f6a70819] [orders-1@311] 13418 --- [ntainer#0-1-C-1] o.s.k.l.DefaultErrorHandler    : Record in retry, will be re-delivered

이제 kcc --topic orders --partition 1 --offset 311 --max-messages 1 로 문제의 원본 레코드를 바로 꺼낼 수 있습니다.

⚠️ 함정 — afterRecord 에서 MDC 를 안 지우면 남의 로그에 남의 오프셋이 찍힌다 리스너 스레드(ntainer#0-1-C-1)는 재사용됩니다. MDC.remove 를 빠뜨리면 다음 레코드가 예외로 죽어 인터셉터를 못 타는 경우 직전 레코드의 오프셋이 그대로 찍힙니다. 로그는 orders-1@311 을 가리키는데 실제 문제는 orders-1@312 인 상황이 만들어집니다. 틀린 관측 정보는 관측 정보가 없는 것보다 나쁩니다. afterRecord 는 성공/실패와 무관하게 호출되므로 여기서 지우는 것이 안전합니다. 배치 리스너를 쓴다면 BatchInterceptor 로 같은 일을 하되, 배치 전체의 범위(orders-1@300..349)를 넣습니다.

💡 실무 팁 — 로그 규약 세 줄 ① 리스너의 모든 로그에 topic-partition@offset. ② 예외 로그에는 메시지 키도. ③ DLT 발행 시 원본 좌표를 WARN 으로 한 줄. 이 세 줄이면 "어떤 주문이 어디서 사라졌는가"에 대부분 답할 수 있습니다.


12-10. Actuator 헬스 — 브로커가 살아 있다는 것이 전부가 아니다

spring-kafka 가 클래스패스에 있으면 Boot 는 KafkaHealthIndicator 를 자동 등록합니다.

curl -s localhost:8080/actuator/health/kafka | jq

결과

{ "status": "UP",
  "details": { "clusterId": "4L6g3nShT-eMCtK--X86sw", "brokerId": "1", "nodes": 1 } }

이 헬스가 하는 일은 KafkaAdmin.describeCluster() 호출 한 번입니다. 브로커에 붙을 수 있느냐만 봅니다. 12-4 처럼 리스너 컨테이너를 전부 멈춘 뒤 다시 조회해도 결과는 그대로 UP 입니다.

curl -X POST localhost:8080/admin/listeners/stop
curl -s localhost:8080/actuator/health/kafka | jq -r '.status'    # → UP

⚠️ 함정 — 컨슈머가 한 건도 소비하지 않는데 헬스는 UP 이다 KafkaHealthIndicator 는 리스너 컨테이너를 보지 않습니다. 컨테이너가 죽었든, 파티션을 하나도 할당받지 못했든 UP 입니다. 쿠버네티스 readiness probe 를 /actuator/health 로 걸어 두면, 아무 일도 안 하는 파드가 정상 판정을 받아 계속 떠 있습니다. 랙은 무한히 늘고, 12-4 때문에 랙 지표는 사라진 상태입니다. 세 개의 관측 장치가 동시에 눈을 감습니다. 해결은 컨테이너 상태를 보는 HealthIndicator 를 직접 만드는 것입니다.

@Component("kafkaListeners")            // ← 빈 이름이 곧 /actuator/health/<이름>
@Profile("step12-health")
@RequiredArgsConstructor
public class KafkaListenerHealthIndicator implements HealthIndicator {

    private final KafkaListenerEndpointRegistry registry;

    @Override
    public Health health() {
        Collection<MessageListenerContainer> containers = registry.getListenerContainers();
        if (containers.isEmpty()) {
            return Health.down().withDetail("reason", "등록된 리스너 컨테이너가 없음").build();
        }
        Map<String, Object> details = new LinkedHashMap<>();
        boolean healthy = true;

        for (MessageListenerContainer c : containers) {
            Collection<TopicPartition> assigned = c.getAssignedPartitions();
            boolean running  = c.isRunning();
            boolean hasParts = assigned != null && !assigned.isEmpty();
            healthy &= running && hasParts;

            details.put(c.getListenerId(), Map.of(
                    "running", running,
                    "assignedPartitions", assigned == null ? List.of() : assigned.stream().map(TopicPartition::toString).toList()
            ));
        }
        return (healthy ? Health.up() : Health.down()).withDetails(details).build();
    }
}
curl -s localhost:8080/actuator/health/kafkaListeners | jq -c

결과 (정상 / 컨테이너 정지 후)

{"status":"UP","details":{
  "s12-inventory-0":{"running":true,"assignedPartitions":["orders-0"]},
  "s12-inventory-1":{"running":true,"assignedPartitions":["orders-1"]},
  "s12-inventory-2":{"running":true,"assignedPartitions":["orders-2"]}}}

{"status":"DOWN","details":{
  "s12-inventory-0":{"running":false,"assignedPartitions":[]},
  "s12-inventory-1":{"running":false,"assignedPartitions":[]},
  "s12-inventory-2":{"running":false,"assignedPartitions":[]}}}

이제 /actuator/health 전체 상태도 DOWN 이 되고, 12-0 에서 만든 readiness 그룹이 이 지표만 봅니다.

💡 실무 팁 — liveness 와 readiness 를 구분하세요. 리스너 컨테이너가 죽었다고 파드를 죽이면(liveness) 재시작 → 리밸런스 → 다시 죽음의 루프에 빠질 수 있습니다. readiness 에만 넣어 트래픽에서 빼고 알람을 울리는 편이 안전합니다. Kafka 컨슈머는 HTTP 트래픽을 받지 않으므로 readiness DOWN 의 실질적 효과는 "알람" 입니다. 그것으로 충분합니다.

💡 getAssignedPartitions() 는 리밸런스 진행 중 잠깐 비어 있을 수 있습니다. 운영에서는 연속 N회 DOWN 일 때만 알람을 울리세요. 리밸런스 프로토콜의 단계별 상태는 Kafka 코스 Step 05 를 참고하세요.


12-11. 알람 기준 — 어떤 지표에 무슨 임계값을 걸 것인가

지표임계값(예)왜 이 값인가심각도
kafka_consumer_group_lag 절대값5분 연속 > 10,000정상 처리량(1,000 msg/s)의 10초분. 순간 스파이크는 넘김Warning
증가율deriv(lag[10m]) > 0 이 15분 지속절대값이 작아도 줄지 않으면 결국 터짐. 트래픽이 적은 서비스에 필수Critical
랙 소진 예상 시간lag / rate(consumed[5m]) > 1800"다 처리하는 데 30분 이상" — 절대값보다 직관적Critical
DLT 유입률rate(dlt_records[5m]) > 0DLT 는 0 이 정상. 한 건이라도 들어오면 사람이 봐야 함Critical
kafka_producer_record_error_totalincrease(...[5m]) > 0프로듀서 실패는 곧 이벤트 유실(Step 02)Critical
리밸런스 횟수increase(rebalance_total[10m]) > 3배포 시 1~2회는 정상. 3회 초과는 컨슈머가 쫓겨나는 중Warning
/actuator/health/kafkaListeners3회 연속 DOWN리밸런스 중 순간 DOWN 을 걸러 냄Critical
랙 지표 부재absent(kafka_consumer_group_lag) 2분12-4 의 함정 대비. 앱 자체가 죽은 경우Critical

랙 0 이 정상을 뜻하지 않습니다. Step 06 에서 본 것처럼, AckMode 를 잘못 잡으면 처리에 실패한 메시지의 오프셋까지 커밋됩니다. 그러면 랙은 완벽하게 0 이고, 대시보드는 초록색이고, 주문은 사라집니다. 랙 알람은 "밀리고 있다" 만 잡습니다. "잘못 처리했다" 는 DLT 유입률과 비즈니스 지표(주문 수 대비 재고 차감 수)로 잡아야 합니다.

# 발행량과 처리량의 괴리 — 조용한 유실을 잡는 가장 단순한 룰
sum(rate(kafka_producer_record_send_total{topic="orders"}[10m]))
  - sum(rate(order_inventory_deduct_seconds_count[10m])) > 1

12-12. 운영 CLI 치트시트

장애 상황에서 대시보드보다 빠른 것이 CLI 입니다. 별칭은 실습 프로젝트 셋업 P-8 에 있습니다.

① 랙 확인 — 가장 먼저 치는 명령

kcg --describe --group s12-inventory

결과

GROUP          TOPIC   PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG   CONSUMER-ID                    HOST         CLIENT-ID
s12-inventory  orders  0          1204            1595            391   consumer-s12-inventory-1-4f2a  /172.19.0.1  consumer-s12-inventory-1
s12-inventory  orders  1          1193            1595            402   consumer-s12-inventory-2-9b71  /172.19.0.1  consumer-s12-inventory-2
s12-inventory  orders  2          1212            1595            383   consumer-s12-inventory-3-c033  /172.19.0.1  consumer-s12-inventory-3

CONSUMER-ID- 로 나오면 그 파티션을 아무도 맡고 있지 않다는 뜻입니다. 12-4 의 상황입니다.

② 그룹 상태와 멤버별 할당

kcg --describe --group s12-inventory --state
kcg --describe --group s12-inventory --members --verbose

결과

GROUP          COORDINATOR (ID)      ASSIGNMENT-STRATEGY  STATE        #MEMBERS
s12-inventory  127.0.0.1:9092 (1)    range                Stable       3

GROUP          CONSUMER-ID                    HOST         CLIENT-ID                 #PARTITIONS  ASSIGNMENT
s12-inventory  consumer-s12-inventory-1-4f2a  /172.19.0.1  consumer-s12-inventory-1  1            orders(0)
s12-inventory  consumer-s12-inventory-2-9b71  /172.19.0.1  consumer-s12-inventory-2  1            orders(1)
s12-inventory  consumer-s12-inventory-3-c033  /172.19.0.1  consumer-s12-inventory-3  1            orders(2)

STATEStable(정상) / PreparingRebalance·CompletingRebalance(리밸런스 중, 오래 지속되면 문제) / Empty(멤버 없음, 커밋된 오프셋만 남음) / Dead 중 하나입니다. #PARTITIONS0 인 멤버가 있으면 concurrency 가 파티션 수보다 큰 것입니다(Step 03).

③ 오프셋 리셋 — 반드시 --dry-run 먼저

# 앱을 먼저 종료할 것. --dry-run 으로 NEW-OFFSET 을 확인한 뒤에만 --execute 로 바꾼다
kcg --group s12-inventory --topic orders --reset-offsets --to-datetime 2025-01-01T00:00:00.000 --dry-run

결과

GROUP          TOPIC   PARTITION  NEW-OFFSET
s12-inventory  orders  0          842
s12-inventory  orders  1          851
s12-inventory  orders  2          839
옵션의미
--to-earliest / --to-latest처음 / 끝으로
--to-offset 500절대 오프셋(모든 파티션 동일)
--shift-by -100현재 위치에서 상대 이동. 되감기에 가장 안전
--to-datetime <ISO>시각 기준. 장애 시작 시각으로 되감을 때

④ 메시지 개수 세기와 DLT 확인

# 파티션별 마지막 오프셋(-1). --time -2 는 시작 오프셋. 둘의 차이가 실제 남아 있는 메시지 수
docker exec -it learn-kafka /opt/kafka/bin/kafka-get-offsets.sh \
  --bootstrap-server localhost:9092 --topic orders --time -1

kcc --topic orders.DLT --from-beginning --property print.headers=true --max-messages 3

결과

orders:0:1595
orders:1:1595
orders:2:1595

kafka_dlt-original-topic:orders,kafka_dlt-exception-fqcn:java.lang.IllegalStateException	{"orderId":"ORD-0311",...}

💡 3.7 에서는 kafka-get-offsets.sh 를 쓰지만, 예전 문서의 kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list ... 도 그대로 동작합니다. DLT 헤더 해석은 Step 07 에 정리돼 있습니다.


정리

개념핵심
4대 신호컨슈머 랙 / 처리 지연 / 에러·DLT / 리밸런스. 네 개를 함께 봐야 원인이 갈린다
Micrometer 자동 연동MicrometerConsumer/ProducerListener 가 Boot 자동 설정으로 등록. 코드 0줄
/actuator/metrics 합산같은 이름 미터를 합산해 보여 줌. concurrency=3 이면 랙이 3배로 보임
클라이언트 랙 지표컨슈머가 죽으면 지표도 사라짐. 가장 위험한 순간에 알람이 안 울린다
AdminClient 랙 게이지listConsumerGroupOffsets + endOffsetsMultiGauge. 알람은 이쪽에 건다
랙 조회 전용 컨슈머실제 그룹 ID 로 subscribe 하면 리밸런스를 일으킨다. endOffsets 만 호출할 것
spring.kafka.listener 타이머micrometerEnabled 기본 true. @TimedTimedAspect 빈이 없으면 조용히 안 나옴
observation-enabledtemplatelistener 둘 다 켜야 함. 프로듀서가 traceparent 를 심는다
끊긴 trace프로듀서 미설정 / 옛 메시지 / 타 서비스 발행. 헤더를 직접 보는 것이 유일한 확인법
MDC partition@offsetRecordInterceptor 로 주입, afterRecord 에서 반드시 제거
KafkaHealthIndicatordescribeCluster 만 확인. 컨테이너가 죽어도 UP
커스텀 HealthIndicatorKafkaListenerEndpointRegistry.getListenerContainers()isRunning() + 할당 파티션
랙 0정상을 뜻하지 않는다. 잘못된 커밋이면 랙 0 인 채로 유실(Step 06)

연습문제

Exercise.java 에 6문제가 있습니다. 정답은 Solution.java. 반드시 실제로 curl 을 쳐서 출력을 확인하세요.

  1. /actuator/prometheus 에서 랙 관련 시계열을 전부 찾아, records.lagrecords.lag.max 의 라벨 차이를 표로 정리하기
  2. AdminClient + MultiGauges12-ex-inventory 그룹의 파티션별 랙 게이지를 등록하고, 리스너를 멈춘 뒤에도 값이 올라가는 것을 확인하기
  3. 리스너의 비즈니스 구간에 Timer.Sample 을 붙여 order.inventory.deduct 타이머를 만들고, spring.kafka.listenerTOTAL_TIME 과 비교해 프레임워크 오버헤드를 계산하기
  4. template.observation-enabledlistener.observation-enabled 를 각각 켜고 끄며 4가지 조합의 로그 traceId 를 기록하고, traceparent 헤더 유무를 kcc 로 대조하기
  5. KafkaListenerEndpointRegistry 를 주입받아 컨테이너별 running / 할당 파티션을 노출하는 HealthIndicator 를 만들고, 컨테이너를 멈춰 DOWN 을 확인하기
  6. RecordInterceptortopic-partition@offset 을 MDC 에 넣고, afterRecordMDC.remove일부러 빼서 잘못된 오프셋이 로그에 남는 것을 재현하기

다음 단계

이제 무엇이 잘못되고 있는지 보이는 상태가 됐습니다. 랙은 브로커 기준으로 잡히고, trace 는 HTTP 요청부터 컨슈머까지 이어지고, 로그 한 줄로 원본 메시지를 꺼낼 수 있습니다.

남은 것은 애초에 잘못되지 않게 만드는 설계입니다. 마지막 스텝에서는 재처리해도 안전한 멱등 컨슈머, DB 와 Kafka 를 진짜로 묶는 Transactional Outbox, 그리고 순서 보장이 필요한 구간을 좁히는 전략을 구현하고, 지금까지의 열두 스텝을 하나의 주문 서비스로 합칩니다.

Step 13 — 실전 패턴과 최종 프로젝트


실습 파일

이 스텝은 curl 을 치는 시간이 코드를 쓰는 시간보다 깁니다. Practice.java 를 프로필로 켜 두고, 터미널을 하나 더 띄워 /actuator/metrics/actuator/prometheus 를 계속 두드리세요. 특히 12-4 → 12-5 는 같은 실험을 두 번 하는 구성입니다. 리스너를 멈췄을 때 한쪽 지표는 사라지고 다른 쪽은 살아 있는 것을 눈으로 봐야 이 스텝의 요점이 전달됩니다. 그다음 Exercise.java 의 6문제를 풀고 Solution.java 로 대조합니다. 세 파일 모두 com.example.order.step12 패키지에 둡니다.

Practice.java

본문 12-4 ~ 12-10 의 예제를 절 번호 주석과 함께 nested static class 로 담은 실행 파일입니다.

  • 프로필 4개(step12-lag, step12-trace, step12-mdc, step12-health)로 나뉘어 있지만, 이 스텝은 여러 개를 함께 켜도 됩니다. 관측 코드끼리는 간섭하지 않습니다. 전부 보려면 --spring.profiles.active=step12,step12-lag,step12-trace,step12-mdc,step12-health 로 실행하세요.
  • [12-5] ConsumerGroupLagMetrics 가 핵심입니다. cf.createConsumer("lag-probe", "-probe") 로 만든 프로브 컨슈머는 subscribeassign 도 하지 않습니다. endOffsets 는 메타데이터 조회라 그것만으로 동작합니다. 실수로 실제 그룹 ID 를 넘기면 리밸런스가 발생하니 그룹 ID 인자를 절대 바꾸지 마세요.
  • [12-4] ListenerAdminController/admin/listeners/stop, /start, /status 세 엔드포인트를 엽니다. 이게 있어야 12-4 와 12-10 의 "컨테이너만 죽은 상태"를 안전하게 재현할 수 있습니다. Step 10 의 KafkaListenerEndpointRegistry 를 그대로 씁니다.
  • [12-7] TracedProducer / TracedListener 는 로그 한 줄만 찍습니다. 볼 것은 로그 본문이 아니라 대괄호 안의 traceId 입니다. application-step12-trace.ymltemplate.observation-enabled: true 가 들어 있는지 반드시 확인하세요. 빠지면 12-8 의 끊긴 trace 가 재현됩니다.
  • [12-9] MdcRecordInterceptor 에는 afterRecordMDC.remove 를 주석 처리할 수 있게 표시해 두었습니다. 주석을 풀었다 막았다 하며 로그의 오프셋이 어긋나는 것을 직접 보세요.
  • LoadGenerator 는 초당 5건씩 orders 로 계속 발행합니다. 랙을 벌리려면 이게 떠 있어야 합니다. app.load.enabled=false 로 끌 수 있습니다.
package com.example.order.step12;

/*
 * ============================================================================
 * Step 12 — 관측성과 운영 : Practice
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step12/Practice.java
 *
 * 실행
 *   이 스텝의 예제들은 서로 간섭하지 않습니다. 전부 함께 켜도 됩니다.
 *
 *   ./gradlew bootRun --args='--spring.profiles.active=step12,step12-lag,step12-trace,step12-mdc,step12-health'
 *
 *   개별로 보고 싶다면:
 *     step12-lag     → [12-4][12-5] AdminClient 랙 게이지 + 컨테이너 stop/start 엔드포인트
 *     step12-trace   → [12-7][12-8] observation 으로 trace 연결. POST /orders?seq=7
 *     step12-mdc     → [12-9]       RecordInterceptor 로 partition@offset 을 MDC 에
 *     step12-health  → [12-10]      리스너 컨테이너 상태 HealthIndicator
 *
 * ★ 실행 전 오프셋 리셋 (앱을 먼저 종료할 것)
 *   kcg --group s12-inventory --topic orders --reset-offsets --to-earliest --execute
 *
 * ★ application-step12-trace.yml 에 아래 두 줄이 반드시 있어야 합니다.
 *   spring.kafka.template.observation-enabled: true
 *   spring.kafka.listener.observation-enabled: true
 *   template 쪽을 빼면 12-8 의 "끊긴 trace" 가 재현됩니다.
 *
 * ★ logging.pattern.level 에 MDC 키를 넣어야 12-9 가 보입니다.
 *   logging.pattern.level: "%5p [${spring.application.name:},%X{traceId:-},%X{spanId:-}] [%X{kafka:-}]"
 *
 * 실행 중에 두드릴 것
 *   curl -s localhost:8080/actuator/metrics | jq -r '.names[]' | grep '^kafka'
 *   curl -s localhost:8080/actuator/metrics/kafka.consumer.group.lag | jq
 *   curl -s localhost:8080/actuator/metrics/spring.kafka.listener | jq
 *   curl -s localhost:8080/actuator/prometheus | grep '^kafka_consumer_group_lag'
 *   curl -s localhost:8080/actuator/health/kafkaListeners | jq
 *   curl -X POST localhost:8080/admin/listeners/stop      ← ★ 12-4 의 핵심 실험
 *   curl -X POST localhost:8080/admin/listeners/start
 *
 * 실행 중에 확인할 CLI
 *   kcg --describe --group s12-inventory
 *   kcg --describe --group s12-inventory --state
 *   kcc --topic orders --partition 1 --offset 318 --max-messages 1 --property print.headers=true
 * ============================================================================
 */

import com.example.order.domain.OrderCreated;

import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.MultiGauge;
import io.micrometer.core.instrument.Tags;
import io.micrometer.core.instrument.Timer;

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.boot.actuate.health.Health;
import org.springframework.boot.actuate.health.HealthIndicator;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.http.MediaType;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerEndpointRegistry;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.RecordInterceptor;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;

import jakarta.annotation.PreDestroy;

import java.util.ArrayList;
import java.util.Collection;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

public final class Practice {

    private Practice() { }

    /* =======================================================================
     * [12-0] 공통 — 스케줄러 활성화
     * -----------------------------------------------------------------------
     * 12-5 의 랙 게이지가 @Scheduled 로 갱신되므로 반드시 필요합니다.
     * 빠뜨리면 게이지가 등록조차 되지 않고, 에러도 나지 않습니다.
     * ===================================================================== */
    @Configuration
    @Profile("step12")
    @EnableScheduling
    public static class SchedulingConfig {
    }

    /* =======================================================================
     * [12-0] 부하 생성기 — 랙을 벌리려면 계속 발행해야 한다
     * -----------------------------------------------------------------------
     * 초당 5건씩 orders 로 발행합니다. app.load.enabled=false 로 끕니다.
     * 12-4 의 실험(리스너를 멈추면 랙이 무한히 커진다)은 이게 떠 있어야 성립합니다.
     * ===================================================================== */
    @Component
    @Profile("step12")
    public static class LoadGenerator {

        private static final Logger log = LoggerFactory.getLogger(LoadGenerator.class);

        private final KafkaTemplate<String, Object> template;
        private final AtomicInteger seq = new AtomicInteger(1);
        private final boolean enabled;

        public LoadGenerator(KafkaTemplate<String, Object> template,
                             @org.springframework.beans.factory.annotation.Value("${app.load.enabled:true}") boolean enabled) {
            this.template = template;
            this.enabled = enabled;
        }

        @Scheduled(fixedDelay = 200L, initialDelay = 3_000L)
        public void publish() {
            if (!enabled) {
                return;
            }
            OrderCreated order = OrderCreated.of(seq.getAndIncrement());
            template.send("orders", order.orderId(), order);
            if (order.orderId().endsWith("00")) {
                log.info("부하 생성 진행 orderId={}", order.orderId());
            }
        }
    }

    /* =======================================================================
     * [12-4] 리스너 컨테이너 stop/start 엔드포인트
     * -----------------------------------------------------------------------
     * "브로커는 살아 있는데 컨슈머만 죽은 상태" 를 안전하게 재현하는 장치입니다.
     * Step 10 의 KafkaListenerEndpointRegistry 를 그대로 씁니다.
     *
     *   curl -X POST localhost:8080/admin/listeners/stop
     *   curl -s     localhost:8080/actuator/prometheus | grep -c records_lag   → 0  ★
     *   curl -s     localhost:8080/actuator/prometheus | grep -c group_lag     → 3  ★
     *
     * 이 두 줄의 대비가 12-4 → 12-5 의 전부입니다.
     * ===================================================================== */
    @RestController
    @Profile("step12-lag")
    public static class ListenerAdminController {

        private static final Logger log = LoggerFactory.getLogger(ListenerAdminController.class);

        private final KafkaListenerEndpointRegistry registry;

        public ListenerAdminController(KafkaListenerEndpointRegistry registry) {
            this.registry = registry;
        }

        @PostMapping("/admin/listeners/stop")
        public String stop() {
            registry.getListenerContainers().forEach(MessageListenerContainer::stop);
            log.warn("모든 리스너 컨테이너를 정지했습니다. 이제 랙 지표가 어떻게 되는지 보세요.");
            return "stopped";
        }

        @PostMapping("/admin/listeners/start")
        public String start() {
            registry.getListenerContainers().forEach(MessageListenerContainer::start);
            log.info("모든 리스너 컨테이너를 재시작했습니다.");
            return "started";
        }

        @GetMapping(value = "/admin/listeners/status", produces = MediaType.APPLICATION_JSON_VALUE)
        public Map<String, Object> status() {
            Map<String, Object> result = new LinkedHashMap<>();
            for (MessageListenerContainer c : registry.getListenerContainers()) {
                Collection<TopicPartition> assigned = c.getAssignedPartitions();
                result.put(c.getListenerId(), Map.of(
                        "running", c.isRunning(),
                        "assigned", assigned == null ? List.of() : assigned.stream().map(TopicPartition::toString).toList()));
            }
            return result;
        }
    }

    /* =======================================================================
     * [12-5] ★핵심★ AdminClient 로 컨슈머 랙을 직접 노출한다
     * -----------------------------------------------------------------------
     * 클라이언트 지표(kafka.consumer.fetch.manager.records.lag)는 컨슈머가 죽으면
     * 함께 사라집니다. 랙이 가장 위험할 때 지표가 없어지는 것입니다(12-4).
     *
     * 여기서는 브로커에 직접 물어봅니다.
     *   ① AdminClient.listConsumerGroupOffsets(group)  → 그룹이 "커밋한" 오프셋
     *   ② KafkaConsumer.endOffsets(partitions)         → 토픽의 마지막 오프셋
     *   ③ lag = ② - ①
     * 이는 kafka-consumer-groups.sh --describe 가 하는 계산과 정확히 같습니다.
     *
     * ⚠️ probe 컨슈머는 절대 subscribe/assign 하지 않습니다.
     *    실제 그룹 ID 로 subscribe 하면 측정용 컨슈머가 그룹에 합류해
     *    리밸런스를 일으키고 파티션을 빼앗아 갑니다.
     * ===================================================================== */
    @Component
    @Profile("step12-lag")
    public static class ConsumerGroupLagMetrics {

        private static final Logger log = LoggerFactory.getLogger(ConsumerGroupLagMetrics.class);
        private static final String GROUP = "s12-inventory";

        private final AdminClient admin;
        private final Consumer<?, ?> probe;
        private final MultiGauge lagGauge;

        public ConsumerGroupLagMetrics(KafkaAdmin kafkaAdmin,
                                       ConsumerFactory<?, ?> consumerFactory,
                                       MeterRegistry registry) {
            this.admin = AdminClient.create(kafkaAdmin.getConfigurationProperties());
            // groupId 를 실제 그룹과 다르게 준다. endOffsets 만 쓸 것이므로 그룹 합류는 일어나지 않는다.
            this.probe = consumerFactory.createConsumer("s12-lag-probe", "-probe");
            this.lagGauge = MultiGauge.builder("kafka.consumer.group.lag")
                    .description("committed offset 기준 파티션별 컨슈머 랙")
                    .baseUnit("records")
                    .register(registry);
        }

        @Scheduled(fixedDelay = 10_000L, initialDelay = 5_000L)
        public void refresh() {
            try {
                Map<TopicPartition, OffsetAndMetadata> committed =
                        admin.listConsumerGroupOffsets(GROUP)
                                .partitionsToOffsetAndMetadata()
                                .get(5, TimeUnit.SECONDS);

                if (committed.isEmpty()) {
                    // 그룹이 아직 한 번도 커밋하지 않았다. 남아 있는 행을 비운다.
                    lagGauge.register(List.of(), true);
                    log.info("group={} 커밋된 오프셋이 없습니다", GROUP);
                    return;
                }

                Map<TopicPartition, Long> endOffsets = probe.endOffsets(committed.keySet());

                List<MultiGauge.Row<?>> rows = new ArrayList<>();
                long total = 0L;
                for (Map.Entry<TopicPartition, OffsetAndMetadata> e : committed.entrySet()) {
                    TopicPartition tp = e.getKey();
                    long end = endOffsets.getOrDefault(tp, 0L);
                    long lag = Math.max(0L, end - e.getValue().offset());
                    total += lag;
                    rows.add(MultiGauge.Row.of(
                            Tags.of("group", GROUP,
                                    "topic", tp.topic(),
                                    "partition", String.valueOf(tp.partition())),
                            lag));
                }
                // overwrite=true 가 핵심. false 로 두면 값이 첫 측정치에서 얼어붙는다.
                lagGauge.register(rows, true);
                log.info("group={} totalLag={} partitions={}", GROUP, total, rows.size());

            } catch (Exception ex) {
                // 랙 수집 실패가 애플리케이션을 죽여서는 안 된다. 로그만 남기고 다음 주기를 기다린다.
                log.warn("랙 수집 실패 group={} : {}", GROUP, ex.toString());
            }
        }

        @PreDestroy
        public void close() {
            probe.close();
            admin.close();
        }
    }

    /* =======================================================================
     * [12-6] 리스너 처리 시간 측정
     * -----------------------------------------------------------------------
     * ContainerProperties.setMicrometerEnabled(true) 는 기본값 true 이므로
     * spring.kafka.listener 타이머는 아무것도 안 해도 생깁니다.
     * 다만 그 타이머는 "리스너 메서드 호출 구간" 전체를 잽니다.
     * 비즈니스 구간만 따로 재려면 Timer.Sample 을 직접 씁니다.
     *
     * ⚠️ @Timed 는 TimedAspect 빈이 없으면 조용히 아무 지표도 만들지 않습니다.
     *    리스너 안에서는 Timer.Sample 이 확실합니다.
     * ===================================================================== */
    @Component
    @Profile("step12-lag")
    public static class TimedInventoryListener {

        private static final Logger log = LoggerFactory.getLogger(TimedInventoryListener.class);

        private final MeterRegistry registry;

        public TimedInventoryListener(MeterRegistry registry) {
            this.registry = registry;
        }

        @KafkaListener(id = "s12-inventory", topics = "orders", groupId = "s12-inventory")
        public void onOrder(OrderCreated order,
                            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                            @Header(KafkaHeaders.OFFSET) long offset) {
            Timer.Sample sample = Timer.start(registry);
            String outcome = "ok";
            try {
                deduct(order);
            } catch (RuntimeException ex) {
                outcome = "error";
                throw ex;
            } finally {
                sample.stop(registry.timer("order.inventory.deduct", "outcome", outcome));
            }
            if (offset % 100 == 0) {
                log.debug("재고 차감 orderId={} orders-{}@{}", order.orderId(), partition, offset);
            }
        }

        /** 처리 시간을 눈에 보이게 하려고 일부러 느리게 만든 구간입니다. */
        private void deduct(OrderCreated order) {
            try {
                Thread.sleep(20L + (order.quantity() * 3L));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }

    /* =======================================================================
     * [12-7] observation 으로 분산 추적 — HTTP → 프로듀서 → 컨슈머
     * -----------------------------------------------------------------------
     *   curl -X POST 'localhost:8080/orders?seq=7'
     *
     * 세 줄의 로그에서 대괄호 안 가운데 값(traceId)이 같아야 성공입니다.
     *   INFO [order-service,6f8c2a1b9d4e5f30,3a1b2c4d5e6f7081] ... HTTP 요청 수신
     *   INFO [order-service,6f8c2a1b9d4e5f30,9c7d1e2f3a4b5c60] ... sent orders-1@318
     *   INFO [order-service,6f8c2a1b9d4e5f30,b2c3d4e5f6a70819] ... 재고 차감
     *
     * traceId 가 비어 있거나(=[order-service,,]) 다르면 12-8 의 함정입니다.
     * template.observation-enabled 가 켜져 있는지 확인하세요.
     * ===================================================================== */
    @RestController
    @Profile("step12-trace")
    public static class TracedProducer {

        private static final Logger log = LoggerFactory.getLogger(TracedProducer.class);

        private final KafkaTemplate<String, Object> template;

        public TracedProducer(KafkaTemplate<String, Object> template) {
            this.template = template;
        }

        @PostMapping("/orders")
        public String create(@RequestParam int seq) {
            OrderCreated order = OrderCreated.of(seq);
            log.info("HTTP 요청 수신 orderId={}", order.orderId());
            template.send("orders", order.orderId(), order)
                    .thenAccept(result -> log.info("sent {}-{}@{} orderId={}",
                            result.getRecordMetadata().topic(),
                            result.getRecordMetadata().partition(),
                            result.getRecordMetadata().offset(),
                            order.orderId()));
            return order.orderId();
        }
    }

    @Component
    @Profile("step12-trace")
    public static class TracedListener {

        private static final Logger log = LoggerFactory.getLogger(TracedListener.class);

        @KafkaListener(id = "s12-trace-inventory", topics = "orders", groupId = "s12-trace-inventory")
        public void onOrder(OrderCreated order,
                            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                            @Header(KafkaHeaders.OFFSET) long offset) {
            log.info("재고 차감 orderId={} orders-{}@{}", order.orderId(), partition, offset);
        }
    }

    /* =======================================================================
     * [12-9] 로그에 topic-partition@offset 남기기
     * -----------------------------------------------------------------------
     * RecordInterceptor 로 MDC 에 넣으면 리스너 코드를 한 줄도 안 고쳐도 됩니다.
     *
     * ⚠️ afterRecord 의 MDC.remove 를 빼면 스레드가 재사용되면서
     *    직전 레코드의 오프셋이 다음 로그에 그대로 남습니다.
     *    로그는 orders-1@311 을 가리키는데 실제 문제는 orders-1@312 인 상황이 됩니다.
     *    ★ 아래 REMOVE_MDC 를 false 로 바꿔 직접 재현해 보세요. ★
     * ===================================================================== */
    public static class MdcRecordInterceptor implements RecordInterceptor<Object, Object> {

        /** ★ 실험용 스위치. false 로 두면 12-9 의 함정이 재현됩니다. */
        public static final boolean REMOVE_MDC = true;

        @Override
        public ConsumerRecord<Object, Object> intercept(ConsumerRecord<Object, Object> record,
                                                        Consumer<Object, Object> consumer) {
            MDC.put("kafka", record.topic() + "-" + record.partition() + "@" + record.offset());
            MDC.put("kafkaKey", String.valueOf(record.key()));
            return record;
        }

        @Override
        public void afterRecord(ConsumerRecord<Object, Object> record, Consumer<Object, Object> consumer) {
            if (REMOVE_MDC) {
                MDC.remove("kafka");
                MDC.remove("kafkaKey");
            }
        }
    }

    @Configuration
    @Profile("step12-mdc")
    public static class MdcListenerConfig {

        /**
         * MDC 인터셉터를 붙인 전용 컨테이너 팩토리.
         * @KafkaListener 에서 containerFactory = "mdcFactory" 로 지정해 씁니다.
         */
        @Bean("mdcFactory")
        public ConcurrentKafkaListenerContainerFactory<Object, Object> mdcFactory(ConsumerFactory<?, ?> cf) {
            ConcurrentKafkaListenerContainerFactory<Object, Object> factory =
                    new ConcurrentKafkaListenerContainerFactory<>();
            @SuppressWarnings("unchecked")
            ConsumerFactory<Object, Object> typed = (ConsumerFactory<Object, Object>) cf;
            factory.setConsumerFactory(typed);
            factory.setConcurrency(3);
            factory.setRecordInterceptor(new MdcRecordInterceptor());
            return factory;
        }
    }

    @Component
    @Profile("step12-mdc")
    public static class MdcInventoryListener {

        private static final Logger log = LoggerFactory.getLogger(MdcInventoryListener.class);

        @KafkaListener(id = "s12-mdc-inventory", topics = "orders",
                       groupId = "s12-mdc-inventory", containerFactory = "mdcFactory")
        public void onOrder(OrderCreated order) {
            // 로그 본문에 파티션/오프셋을 쓰지 않았는데도 패턴의 %X{kafka} 로 찍힙니다.
            if (order.orderId().endsWith("11")) {
                log.warn("재고 차감 실패 orderId={}", order.orderId());
                throw new IllegalStateException("재고 부족: " + order.orderId());
            }
            log.info("재고 차감 orderId={}", order.orderId());
        }
    }

    /* =======================================================================
     * [12-10] 리스너 컨테이너 상태 HealthIndicator
     * -----------------------------------------------------------------------
     * Boot 의 KafkaHealthIndicator 는 describeCluster() 만 봅니다.
     * 브로커에 붙기만 하면 컨테이너가 전부 죽어 있어도 UP 입니다.
     *
     *   curl -s localhost:8080/actuator/health/kafka          → UP  (컨테이너가 죽어도)
     *   curl -s localhost:8080/actuator/health/kafkaListeners → DOWN ★
     *
     * 빈 이름 "kafkaListeners" 가 그대로 엔드포인트 경로가 됩니다.
     * ===================================================================== */
    @Component("kafkaListeners")
    @Profile("step12-health")
    public static class KafkaListenerHealthIndicator implements HealthIndicator {

        private final KafkaListenerEndpointRegistry registry;

        public KafkaListenerHealthIndicator(KafkaListenerEndpointRegistry registry) {
            this.registry = registry;
        }

        @Override
        public Health health() {
            Collection<MessageListenerContainer> containers = registry.getListenerContainers();
            if (containers.isEmpty()) {
                return Health.down().withDetail("reason", "등록된 리스너 컨테이너가 없음").build();
            }

            Map<String, Object> details = new LinkedHashMap<>();
            boolean healthy = true;

            for (MessageListenerContainer c : containers) {
                Collection<TopicPartition> assigned = c.getAssignedPartitions();
                boolean running = c.isRunning();
                // running=true 인데 파티션이 하나도 없는 "좀비" 상태도 비정상으로 본다.
                boolean hasPartitions = assigned != null && !assigned.isEmpty();
                healthy = healthy && running && hasPartitions;

                details.put(c.getListenerId(), Map.of(
                        "running", running,
                        "assignedPartitions",
                        assigned == null ? List.of() : assigned.stream().map(TopicPartition::toString).toList()));
            }
            return (healthy ? Health.up() : Health.down()).withDetails(details).build();
        }
    }

    /* =======================================================================
     * [12-11] 알람용 보조 지표 — DLT 유입 카운터
     * -----------------------------------------------------------------------
     * DLT 는 "0 이 정상" 인 토픽입니다. 한 건이라도 들어오면 사람이 봐야 합니다.
     * 리스너를 하나 붙여 Counter 로 세면 알람 룰이 단순해집니다.
     *   rate(order_dlt_received_total[5m]) > 0
     * ===================================================================== */
    @Component
    @Profile("step12-lag")
    public static class DltCounter {

        private static final Logger log = LoggerFactory.getLogger(DltCounter.class);

        private final MeterRegistry registry;

        public DltCounter(MeterRegistry registry) {
            this.registry = registry;
        }

        @KafkaListener(id = "s12-dlt-monitor", topics = "orders.DLT", groupId = "s12-dlt-monitor")
        public void onDlt(ConsumerRecord<String, Object> record,
                          @Header(name = "kafka_dlt-original-topic", required = false) byte[] originalTopic) {
            String origin = originalTopic == null ? "unknown" : new String(originalTopic);
            registry.counter("order.dlt.received", "originalTopic", origin).increment();
            log.warn("DLT 유입 key={} {}-{}@{}", record.key(), record.topic(), record.partition(), record.offset());
        }
    }

    /* =======================================================================
     * [12-0] 부하 생성기 설정 바인딩용 (선택)
     * ===================================================================== */
    @ConfigurationProperties(prefix = "app.load")
    public static class LoadProperties {

        private boolean enabled = true;

        public boolean isEnabled() {
            return enabled;
        }

        public void setEnabled(boolean enabled) {
            this.enabled = enabled;
        }
    }
}

Exercise.java

6문제의 문제지입니다. 각 문제는 // 여기에 작성: 자리를 비워 두었고, 컴파일은 되도록 뼈대를 남겨 두었습니다.

  • 문제 1·4 는 코드를 거의 안 씁니다. // 관측 기록: 주석의 표를 curl 출력으로 채우는 문제입니다. 손으로 채우지 말고 실제 출력을 붙여 넣으세요.
  • 문제 2 는 이 스텝의 핵심 실습입니다. MultiGauge.register(rows, true) 의 두 번째 인자를 false 로 두면 어떻게 되는지도 함께 확인하도록 // 실험: 주석을 달아 두었습니다.
  • 문제 3spring.kafka.listenerTOTAL_TIME 에서 order.inventory.deductTOTAL_TIME 을 뺀 값을 계산합니다. 두 타이머의 COUNT 가 같은지 먼저 확인하세요. 다르면 비교 자체가 무의미합니다.
  • 문제 6일부러 버그를 만드는 문제입니다. MDC.remove 를 뺀 상태로 예외를 던지는 레코드를 하나 섞어, 다음 레코드의 로그에 이전 오프셋이 남는 것을 관측합니다. 관측이 끝나면 반드시 원복하세요.
  • 각 문제 끝의 // 확인: 주석에 기대 출력 한 줄이 적혀 있습니다. 문제 2 의 확인 줄은 kafka_consumer_group_lag{...,partition="1",topic="orders",} 402.0 형태입니다.
package com.example.order.step12;

/*
 * ============================================================================
 * Step 12 — 관측성과 운영 : Exercise (6문제)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step12/Exercise.java
 *
 * 실행
 *   ./gradlew bootRun --args='--spring.profiles.active=step12,step12-ex'
 *
 * ★ 실행 전 오프셋 리셋 (앱을 먼저 종료할 것)
 *   kcg --group s12-ex-inventory --topic orders --reset-offsets --to-earliest --execute
 *
 * ★ 이 스텝은 "코드를 쓰는 문제" 보다 "curl 을 쳐서 출력을 기록하는 문제" 가 많습니다.
 *   // 관측 기록: 주석의 표는 손으로 지어내지 말고 실제 출력을 붙여 넣으세요.
 *
 * 두드릴 엔드포인트
 *   curl -s localhost:8080/actuator/metrics | jq -r '.names[]' | grep '^kafka'
 *   curl -s localhost:8080/actuator/prometheus | grep -E 'records_lag|group_lag'
 *   curl -s localhost:8080/actuator/metrics/spring.kafka.listener | jq
 *   curl -s localhost:8080/actuator/health/kafkaListeners | jq
 *   curl -X POST localhost:8080/admin/listeners/stop
 * ============================================================================
 */

import com.example.order.domain.OrderCreated;

import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.MultiGauge;
import io.micrometer.core.instrument.Tags;
import io.micrometer.core.instrument.Timer;

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.boot.actuate.health.Health;
import org.springframework.boot.actuate.health.HealthIndicator;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerEndpointRegistry;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.RecordInterceptor;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;

import jakarta.annotation.PreDestroy;

import java.util.ArrayList;
import java.util.Collection;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;

public final class Exercise {

    private Exercise() { }

    /* =======================================================================
     * 문제 1. /actuator/prometheus 에서 랙 지표 찾기
     * -----------------------------------------------------------------------
     * 요구사항
     *   - 앱을 띄우고 아래 명령으로 랙 관련 시계열을 전부 뽑는다.
     *       curl -s localhost:8080/actuator/prometheus | grep -E '^kafka_consumer_fetch_manager_records_lag'
     *   - kafka_consumer_fetch_manager_records_lag (단수) 와
     *     kafka_consumer_fetch_manager_records_lag_max (max) 의 라벨 차이를 표로 정리한다.
     *   - "파티션 단위로 알람을 걸려면 둘 중 무엇을 써야 하는가" 에 답한다.
     *
     * // 관측 기록: 실제 출력을 붙여 넣으세요
     * // ┌─────────────────────────────────┬──────────────────────────────────┐
     * // │ 지표 이름                        │ 라벨                              │
     * // ├─────────────────────────────────┼──────────────────────────────────┤
     * // │ ..._records_lag                 │                                  │
     * // │ ..._records_lag_max             │                                  │
     * // └─────────────────────────────────┴──────────────────────────────────┘
     * //
     * // 파티션 단위 알람에 써야 할 지표:
     * // 그렇게 판단한 이유:
     *
     * 확인: records_lag 쪽에만 topic= / partition= 라벨이 보여야 합니다.
     * ===================================================================== */

    /* =======================================================================
     * 문제 2. AdminClient 로 랙 게이지 등록하기  ★핵심★
     * -----------------------------------------------------------------------
     * 요구사항
     *   - 그룹 "s12-ex-inventory" 의 파티션별 랙을 MultiGauge 로 등록한다.
     *   - 지표 이름은 "kafka.consumer.group.lag", 태그는 group/topic/partition.
     *   - AdminClient.listConsumerGroupOffsets 로 커밋 오프셋을,
     *     probe 컨슈머의 endOffsets 로 마지막 오프셋을 가져와 뺀다.
     *   - @Scheduled(fixedDelay = 10_000) 로 갱신한다.
     *   - 등록 후 curl -X POST localhost:8080/admin/listeners/stop 으로 리스너를 멈추고,
     *     30초 뒤에도 값이 "계속 올라가는지" 확인한다.
     *
     * ⚠️ probe 컨슈머의 groupId 는 반드시 실제 그룹과 다르게 줄 것.
     *    같은 그룹으로 subscribe 하면 리밸런스가 발생합니다.
     *
     * // 실험: lagGauge.register(rows, false) 로 바꿔서 한 번 더 돌려 보세요.
     * //       값이 어떻게 되는지, 에러가 나는지 기록하세요.
     * // 관측 기록:
     *
     * 확인: curl -s localhost:8080/actuator/prometheus | grep '^kafka_consumer_group_lag'
     *       kafka_consumer_group_lag{...,partition="1",topic="orders",} 402.0
     * ===================================================================== */
    @Component
    @Profile("step12-ex")
    public static class Ex2LagMetrics {

        private static final Logger log = LoggerFactory.getLogger(Ex2LagMetrics.class);
        private static final String GROUP = "s12-ex-inventory";

        private final AdminClient admin;
        private final Consumer<?, ?> probe;
        private final MultiGauge lagGauge;

        public Ex2LagMetrics(KafkaAdmin kafkaAdmin,
                             ConsumerFactory<?, ?> consumerFactory,
                             MeterRegistry registry) {
            this.admin = AdminClient.create(kafkaAdmin.getConfigurationProperties());
            this.probe = consumerFactory.createConsumer("s12-ex-lag-probe", "-probe");

            // 여기에 작성: MultiGauge 를 만들어 lagGauge 에 대입하세요.
            //   MultiGauge.builder("kafka.consumer.group.lag")
            //             .description(...).baseUnit("records").register(registry)
            this.lagGauge = null;
        }

        @Scheduled(fixedDelay = 10_000L, initialDelay = 5_000L)
        public void refresh() {
            try {
                Map<TopicPartition, OffsetAndMetadata> committed =
                        admin.listConsumerGroupOffsets(GROUP)
                                .partitionsToOffsetAndMetadata()
                                .get(5, TimeUnit.SECONDS);

                if (committed.isEmpty()) {
                    return;
                }

                // 여기에 작성:
                //   ① probe.endOffsets(committed.keySet()) 로 마지막 오프셋을 가져오고
                //   ② 파티션마다 lag 을 계산해 MultiGauge.Row 를 만들고
                //   ③ lagGauge.register(rows, true) 로 등록하고
                //   ④ 총합을 log.info 로 남긴다

                List<MultiGauge.Row<?>> rows = new ArrayList<>();
                log.info("group={} rows={}", GROUP, rows.size());

            } catch (Exception ex) {
                log.warn("랙 수집 실패: {}", ex.toString());
            }
        }

        @PreDestroy
        public void close() {
            probe.close();
            admin.close();
        }
    }

    /* =======================================================================
     * 문제 3. 리스너 처리 시간 타이머 추가하기
     * -----------------------------------------------------------------------
     * 요구사항
     *   - 아래 리스너의 비즈니스 구간(deduct 호출)만 Timer.Sample 로 측정한다.
     *   - 지표 이름은 "order.inventory.deduct", 태그는 outcome=ok|error.
     *   - 예외가 나도 타이머가 기록되도록 finally 에서 stop 한다.
     *   - 그런 다음 아래 두 값을 비교해 프레임워크 오버헤드를 계산한다.
     *       curl -s localhost:8080/actuator/metrics/spring.kafka.listener | jq '.measurements'
     *       curl -s localhost:8080/actuator/metrics/order.inventory.deduct | jq '.measurements'
     *
     * // 관측 기록:
     * // spring.kafka.listener      COUNT=____  TOTAL_TIME=____  MAX=____
     * // order.inventory.deduct     COUNT=____  TOTAL_TIME=____  MAX=____
     * // 프레임워크 오버헤드(초)    = ____
     * // 두 COUNT 가 다르다면 그 이유는?
     *
     * 확인: order.inventory.deduct 의 availableTags 에 outcome 이 보여야 합니다.
     * ===================================================================== */
    @Component
    @Profile("step12-ex")
    public static class Ex3TimedListener {

        private static final Logger log = LoggerFactory.getLogger(Ex3TimedListener.class);

        private final MeterRegistry registry;

        public Ex3TimedListener(MeterRegistry registry) {
            this.registry = registry;
        }

        @KafkaListener(id = "s12-ex-inventory", topics = "orders", groupId = "s12-ex-inventory")
        public void onOrder(OrderCreated order,
                            @Header(KafkaHeaders.OFFSET) long offset) {
            // 여기에 작성: Timer.Sample 로 deduct(order) 구간만 측정하세요.
            deduct(order);

            if (offset % 200 == 0) {
                log.debug("재고 차감 orderId={} offset={}", order.orderId(), offset);
            }
        }

        private void deduct(OrderCreated order) {
            try {
                Thread.sleep(20L + (order.quantity() * 3L));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            if (order.orderId().endsWith("77")) {
                throw new IllegalStateException("재고 부족: " + order.orderId());
            }
        }
    }

    /* =======================================================================
     * 문제 4. observation 을 켜고 trace 연결 확인하기
     * -----------------------------------------------------------------------
     * 요구사항
     *   - application-step12-ex.yml 의 아래 두 값을 바꿔 가며 4가지 조합을 시험한다.
     *       spring.kafka.template.observation-enabled
     *       spring.kafka.listener.observation-enabled
     *   - 조합마다 POST /orders?seq=7 을 치고 로그의 traceId 를 기록한다.
     *   - 조합마다 kcc 로 레코드 헤더에 traceparent 가 있는지 확인한다.
     *       kcc --topic orders --partition 1 --offset <오프셋> --max-messages 1 --property print.headers=true
     *
     * // 관측 기록:
     * // ┌──────────┬──────────┬───────────────┬───────────────┬────────────────┐
     * // │ template │ listener │ HTTP traceId  │ 컨슈머 traceId │ traceparent 헤더│
     * // ├──────────┼──────────┼───────────────┼───────────────┼────────────────┤
     * // │ true     │ true     │               │               │                │
     * // │ true     │ false    │               │               │                │
     * // │ false    │ true     │               │               │                │
     * // │ false    │ false    │               │               │                │
     * // └──────────┴──────────┴───────────────┴───────────────┴────────────────┘
     * //
     * // trace 가 하나로 이어진 조합은?
     * // "헤더는 있는데 trace 가 안 이어진" 조합이 있었는가? 왜인가?
     *
     * 확인: 둘 다 true 일 때만 HTTP 와 컨슈머의 traceId 가 같아야 합니다.
     * ===================================================================== */

    /* =======================================================================
     * 문제 5. 커스텀 HealthIndicator 로 컨테이너 상태 노출하기
     * -----------------------------------------------------------------------
     * 요구사항
     *   - 빈 이름을 "exKafkaListeners" 로 두어 /actuator/health/exKafkaListeners 로 노출한다.
     *   - KafkaListenerEndpointRegistry.getListenerContainers() 를 순회하며
     *     isRunning() 과 getAssignedPartitions() 를 details 에 담는다.
     *   - 하나라도 running=false 이거나 할당 파티션이 비어 있으면 DOWN.
     *   - 컨테이너가 하나도 등록되지 않았다면 reason 을 담아 DOWN.
     *   - 만든 뒤 POST /admin/listeners/stop 을 치고 DOWN 으로 바뀌는지 확인한다.
     *   - 같은 시각 /actuator/health/kafka 는 무엇인지도 함께 기록한다.
     *
     * // 관측 기록:
     * // 정상 시   /actuator/health/kafka = ____   /actuator/health/exKafkaListeners = ____
     * // 정지 후   /actuator/health/kafka = ____   /actuator/health/exKafkaListeners = ____
     *
     * 확인: 정지 후 kafka 는 UP, exKafkaListeners 는 DOWN 이어야 합니다.
     * ===================================================================== */
    @Component("exKafkaListeners")
    @Profile("step12-ex")
    public static class Ex5HealthIndicator implements HealthIndicator {

        private final KafkaListenerEndpointRegistry registry;

        public Ex5HealthIndicator(KafkaListenerEndpointRegistry registry) {
            this.registry = registry;
        }

        @Override
        public Health health() {
            Collection<MessageListenerContainer> containers = registry.getListenerContainers();
            Map<String, Object> details = new LinkedHashMap<>();

            // 여기에 작성:
            //   ① containers 가 비었으면 down + reason
            //   ② 순회하며 running / assignedPartitions 를 details 에 담고
            //   ③ 하나라도 비정상이면 DOWN

            return Health.unknown().withDetails(details).build();
        }
    }

    /* =======================================================================
     * 문제 6. 리스너 로그에 partition@offset MDC 추가하기 (그리고 일부러 망가뜨리기)
     * -----------------------------------------------------------------------
     * 요구사항
     *   - RecordInterceptor 를 구현해 intercept 에서 MDC 에
     *       "kafka" = topic-partition@offset
     *     를 넣는다.
     *   - 전용 컨테이너 팩토리 "exMdcFactory" 에 붙이고, 아래 리스너에서 쓴다.
     *   - logging.pattern.level 에 [%X{kafka:-}] 를 추가해 로그를 확인한다.
     *   - ★ 그다음 afterRecord 의 MDC.remove 를 "빼고" 다시 돌린다.
     *     ORD-xx11 에서 예외가 나므로, 그 다음 레코드의 로그에
     *     이전 오프셋이 그대로 남는 것을 관측한다.
     *   - 관측이 끝나면 반드시 MDC.remove 를 원복한다.
     *
     * // 관측 기록 (remove 를 뺀 상태):
     * // 예외가 난 레코드의 실제 좌표    : orders-__@____
     * // 그 다음 레코드의 실제 좌표      : orders-__@____
     * // 그 다음 레코드 로그에 찍힌 좌표 : orders-__@____   ← 어긋났는가?
     *
     * 확인: remove 를 넣은 상태에서는 모든 로그의 [orders-N@M] 이
     *       그 로그를 만든 레코드의 좌표와 정확히 일치해야 합니다.
     * ===================================================================== */
    public static class Ex6MdcInterceptor implements RecordInterceptor<Object, Object> {

        @Override
        public ConsumerRecord<Object, Object> intercept(ConsumerRecord<Object, Object> record,
                                                        Consumer<Object, Object> consumer) {
            // 여기에 작성: MDC.put("kafka", ...)
            return record;
        }

        @Override
        public void afterRecord(ConsumerRecord<Object, Object> record, Consumer<Object, Object> consumer) {
            // 여기에 작성: MDC.remove("kafka")
            // ★ 실험할 때는 이 줄을 주석 처리하세요.
        }
    }

    @Configuration
    @Profile("step12-ex")
    public static class Ex6Config {

        @Bean("exMdcFactory")
        public ConcurrentKafkaListenerContainerFactory<Object, Object> exMdcFactory(ConsumerFactory<?, ?> cf) {
            ConcurrentKafkaListenerContainerFactory<Object, Object> factory =
                    new ConcurrentKafkaListenerContainerFactory<>();
            @SuppressWarnings("unchecked")
            ConsumerFactory<Object, Object> typed = (ConsumerFactory<Object, Object>) cf;
            factory.setConsumerFactory(typed);
            factory.setConcurrency(3);

            // 여기에 작성: factory.setRecordInterceptor(...)

            return factory;
        }
    }

    @Component
    @Profile("step12-ex")
    public static class Ex6Listener {

        private static final Logger log = LoggerFactory.getLogger(Ex6Listener.class);

        @KafkaListener(id = "s12-ex-mdc", topics = "orders",
                       groupId = "s12-ex-mdc", containerFactory = "exMdcFactory")
        public void onOrder(OrderCreated order) {
            if (order.orderId().endsWith("11")) {
                log.warn("재고 차감 실패 orderId={}", order.orderId());
                throw new IllegalStateException("재고 부족: " + order.orderId());
            }
            log.info("재고 차감 orderId={}", order.orderId());
        }
    }

    /** 문제 2 에서 쓰는 태그 헬퍼. 그대로 두세요. */
    static Tags lagTags(String group, TopicPartition tp) {
        return Tags.of("group", group, "topic", tp.topic(), "partition", String.valueOf(tp.partition()));
    }

    /** 문제 3 에서 쓰는 타이머 헬퍼. 그대로 두세요. */
    static Timer deductTimer(MeterRegistry registry, String outcome) {
        return registry.timer("order.inventory.deduct", "outcome", outcome);
    }
}

Solution.java

6문제의 정답 코드와, "왜 그 답인가"를 설명하는 긴 블록 주석이 함께 들어 있습니다. 풀어 본 뒤에 여세요.

  • 정답 1records.lagtopic/partition 라벨이 있고 records.lag.max 에는 없다는 차이를 정리합니다. 그래서 파티션 단위 알람은 records.lag 로만 가능하고, lag.maxclient_id 단위 요약에 쓴다는 결론까지 적었습니다.
  • 정답 2 의 포인트는 MultiGauge.register(rows, true)overwrite=true 입니다. false 로 두면 이미 등록된 태그 조합의 값이 갱신되지 않아 랙이 첫 측정값에서 얼어붙습니다. 값이 안 변하는데 에러도 없는 전형적인 조용한 버그라 실험으로 확인시킵니다.
  • 정답 3 은 두 타이머의 차이가 곧 프레임워크 오버헤드(역직렬화 + 인터셉터 + Observation)라는 계산과, COUNT 가 다를 수 있는 이유(재시도로 리스너가 여러 번 호출되면 spring.kafka.listener 의 count 만 늘어남)를 설명합니다.
  • 정답 4 는 네 가지 조합의 결과표입니다. 프로듀서만 켠 경우에도 컨슈머 로그의 traceId 가 비는 것이 핵심입니다. 헤더는 심겼는데 읽는 쪽이 없으므로, 헤더 유무와 trace 연결 여부가 별개라는 것을 보여 줍니다.
  • 정답 5isRunning() 만으로 부족한 이유를 설명합니다. 컨테이너는 running=true 인데 getAssignedPartitions() 가 비어 있는 상태(파티션을 전부 뺏긴 좀비)가 실제로 존재하므로, 두 조건을 && 로 묶어야 합니다. 다만 리밸런스 중 순간 DOWN 을 피하려고 알람은 연속 N회 조건을 걸라는 단서를 답니다.
  • 정답 6MDC.remove 누락의 재현 결과입니다. orders-1@311 에서 예외가 나고 다음 레코드 orders-1@312 의 로그에도 311 이 찍히는 로그를 그대로 실었습니다. afterRecord 가 성공·실패와 무관하게 호출되므로 여기가 유일하게 안전한 제거 위치라는 결론입니다.
package com.example.order.step12;

/*
 * ============================================================================
 * Step 12 — 관측성과 운영 : Solution (6문제 정답)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step12/Solution.java
 *
 * 실행
 *   ./gradlew bootRun --args='--spring.profiles.active=step12,step12-sol'
 *
 * ★ 실행 전 오프셋 리셋
 *   kcg --group s12-sol-inventory --topic orders --reset-offsets --to-earliest --execute
 *
 * 반드시 직접 풀어 본 뒤에 여세요.
 * ============================================================================
 */

import com.example.order.domain.OrderCreated;

import io.micrometer.core.instrument.MeterRegistry;
import io.micrometer.core.instrument.MultiGauge;
import io.micrometer.core.instrument.Tags;
import io.micrometer.core.instrument.Timer;

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.consumer.Consumer;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.boot.actuate.health.Health;
import org.springframework.boot.actuate.health.HealthIndicator;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.config.KafkaListenerEndpointRegistry;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.KafkaAdmin;
import org.springframework.kafka.listener.MessageListenerContainer;
import org.springframework.kafka.listener.RecordInterceptor;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;

import jakarta.annotation.PreDestroy;

import java.util.ArrayList;
import java.util.Collection;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;

public final class Solution {

    private Solution() { }

    /* =======================================================================
     * 정답 1. /actuator/prometheus 의 랙 지표 두 종류
     * =======================================================================
     *
     * $ curl -s localhost:8080/actuator/prometheus | grep -E '^kafka_consumer_fetch_manager_records_lag'
     *
     * kafka_consumer_fetch_manager_records_lag_max{application="order-service",
     *   client_id="consumer-s12-sol-inventory-2",kafka_version="3.6.1",
     *   spring_id="consumerFactory.s12-sol-inventory-1",} 31.0
     * kafka_consumer_fetch_manager_records_lag{application="order-service",
     *   client_id="consumer-s12-sol-inventory-2",partition="1",topic="orders",} 31.0
     *
     * ┌────────────────────────┬──────────────────────────────────────────────┐
     * │ 지표                    │ 라벨                                          │
     * ├────────────────────────┼──────────────────────────────────────────────┤
     * │ ..._records_lag        │ client_id, topic, partition                  │
     * │ ..._records_lag_max    │ client_id, spring_id, kafka_version (파티션 없음)│
     * └────────────────────────┴──────────────────────────────────────────────┘
     *
     * 파티션 단위 알람에는 records_lag(단수)를 씁니다. lag_max 는 파티션 라벨이
     * 없어서 "어느 파티션이 밀렸는지" 를 알 수 없고, 이름 그대로 그 컨슈머가
     * 담당한 파티션들 중 최댓값만 줍니다.
     *
     * 그러면 lag_max 는 왜 있는가? 파티션이 수백 개인 토픽에서 시계열 카디널리티를
     * 줄이려고 씁니다. 대시보드 상단의 "요약 한 줄" 용도입니다.
     *
     * ⚠️ 그리고 이 둘은 12-4 의 함정을 똑같이 갖고 있습니다. 컨슈머가 죽으면
     *    둘 다 사라집니다. 그래서 정답 2 의 게이지가 필요합니다.
     * ===================================================================== */

    /* =======================================================================
     * 정답 2. AdminClient + MultiGauge 로 랙 게이지 등록
     * =======================================================================
     *
     * 왜 이 답인가
     *
     * ① 커밋 오프셋은 AdminClient 로만 얻을 수 있습니다.
     *    컨슈머 객체의 position() 은 "그 컨슈머가 읽은 위치" 이지 "그룹이 커밋한
     *    위치" 가 아닙니다. 컨슈머가 죽어 있으면 position() 자체가 없습니다.
     *    listConsumerGroupOffsets 는 __consumer_offsets 를 읽으므로 컨슈머의
     *    생사와 무관합니다. 이것이 이 문제의 전부입니다.
     *
     * ② endOffsets 는 그룹 멤버십이 필요 없습니다.
     *    probe 컨슈머는 subscribe 도 assign 도 하지 않습니다. endOffsets 는
     *    ListOffsets API 를 직접 부르는 메타데이터 조회입니다.
     *    ⚠️ 여기서 실제 그룹 ID 로 subscribe 를 걸면 측정용 컨슈머가 그룹에
     *       합류해 리밸런스를 일으키고, 진짜 컨슈머에게서 파티션을 빼앗습니다.
     *       "관측하려다 장애를 만드는" 전형적인 사고입니다.
     *
     * ③ register(rows, true) 의 overwrite=true 가 핵심입니다.
     *    false 로 두면 이미 등록된 태그 조합(group/topic/partition)의 값이
     *    갱신되지 않습니다. 첫 측정값 그대로 얼어붙고, 예외도 경고도 없습니다.
     *
     *    실험 결과 (register(rows, false)):
     *      10:00:05  totalLag=0     → 게이지 0
     *      10:00:15  totalLag=284   → 게이지 여전히 0   ★
     *      10:00:25  totalLag=1176  → 게이지 여전히 0   ★
     *    로그의 totalLag 은 늘어나는데 /actuator/metrics 는 0 입니다.
     *    이 코스가 말하는 "에러 없이 조용히 틀린" 코드의 표본입니다.
     *
     * ④ Math.max(0, ...) 로 음수를 막습니다.
     *    커밋 오프셋과 endOffsets 를 서로 다른 시점에 읽으므로, 처리가 빠를 때
     *    아주 잠깐 음수가 나올 수 있습니다. 음수 랙은 대시보드를 망칩니다.
     *
     * ⑤ 예외를 삼키고 로그만 남깁니다.
     *    @Scheduled 메서드에서 예외를 던지면 스케줄러 로그에만 남고 다음 주기는
     *    계속 돕니다. 하지만 브로커 일시 장애 때마다 ERROR 스택트레이스가
     *    쏟아지므로 WARN 한 줄로 줄이는 편이 낫습니다.
     *
     * ⑥ 리스너를 멈춘 뒤 확인:
     *    $ curl -X POST localhost:8080/admin/listeners/stop && sleep 30
     *    $ curl -s localhost:8080/actuator/prometheus | grep -c records_lag   → 0
     *    $ curl -s localhost:8080/actuator/prometheus | grep -c group_lag     → 3  ★
     *    kafka_consumer_group_lag{...,partition="1",topic="orders",} 1284.0
     * ===================================================================== */
    @Component
    @Profile("step12-sol")
    public static class Sol2LagMetrics {

        private static final Logger log = LoggerFactory.getLogger(Sol2LagMetrics.class);
        private static final String GROUP = "s12-sol-inventory";

        private final AdminClient admin;
        private final Consumer<?, ?> probe;
        private final MultiGauge lagGauge;

        public Sol2LagMetrics(KafkaAdmin kafkaAdmin,
                              ConsumerFactory<?, ?> consumerFactory,
                              MeterRegistry registry) {
            this.admin = AdminClient.create(kafkaAdmin.getConfigurationProperties());
            this.probe = consumerFactory.createConsumer("s12-sol-lag-probe", "-probe");
            this.lagGauge = MultiGauge.builder("kafka.consumer.group.lag")
                    .description("committed offset 기준 파티션별 컨슈머 랙")
                    .baseUnit("records")
                    .register(registry);
        }

        @Scheduled(fixedDelay = 10_000L, initialDelay = 5_000L)
        public void refresh() {
            try {
                Map<TopicPartition, OffsetAndMetadata> committed =
                        admin.listConsumerGroupOffsets(GROUP)
                                .partitionsToOffsetAndMetadata()
                                .get(5, TimeUnit.SECONDS);

                if (committed.isEmpty()) {
                    lagGauge.register(List.of(), true);
                    log.info("group={} 커밋된 오프셋이 없습니다", GROUP);
                    return;
                }

                Map<TopicPartition, Long> endOffsets = probe.endOffsets(committed.keySet());

                List<MultiGauge.Row<?>> rows = new ArrayList<>();
                long total = 0L;
                for (Map.Entry<TopicPartition, OffsetAndMetadata> e : committed.entrySet()) {
                    TopicPartition tp = e.getKey();
                    long lag = Math.max(0L, endOffsets.getOrDefault(tp, 0L) - e.getValue().offset());
                    total += lag;
                    rows.add(MultiGauge.Row.of(
                            Tags.of("group", GROUP,
                                    "topic", tp.topic(),
                                    "partition", String.valueOf(tp.partition())),
                            lag));
                }
                lagGauge.register(rows, true);
                log.info("group={} totalLag={} partitions={}", GROUP, total, rows.size());

            } catch (Exception ex) {
                log.warn("랙 수집 실패 group={} : {}", GROUP, ex.toString());
            }
        }

        @PreDestroy
        public void close() {
            probe.close();
            admin.close();
        }
    }

    /* =======================================================================
     * 정답 3. Timer.Sample 로 비즈니스 구간만 측정
     * =======================================================================
     *
     * 왜 이 답인가
     *
     * ① finally 에서 stop 합니다.
     *    예외가 나도 시간이 기록되어야 "실패한 처리가 얼마나 오래 걸렸는지" 를
     *    볼 수 있습니다. 느려서 타임아웃으로 실패하는 경우가 가장 흔한데,
     *    catch 안에서만 stop 하면 그 데이터가 안 남습니다.
     *
     * ② outcome 태그를 예외 발생 시 "error" 로 바꾼 뒤 rethrow 합니다.
     *    rethrow 를 빼면 DefaultErrorHandler 가 예외를 못 받아 재시도도 DLT 도
     *    동작하지 않습니다(Step 07 의 함정과 같은 구조).
     *
     * ③ 실측 (10분 부하 후)
     *      spring.kafka.listener    COUNT=1842  TOTAL_TIME=41.196  MAX=1.084
     *      order.inventory.deduct   COUNT=1842  TOTAL_TIME=38.940  MAX=1.061
     *      프레임워크 오버헤드 = 41.196 - 38.940 = 2.256초 (건당 약 1.2ms)
     *
     *    2.256초의 정체는 역직렬화(JsonDeserializer) + RecordInterceptor +
     *    Observation 생성 + 메시지 변환입니다. 전체의 5.5% 입니다.
     *    "Kafka 가 느리다" 는 신고의 대부분은 이 5% 가 아니라 나머지 95% 입니다.
     *
     * ④ 두 COUNT 가 다를 수 있는 이유
     *    재시도가 걸리면 같은 레코드로 리스너가 여러 번 호출됩니다. 두 타이머
     *    모두 호출 횟수만큼 늘어나므로 보통은 같습니다. 하지만
     *    - 리스너 진입 전(역직렬화 단계)에서 실패하면 spring.kafka.listener 만 증가
     *    - deduct 호출 전에 return 하는 분기가 있으면 order.inventory.deduct 만 미증가
     *    두 경우에 어긋납니다. COUNT 가 다르면 TOTAL_TIME 비교는 무의미하므로
     *    반드시 먼저 확인해야 합니다.
     *
     * ⑤ @Timed 를 쓰지 않은 이유
     *    @Timed 는 TimedAspect 빈을 직접 등록해야 동작하고, @KafkaListener
     *    메서드는 프록시 대상이 되면 리스너 등록 자체가 꼬일 수 있습니다.
     *    리스너 내부에서는 Timer.Sample 이 정공법입니다.
     * ===================================================================== */
    @Component
    @Profile("step12-sol")
    public static class Sol3TimedListener {

        private static final Logger log = LoggerFactory.getLogger(Sol3TimedListener.class);

        private final MeterRegistry registry;

        public Sol3TimedListener(MeterRegistry registry) {
            this.registry = registry;
        }

        @KafkaListener(id = "s12-sol-inventory", topics = "orders", groupId = "s12-sol-inventory")
        public void onOrder(OrderCreated order,
                            @Header(KafkaHeaders.OFFSET) long offset) {
            Timer.Sample sample = Timer.start(registry);
            String outcome = "ok";
            try {
                deduct(order);
            } catch (RuntimeException ex) {
                outcome = "error";
                throw ex;                       // ★ 반드시 다시 던진다
            } finally {
                sample.stop(registry.timer("order.inventory.deduct", "outcome", outcome));
            }
            if (offset % 200 == 0) {
                log.debug("재고 차감 orderId={} offset={}", order.orderId(), offset);
            }
        }

        private void deduct(OrderCreated order) {
            try {
                Thread.sleep(20L + (order.quantity() * 3L));
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
            if (order.orderId().endsWith("77")) {
                throw new IllegalStateException("재고 부족: " + order.orderId());
            }
        }
    }

    /* =======================================================================
     * 정답 4. observation-enabled 4가지 조합
     * =======================================================================
     *
     * ┌──────────┬──────────┬──────────────────┬──────────────────┬─────────────┐
     * │ template │ listener │ HTTP traceId     │ 컨슈머 traceId    │ traceparent │
     * ├──────────┼──────────┼──────────────────┼──────────────────┼─────────────┤
     * │ true     │ true     │ 6f8c2a1b9d4e5f30 │ 6f8c2a1b9d4e5f30 │ 있음  ★연결 │
     * │ true     │ false    │ 6f8c2a1b9d4e5f30 │ (비어 있음)       │ 있음        │
     * │ false    │ true     │ 6f8c2a1b9d4e5f30 │ d41d8cd98f00b204 │ 없음        │
     * │ false    │ false    │ 6f8c2a1b9d4e5f30 │ (비어 있음)       │ 없음        │
     * └──────────┴──────────┴──────────────────┴──────────────────┴─────────────┘
     *
     * 왜 이런가
     *
     * ① 둘 다 true 일 때만 이어집니다. 프로듀서가 traceparent 를 "심고",
     *    컨슈머가 그것을 "읽어야" 하나의 trace 가 됩니다. 한쪽만으로는 안 됩니다.
     *
     * ② template=true / listener=false 는 이 문제의 핵심 함정입니다.
     *    헤더는 멀쩡히 심겼는데 읽는 쪽이 없어서 컨슈머 로그의 traceId 가
     *    아예 비어 있습니다. 즉 "헤더 유무" 와 "trace 연결" 은 별개입니다.
     *    헤더만 보고 "심겼으니 됐다" 고 판단하면 안 됩니다.
     *
     * ③ template=false / listener=true 는 컨슈머가 새 trace 를 시작합니다.
     *    traceId 가 비어 있지 않으므로 로그만 보면 정상처럼 보입니다.
     *    HTTP 요청의 traceId 와 대조해야만 끊긴 것을 알 수 있습니다.
     *    이것이 가장 발견하기 어려운 조합입니다.
     *
     * ④ 헤더 확인 명령
     *    kcc --topic orders --partition 1 --offset 318 --max-messages 1 \
     *        --property print.headers=true
     *    → traceparent:00-6f8c2a1b9d4e5f306f8c2a1b9d4e5f30-9c7d1e2f3a4b5c60-01
     *    끝의 -01 은 "샘플링됨" 입니다. -00 이면 수집기가 버립니다.
     *    management.tracing.sampling.probability 를 0.1 로 두면 -00 이 90% 입니다.
     *
     * ⑤ 설정을 켜기 전에 발행된 메시지에는 헤더가 없습니다. 랙이 큰 상태에서
     *    배포하면 한동안 끊긴 trace 가 대량으로 섞이는데, 장애가 아닙니다.
     * ===================================================================== */

    /* =======================================================================
     * 정답 5. 리스너 컨테이너 HealthIndicator
     * =======================================================================
     *
     * 왜 이 답인가
     *
     * ① isRunning() 만으로는 부족합니다.
     *    컨테이너는 running=true 인데 getAssignedPartitions() 가 비어 있는
     *    상태가 실제로 존재합니다. concurrency 를 파티션 수보다 크게 잡은 경우
     *    (Step 03), 또는 리밸런스로 파티션을 전부 뺏긴 "좀비" 컨테이너입니다.
     *    이 컨테이너는 살아 있지만 한 건도 처리하지 않습니다.
     *    그래서 running && !assigned.isEmpty() 두 조건을 && 로 묶습니다.
     *
     * ② getAssignedPartitions() 는 null 을 반환할 수 있습니다.
     *    MessageListenerContainer 의 기본 구현이 null 입니다. NPE 를 막으려면
     *    null 검사가 필요합니다. HealthIndicator 안에서 예외가 나면 Actuator 가
     *    500 을 내려 "헬스체크 자체가 장애" 가 됩니다.
     *
     * ③ 빈 이름이 곧 엔드포인트 경로입니다.
     *    @Component("exKafkaListeners") → /actuator/health/exKafkaListeners
     *    Boot 는 빈 이름에서 "HealthIndicator" 접미사만 떼고 그대로 씁니다.
     *
     * ④ 실측
     *      정상 시   /actuator/health/kafka = UP    /exKafkaListeners = UP
     *      정지 후   /actuator/health/kafka = UP ★  /exKafkaListeners = DOWN
     *    브로커는 멀쩡하니 KafkaHealthIndicator 는 끝까지 UP 입니다.
     *    이 UP 을 readiness probe 로 쓰면 아무 일도 안 하는 파드가 정상 판정을
     *    받고 계속 떠 있습니다. 그동안 랙 지표는 12-4 때문에 사라져 있습니다.
     *
     * ⑤ 알람은 연속 N회 DOWN 일 때만 울리세요.
     *    리밸런스가 진행되는 수 초 동안 assignedPartitions 가 비므로,
     *    단발 DOWN 으로 알람을 걸면 배포할 때마다 울립니다.
     *
     * ⑥ liveness 가 아니라 readiness 에 넣습니다.
     *    liveness 로 걸면 컨테이너 정지 → 파드 재시작 → 리밸런스 → 또 정지의
     *    루프가 생길 수 있습니다. 컨슈머는 HTTP 트래픽을 안 받으므로 readiness
     *    DOWN 의 실질 효과는 "알람" 이고, 그것으로 충분합니다.
     * ===================================================================== */
    @Component("exKafkaListeners")
    @Profile("step12-sol")
    public static class Sol5HealthIndicator implements HealthIndicator {

        private final KafkaListenerEndpointRegistry registry;

        public Sol5HealthIndicator(KafkaListenerEndpointRegistry registry) {
            this.registry = registry;
        }

        @Override
        public Health health() {
            Collection<MessageListenerContainer> containers = registry.getListenerContainers();
            if (containers.isEmpty()) {
                return Health.down().withDetail("reason", "등록된 리스너 컨테이너가 없음").build();
            }

            Map<String, Object> details = new LinkedHashMap<>();
            boolean healthy = true;

            for (MessageListenerContainer c : containers) {
                Collection<TopicPartition> assigned = c.getAssignedPartitions();
                boolean running = c.isRunning();
                boolean hasPartitions = assigned != null && !assigned.isEmpty();
                healthy = healthy && running && hasPartitions;

                details.put(c.getListenerId(), Map.of(
                        "running", running,
                        "assignedPartitions",
                        assigned == null ? List.of() : assigned.stream().map(TopicPartition::toString).toList()));
            }
            return (healthy ? Health.up() : Health.down()).withDetails(details).build();
        }
    }

    /* =======================================================================
     * 정답 6. RecordInterceptor 로 MDC 주입 — 그리고 remove 를 빼면 생기는 일
     * =======================================================================
     *
     * 왜 이 답인가
     *
     * ① intercept 는 리스너 호출 "직전" 에, afterRecord 는 "직후" 에 불립니다.
     *    afterRecord 는 리스너가 성공했든 예외를 던졌든 항상 호출됩니다.
     *    그래서 MDC 제거는 여기가 유일하게 안전한 자리입니다.
     *    리스너 메서드의 finally 에 넣으면, 역직렬화 단계에서 실패해 리스너가
     *    아예 호출되지 않은 경우를 못 잡습니다.
     *
     * ② 리스너 코드를 한 줄도 안 고쳐도 됩니다. 이것이 인터셉터를 쓰는 이유입니다.
     *    로그 규약을 조직 전체에 강제하려면, 각 리스너에 로깅 코드를 넣게 하는
     *    대신 공용 컨테이너 팩토리에 인터셉터 하나를 다는 것이 확실합니다.
     *
     * ③ MDC.remove 를 뺀 상태의 실측 (ORD-0311 에서 예외)
     *
     *    WARN  [order-service,...] [orders-1@311] ... : 재고 차감 실패 orderId=ORD-0311
     *    ERROR [order-service,...] [orders-1@311] ... o.s.k.l.DefaultErrorHandler : Record in retry
     *    INFO  [order-service,...] [orders-1@311] ... : 재고 차감 orderId=ORD-0312   ★ 어긋남
     *
     *    마지막 줄의 실제 레코드는 orders-1@312 인데 로그에는 311 이 찍혔습니다.
     *    ORD-0312 를 조사하려고 orders-1@311 을 꺼내 보면 엉뚱한 메시지가 나옵니다.
     *
     *    ⚠️ 이것이 이 문제의 결론입니다. 틀린 관측 정보는 없는 것보다 나쁩니다.
     *       관측 정보가 없으면 다른 방법을 찾지만, 틀린 정보는 잘못된 방향으로
     *       사람을 몇 시간씩 끌고 갑니다.
     *
     * ④ 스레드가 재사용되기 때문에 생기는 문제입니다. 리스너 스레드
     *    ntainer#0-1-C-1 은 파티션 1의 모든 레코드를 순차 처리합니다.
     *    MDC 는 ThreadLocal 이므로 지우지 않으면 다음 레코드까지 살아남습니다.
     *
     * ⑤ 배치 리스너를 쓴다면 RecordInterceptor 대신 BatchInterceptor 를 씁니다.
     *    이때 MDC 값은 배치 전체의 범위로 넣는 것이 정확합니다.
     *      "orders-1@300..349"
     *    배치 안의 개별 레코드를 특정해야 한다면 리스너 안에서 루프를 돌며
     *    직접 MDC 를 갱신해야 합니다.
     *
     * ⑥ traceId 와 함께 쓰는 것이 최종형입니다.
     *    logging.pattern.level:
     *      "%5p [${spring.application.name:},%X{traceId:-},%X{spanId:-}] [%X{kafka:-}]"
     *    → INFO [order-service,6f8c2a1b9d4e5f30,b2c3d4e5f6a70819] [orders-1@311]
     *    trace 로 흐름을 보고, 좌표로 원본 메시지를 꺼냅니다. 이 둘이면 됩니다.
     * ===================================================================== */
    public static class Sol6MdcInterceptor implements RecordInterceptor<Object, Object> {

        @Override
        public ConsumerRecord<Object, Object> intercept(ConsumerRecord<Object, Object> record,
                                                        Consumer<Object, Object> consumer) {
            MDC.put("kafka", record.topic() + "-" + record.partition() + "@" + record.offset());
            MDC.put("kafkaKey", String.valueOf(record.key()));
            return record;
        }

        @Override
        public void afterRecord(ConsumerRecord<Object, Object> record, Consumer<Object, Object> consumer) {
            // ★ 성공/실패와 무관하게 호출된다. 여기가 유일하게 안전한 제거 위치.
            MDC.remove("kafka");
            MDC.remove("kafkaKey");
        }
    }

    @Configuration
    @Profile("step12-sol")
    public static class Sol6Config {

        @Bean("solMdcFactory")
        public ConcurrentKafkaListenerContainerFactory<Object, Object> solMdcFactory(ConsumerFactory<?, ?> cf) {
            ConcurrentKafkaListenerContainerFactory<Object, Object> factory =
                    new ConcurrentKafkaListenerContainerFactory<>();
            @SuppressWarnings("unchecked")
            ConsumerFactory<Object, Object> typed = (ConsumerFactory<Object, Object>) cf;
            factory.setConsumerFactory(typed);
            factory.setConcurrency(3);
            factory.setRecordInterceptor(new Sol6MdcInterceptor());
            return factory;
        }
    }

    @Component
    @Profile("step12-sol")
    public static class Sol6Listener {

        private static final Logger log = LoggerFactory.getLogger(Sol6Listener.class);

        @KafkaListener(id = "s12-sol-mdc", topics = "orders",
                       groupId = "s12-sol-mdc", containerFactory = "solMdcFactory")
        public void onOrder(OrderCreated order) {
            if (order.orderId().endsWith("11")) {
                log.warn("재고 차감 실패 orderId={}", order.orderId());
                throw new IllegalStateException("재고 부족: " + order.orderId());
            }
            log.info("재고 차감 orderId={}", order.orderId());
        }
    }
}