Step 05 — 메시지 변환과 헤더

학습 목표

  • @KafkaListener 메서드가 받을 수 있는 모든 파라미터 종류를 표로 정리하고, 각각 언제 쓰는지 판단한다
  • 파라미터를 누가 해석하는지(MessagingMessageListenerAdapterDefaultMessageHandlerMethodFactory) 호출 경로를 따라간다
  • Spring Kafka 3.0 에서 이름이 바뀐 KafkaHeaders 상수를 런타임 null 이 나기 전에 잡아낸다
  • @Headerrequired / defaultValue 로 옵셔널 헤더를 안전하게 받는다
  • Serializer 방식과 MessageConverter 방식의 변환 지점 차이를 실측 로그로 구분하고, @KafkaHandler 로 한 토픽의 여러 타입을 분기한다
  • traceId 헤더를 발행 → 소비 → MDC → 재발행까지 전파시키고, 로그 패턴 %X{traceId} 에 실제로 찍히는 것을 확인한다
  • 헤더가 byte 라서 생기는 [B@1a2b3c 사고를 재현하고 고친다

선행 스텝: Step 04 — 직렬화와 역직렬화 예상 소요: 90분


5-0. 실습 준비

Step 04 에서 값(payload)이 어떻게 바이트가 되고 다시 객체가 되는지를 봤습니다. 이 스텝은 그 바깥, 레코드에 함께 실려 오는 나머지 전부 — 키·파티션·오프셋·타임스탬프·헤더 — 를 다룹니다.

이 스텝에서만 쓰는 두 번째 이벤트 타입이 필요합니다. domain 패키지에 하나 추가하세요.

package com.example.order.domain;

import java.time.Instant;

public record OrderCancelled(String orderId, String reason, Instant cancelledAt) {
    public static OrderCancelled of(int seq) {
        return new OrderCancelled(
                "ORD-%04d".formatted(seq),
                seq % 2 == 0 ? "OUT_OF_STOCK" : "USER_REQUEST",
                Instant.parse("2025-01-01T00:00:00Z").plusSeconds(seq * 60L));
    }
}

spring.json.trusted.packages 가 이미 com.example.order.domain 이므로 설정은 건드릴 필요 없습니다. 컨슈머 그룹은 s05- 접두사를 씁니다.

./gradlew bootRun --args='--spring.profiles.active=step05'

결과

INFO 13405 --- [           main] c.e.o.OrderServiceApplication            : The following 1 profile is active: "step05"
INFO 13405 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : s05-inventory: partitions assigned: [orders-0, orders-1, orders-2]

5-1. 리스너 메서드가 받을 수 있는 파라미터 전부

@KafkaListener 메서드의 시그니처는 자유롭습니다. 필요한 것만 골라서 아무 순서로 선언하면 됩니다.

@KafkaListener(topics = "orders", groupId = "s05-inventory")
public void onOrder(
        @Payload OrderCreated event,
        @Header(KafkaHeaders.RECEIVED_KEY) String key,
        @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
        @Header(KafkaHeaders.OFFSET) long offset,
        @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
        @Header(KafkaHeaders.RECEIVED_TIMESTAMP) long timestamp) {

    log.info("topic={} p={} off={} key={} ts={} sku={}",
            topic, partition, offset, key, timestamp, event.sku());
}

결과

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$AllParams          : topic=orders p=1 off=0 key=ORD-0001 ts=1735689660000 sku=SKU-002
INFO 13405 --- [ntainer#0-1-C-1] c.e.o.step05.Practice$AllParams          : topic=orders p=0 off=0 key=ORD-0002 ts=1735689720000 sku=SKU-003

받을 수 있는 것 전체입니다.

파라미터타입언제 쓰는가
(애노테이션 없음)OrderCreated역직렬화된 페이로드가장 흔한 형태. 메타데이터가 필요 없을 때
@PayloadOrderCreated위와 동일파라미터가 2개 이상일 때 명시적으로 페이로드를 지목
@Header(RECEIVED_KEY)String레코드 키키 기반 라우팅·멱등 판정. 키가 null 이면 이 헤더도 없다(5-4)
@Header(RECEIVED_PARTITION)int파티션 번호순서 보장 디버깅, 파티션별 통계
@Header(OFFSET)long오프셋재처리 지점 기록, 중복 판정
@Header(RECEIVED_TOPIC)String토픽명한 리스너가 topics = {"orders", "orders.DLT"} 처럼 여러 토픽을 볼 때
@Header(RECEIVED_TIMESTAMP)longepoch millis지연(lag) 계산: System.currentTimeMillis() - ts
@Header(TIMESTAMP_TYPE)StringCreateTime/LogAppendTime타임스탬프의 출처 판별
@Header(GROUP_ID)String컨슈머 그룹공통 리스너를 여러 그룹이 공유할 때
@Header("traceId")String커스텀 헤더추적 ID 전파(5-8)
@HeadersMap<String,Object>헤더 전체이름을 미리 모를 때. ⚠️ byte 가 섞인다(5-9)
(애노테이션 없음)MessageHeaders헤더 전체(불변)@Headers Map 과 거의 같음. 읽기 전용 의미가 명확
(애노테이션 없음)ConsumerRecord<K,V>원본 레코드메타데이터 전부 + serializedValueSize 등(5-5)
(애노테이션 없음)Message<OrderCreated>Spring Messaging 메시지getPayload() + getHeaders(). 다른 Spring Messaging 컴포넌트로 넘길 때
(애노테이션 없음)Acknowledgment수동 커밋 핸들AckMode.MANUAL* 일 때만 non-null (Step 06)
(애노테이션 없음)Consumer<?,?>카프카 컨슈머pause()/seek()/position() 직접 호출 (Step 10)

💡 파라미터 개수에 성능 차이는 없습니다. 헤더는 이미 MessageHeaders 맵으로 만들어진 뒤 리플렉션으로 꺼내 쓰는 것이라, 5개를 받든 1개를 받든 브로커에서 더 가져오는 건 없습니다. 필요한 걸 다 받으세요.

⚠️ Acknowledgment 를 선언했는데 AckModeBATCH(기본값)면 null 이 주입됩니다. 예외가 아니라 null 입니다. ack.acknowledge() 를 부르는 순간 NullPointerException 이 납니다. Step 06 에서 다룹니다.


5-2. 파라미터 해석은 누가 하는가

@KafkaListener 는 Kafka 기능이 아니라 Spring Messaging 기능입니다. 경로를 따라가면 이렇습니다.

① KafkaListenerAnnotationBeanPostProcessor   @KafkaListener 메서드 → MethodKafkaListenerEndpoint 등록
   (이때 DefaultMessageHandlerMethodFactory 를 주입)

② KafkaListenerEndpointRegistry              엔드포인트마다 컨테이너 생성·기동

③ MessagingMessageListenerAdapter            ConsumerRecord ──RecordMessageConverter──▶ Message<?>
                                             (payload + KafkaHeaders.* + 커스텀 헤더)

④ InvocableHandlerMethod                     파라미터마다 리졸버를 순서대로 물어본다
     · HeaderMethodArgumentResolver    ← @Header
     · HeadersMethodArgumentResolver   ← @Headers, MessageHeaders
     · MessageMethodArgumentResolver   ← Message<?>
     · PayloadMethodArgumentResolver   ← @Payload, 그리고 "남는 것"

⑤ 여러분의 메서드

핵심은 ③ 과 ④ 사이입니다. ④ 는 Kafka 를 전혀 모릅니다. 이미 Message<?> 로 번역된 것만 봅니다. Kafka 의 파티션·오프셋이 @Header 로 꺼내지는 이유는, ③ 이 그것들을 kafka_receivedPartition 같은 헤더 이름으로 심어 두었기 때문입니다.

애노테이션 없는 파라미터가 페이로드가 되는 규칙

PayloadMethodArgumentResolver마지막 순서로 등록되고, useDefaultResolution = true 로 동작합니다. 즉 "앞의 리졸버가 아무도 못 가져간 파라미터"를 전부 페이로드로 간주합니다. 그래서 이런 것들이 성립합니다.

public void a(OrderCreated event) { }                                  // OK — 페이로드
public void b(@Payload OrderCreated event, @Header(OFFSET) long off) { } // OK — 명시
public void c(OrderCreated event, @Header(OFFSET) long off) { }          // OK — off 는 헤더가 가져감

ConsumerRecord / Message / Acknowledgment / Consumer앞쪽 리졸버 또는 어댑터가 특별 취급하므로 페이로드로 오해되지 않습니다.

⚠️ 함정 — 애노테이션 없는 파라미터가 둘이면 조용히 잘못된 쪽이 페이로드가 된다

public void bad(OrderCreated event, String key) { }   // key 에 무엇이 들어올까?

keyPayloadMethodArgumentResolver 가 가져갑니다. 결과는 둘 다 페이로드이며, String 으로 변환하려다 MethodArgumentNotValidException 또는 이해 불가능한 변환 실패가 납니다. 기동 시점에는 아무 경고도 없고, 첫 메시지가 올 때까지 아무 일도 안 일어납니다.

ERROR 13405 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : Error handler threw an exception
org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Caused by: org.springframework.messaging.converter.MessageConversionException: Cannot convert from [com.example.order.domain.OrderCreated] to [java.lang.String] for GenericMessage [payload=OrderCreated[orderId=ORD-0001, ...]]

규칙 하나로 예방됩니다: 파라미터가 2개 이상이면 페이로드에 반드시 @Payload 를 붙이세요.


5-3. ⚠️ KafkaHeaders 상수 이름이 3.0 에서 바뀌었다

이 스텝에서 가장 많이 시간을 잡아먹는 함정입니다. 인터넷의 Spring Kafka 예제는 압도적으로 2.x 기준입니다.

2.x 상수2.x 문자열 값3.x 상수3.x 문자열 값
RECEIVED_MESSAGE_KEYkafka_receivedMessageKeyRECEIVED_KEYkafka_receivedKey
RECEIVED_PARTITION_IDkafka_receivedPartitionIdRECEIVED_PARTITIONkafka_receivedPartition
MESSAGE_KEYkafka_messageKeyKEYkafka_key
PARTITION_IDkafka_partitionIdPARTITIONkafka_partition
RECEIVED_TOPICkafka_receivedTopicRECEIVED_TOPIC (동일)kafka_receivedTopic
OFFSETkafka_offsetOFFSET (동일)kafka_offset
RECEIVED_TIMESTAMPkafka_receivedTimestampRECEIVED_TIMESTAMP (동일)kafka_receivedTimestamp

상수 이름만 바뀐 게 아니라 문자열 값도 바뀌었습니다. 여기서 두 갈래로 갈립니다.

(A) 상수를 쓴 경우 — 컴파일 에러. 다행입니다.

@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key   // 2.x 예제를 복사
> Task :compileJava FAILED
Practice.java:88: error: cannot find symbol
        @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) String key,
                            ^
  symbol:   variable RECEIVED_MESSAGE_KEY   location: class KafkaHeaders

(B) 문자열을 직접 쓴 경우 — 런타임에 실패합니다. 3.x 어댑터는 그 이름의 헤더를 만들지 않습니다.

@Header("kafka_receivedMessageKey") String key          // 상수 대신 문자열
ERROR 13405 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : Error handler threw an exception
org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Caused by: org.springframework.messaging.MessageHandlingException: Missing header 'kafka_receivedMessageKey' for method parameter type [class java.lang.String]

여기까진 그래도 시끄럽습니다. 진짜 위험한 건 required = false 를 함께 복사해 온 경우입니다.

@Header(name = "kafka_receivedMessageKey", required = false) String key   // 영원히 null

예외도, 경고도, 로그도 없습니다. key모든 메시지에서 null 입니다. "키가 안 들어와요" 라는 버그로 며칠을 태우는 전형적인 경로입니다.

결과 (증상)

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$LegacyHeader       : key=null off=0 payload=OrderCreated
INFO 13405 --- [ntainer#0-1-C-1] c.e.o.step05.Practice$LegacyHeader       : key=null off=0 payload=OrderCreated

💡 실무 팁 — 헤더 이름은 절대 문자열로 쓰지 마세요. KafkaHeaders.* 상수만 쓰면 (A) 로 떨어져 빌드가 깨지고, 그게 최선의 결과입니다. 마이그레이션 시 grep -rn '"kafka_' src/ 한 번이면 (B) 를 전부 찾아낼 수 있습니다. 어떤 헤더가 실제로 들어오는지 헷갈리면 5-5 의 ConsumerRecord 덤프를 한 번 찍어 보세요. 추측보다 빠릅니다.


5-4. ⚠️ @Header 가 없을 때 — 조용히 null 인가, 예외인가

@Header 의 기본값은 required = true 입니다. 헤더가 없으면 예외입니다.

@KafkaListener(topics = "orders", groupId = "s05-trace")
public void onOrder(@Payload OrderCreated event,
                    @Header("traceId") String traceId) {   // required 기본 true
    log.info("traceId={} order={}", traceId, event.orderId());
}

traceId 헤더 없이 발행된 메시지가 하나라도 오면:

ERROR 13405 --- [ntainer#0-1-C-1] o.s.k.l.KafkaMessageListenerContainer    : Error handler threw an exception
org.springframework.kafka.listener.ListenerExecutionFailedException: Listener method could not be invoked with the incoming message
Endpoint handler details:
Method [public void com.example.order.step05.Practice$RequiredHeader.onOrder(java.lang.Object,java.lang.String)]
Bean [com.example.order.step05.Practice$RequiredHeader@2c9c4a1e]
	at org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter.invokeHandler(MessagingMessageListenerAdapter.java:376)
	... 15 common frames omitted
Caused by: org.springframework.messaging.MessageHandlingException: Missing header 'traceId' for method parameter type [class java.lang.String]
	at org.springframework.messaging.handler.annotation.support.HeaderMethodArgumentResolver.handleMissingValue(HeaderMethodArgumentResolver.java:100)
	... 17 common frames omitted
WARN 13405 --- [ntainer#0-1-C-1] o.s.k.l.DefaultErrorHandler              : Backoff none exhausted for orders-1@3

이 예외는 재시도해도 절대 성공하지 않습니다. 헤더가 없다는 사실은 재시도로 변하지 않습니다. DLT 가 없으면 DefaultErrorHandler 가 로그만 남기고 넘어가고(기본 FixedBackOff(0, 9) 소진 후 skip), 그 메시지는 처리되지 않은 채 오프셋만 전진합니다. 조용한 유실입니다.

세 가지 선택지가 있습니다.

선언헤더가 없을 때쓰는 상황
@Header("traceId") String t예외 (MessageHandlingException)헤더가 계약상 반드시 있어야 할 때
@Header(name="traceId", required=false) String tnull있으면 쓰고 없으면 넘어감
@Header(name="traceId", defaultValue="none") String t"none"null 체크를 없애고 싶을 때

결과 (required=false 로 고친 뒤. traceId 없이 발행된 ORD-0003 도 통과합니다)

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$OptionalHeader     : traceId=3f7a1c8e retry=0 key=ORD-0001
INFO 13405 --- [ntainer#0-2-C-1] c.e.o.step05.Practice$OptionalHeader     : traceId=- retry=0 key=ORD-0003

⚠️ 함정 — 원시 타입에 required=false 를 붙이면 다른 예외가 난다

@Header(name = "retryCount", required = false) int retryCount   // int, not Integer

헤더가 없으면 리졸버가 null 을 넘기고, 리플렉션 호출에서 IllegalArgumentException: argument type mismatch 가 납니다. 원인 메시지가 헤더와 아무 상관없어 보여서 추적이 어렵습니다. 옵셔널 헤더는 반드시 래퍼 타입(Integer, Long) 이나 defaultValue 와 함께 쓰세요.

💡 RECEIVED_KEY 도 옵셔널로 두는 게 안전합니다. 키 없이 발행된 레코드(send(topic, value))는 이 헤더 자체가 없습니다. 키를 항상 넣는 코드만 보다가, 다른 팀이 키 없이 한 건 넣는 순간 리스너 전체가 죽습니다.


5-5. ConsumerRecord 를 직접 받기

메타데이터를 전부 원본 그대로 보고 싶다면 ConsumerRecord 를 그냥 파라미터로 선언하면 됩니다.

@KafkaListener(topics = "orders", groupId = "s05-audit")
public void audit(ConsumerRecord<String, Object> record) {
    StringBuilder headers = new StringBuilder();
    for (Header h : record.headers()) {                                  // org.apache.kafka...Header
        headers.append(h.key()).append('=')
               .append(new String(h.value(), StandardCharsets.UTF_8)).append(' ');
    }
    log.info("{}-{}@{} key={} ts={}({}) keySize={} valSize={} headers=[{}]",
            record.topic(), record.partition(), record.offset(),
            record.key(), record.timestamp(), record.timestampType(),
            record.serializedKeySize(), record.serializedValueSize(), headers.toString().trim());
}

결과

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$Audit              : orders-1@0 key=ORD-0001 ts=1735689660000(CreateTime) keySize=8 valSize=152 headers=[__TypeId__=com.example.order.domain.OrderCreated traceId=3f7a1c8e spring_json_header_types={"traceId":"java.lang.String"}]

ConsumerRecord 로만 얻을 수 있는 정보들입니다.

메서드용도
serializedKeySize() / serializedValueSize()바이트 수메시지 크기 감시. 5-10 의 헤더 비용 계산
timestampType()CreateTime/LogAppendTime프로듀서 시각인지 브로커 수신 시각인지
headers()Headers (iterable)모든 헤더, 같은 이름 중복 포함
headers().lastHeader(k)Header or null같은 이름이 여럿일 때 마지막 것
leaderEpoch()Optional<Integer>리더 변경 추적(고급)

💡 @Headers Map 은 같은 이름의 헤더가 여럿이면 하나만 남깁니다. Kafka 헤더는 멀티맵입니다 — 같은 키를 여러 번 넣을 수 있습니다. @RetryableTopic 이 재시도 이력을 같은 이름으로 여러 번 쌓는 것이 대표적입니다(Step 08). 전부 봐야 한다면 ConsumerRecord.headers() 를 순회하는 것 말고는 방법이 없습니다.

언제 ConsumerRecord 를 쓰는가: 로깅·감사·재처리 도구처럼 레코드를 내용이 아니라 사물로 다루는 코드입니다. 반대로 비즈니스 로직 리스너는 ConsumerRecord 를 받지 마세요. 도메인 코드에 org.apache.kafka 임포트가 스며듭니다.


5-6. MessageConverterSerializer 와 무엇이 다른가

Step 04 에서는 JsonDeserializer 로 바이트를 객체로 만들었습니다. 그런데 5-2 의 ③ 단계에도 변환기가 하나 더 있습니다. RecordMessageConverter 입니다.

브로커의 바이트 ──① Deserializer──▶ ConsumerRecord<K,V> ──② RecordMessageConverter──▶ Message<?> ──▶ 리스너
                 (kafka-clients, Step 04)                  (Spring Messaging, 이 절)

둘 중 한 곳에서만 JSON→객체 변환을 하면 됩니다. 어느 쪽을 고르냐가 설계 선택입니다.

Serializer 방식 (Step 04)MessageConverter 방식
변환 위치① kafka-clients② Spring Messaging
컨슈머 설정value-deserializer: JsonDeserializervalue-deserializer: StringDeserializer
타입 결정spring.json.value.default.type 또는 __TypeId__ 헤더리스너 메서드의 파라미터 타입
한 팩토리로 여러 타입어렵다 (JsonDeserializer 인스턴스가 타입에 묶임)쉽다 (메서드마다 다른 타입)
@KafkaHandler 분기__TypeId__ 필요파라미터 타입으로 자연스럽게
실패 지점poll 직후. 포이즌 필이 파티션을 막음리스너 호출 시점. DefaultErrorHandler 가 정상 처리
Kafka 표준 컨슈머와 호환그대로그대로

구현체 네 가지입니다. 컨슈머의 value-deserializer 와 짝이 맞아야 합니다.

Converter짝이 되는 Deserializer비고
MessagingMessageConverter아무거나기본값. 페이로드를 변환하지 않고 그대로 전달(헤더 매핑만)
StringJsonMessageConverterStringDeserializer가장 흔함
ByteArrayJsonMessageConverterByteArrayDeserializer문자열 중간 변환을 생략해 조금 빠름
BytesJsonMessageConverterBytesDeserializerorg.apache.kafka.common.utils.Bytes
@Bean
ConcurrentKafkaListenerContainerFactory<String, String> jsonConverterFactory(
        ConsumerFactory<String, String> cf) {

    var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
    factory.setConsumerFactory(cf);
    factory.setRecordMessageConverter(new StringJsonMessageConverter());   // ← 여기
    return factory;
}

@KafkaListener(topics = "orders", groupId = "s05-conv",
               containerFactory = "jsonConverterFactory")
public void onOrder(OrderCreated event) {          // 타입은 이 시그니처가 정한다
    log.info("converted -> {}", event.orderId());
}

결과

INFO 13405 --- [ntainer#0-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : s05-conv: partitions assigned: [orders-0, orders-1, orders-2]
INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$ConverterListener  : converted -> ORD-0001 sku=SKU-002

⚠️ 함정 — Deserializer 와 Converter 를 둘 다 JSON 으로 두면 아무 일도 안 일어난 것처럼 보인다 value-deserializer: JsonDeserializer 를 그대로 둔 채 StringJsonMessageConverter 를 얹으면, ① 에서 이미 OrderCreated 가 되어 있으므로 ② 는 "변환할 게 없네" 하고 그냥 통과시킵니다. 예외도 로그도 없습니다. 그래서 컨버터를 붙였는데 5-7 의 @KafkaHandler 분기가 안 되는데도 원인을 못 찾습니다. 컨버터를 쓰기로 했으면 deserializer 를 반드시 StringDeserializer 로 내리세요.

spring.kafka.consumer.value-deserializer: org.apache.kafka.common.serialization.StringDeserializer

반대로 두 곳 모두 켰는데 deserializer 만 String 이면 정상 동작합니다. 켜진 곳이 정확히 하나여야 합니다.

💡 프로듀서 쪽에도 대칭으로 컨버터를 걸 수 있지만(KafkaTemplate.setMessagingConverter), send(topic, key, value) 는 컨버터를 타지 않고 serializer 로 직행합니다. 컨버터를 타는 건 send(Message<?>) 뿐입니다.


5-7. @KafkaHandler 와 클래스 레벨 @KafkaListener

한 토픽에 여러 이벤트 타입이 섞여 오는 경우가 있습니다(주문 도메인 이벤트를 orders 하나로 묶는 설계). 이때 if (payload instanceof ...) 를 쓰지 말고 클래스 레벨 @KafkaListener + @KafkaHandler 로 분기합니다.

@Component
@Profile("step05")
@KafkaListener(topics = "orders", groupId = "s05-multi",
               containerFactory = "jsonConverterFactory")   // ← 클래스에 붙인다
public static class MultiTypeListener {

    @KafkaHandler
    public void onCreated(OrderCreated event) {
        log.info("[created]   {} sku={}", event.orderId(), event.sku());
    }

    @KafkaHandler
    public void onCancelled(OrderCancelled event) {
        log.info("[cancelled] {} reason={}", event.orderId(), event.reason());
    }

    @KafkaHandler(isDefault = true)
    public void onUnknown(Object payload,
                          @Header(KafkaHeaders.RECEIVED_KEY) String key) {
        log.warn("[unknown]   key={} type={}", key, payload.getClass().getName());
    }
}

결과

INFO 13405 --- [ntainer#0-1-C-1] c.e.o.step05.Practice$MultiTypeListener  : [created]   ORD-0001 sku=SKU-002
INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$MultiTypeListener  : [cancelled] ORD-0002 reason=OUT_OF_STOCK
WARN  13405 --- [ntainer#0-1-C-1] c.e.o.step05.Practice$MultiTypeListener  : [unknown]   key=ORD-0004 type=java.util.LinkedHashMap

마지막 줄이 흥미롭습니다. __TypeId__ 가 없는 레코드는 컨버터가 타입을 못 정해 LinkedHashMap 으로 남고, isDefault=true 핸들러가 받습니다.

⚠️ 함정 — isDefault 핸들러가 없으면 예외, 있으면 조용히 삼킨다 매칭되는 @KafkaHandler 가 없을 때:

Caused by: org.springframework.kafka.KafkaException: No method found for class java.util.LinkedHashMap

시끄럽습니다. 그런데 isDefault=true 핸들러를 만들어 두고 거기서 아무것도 안 하면, 모르는 타입이 전부 그리로 빨려 들어가 조용히 사라집니다. 기본 핸들러는 반드시 WARN 이상으로 로그를 남기거나 DLT 로 넘기세요. 삼키는 용도로 쓰면 안 됩니다.

⚠️ 클래스 레벨 @KafkaListener 를 쓰면서 같은 클래스에 메서드 레벨 @KafkaListener 를 함께 두면 안 됩니다. @KafkaHandler 가 하나도 없으면 기동 시 IllegalStateException: No @KafkaHandler methods found 로 즉시 실패합니다. 이건 다행인 쪽입니다.

💡 프로듀서에서 타입별로 다른 serializer 가 필요하면 DelegatingByTypeSerializer(Map<Class<?>, Serializer<?>>) 를 씁니다. 이 실습은 두 타입 모두 JSON 이라 기본 JsonSerializer 하나로 충분합니다.


5-8. 커스텀 헤더로 추적 ID 전파 (실전)

한 요청이 HTTP → orders → inventory → payments 로 흘러갈 때, 로그를 하나로 꿰려면 traceId 가 메시지를 따라가야 합니다. 페이로드에 넣지 마세요. 도메인 이벤트가 인프라 관심사로 오염됩니다. 헤더가 정답입니다.

프로듀서 — 헤더 붙이기 두 가지 방법

(A) ProducerRecord 에 직접

var record = new ProducerRecord<String, OrderCreated>("orders", event.orderId(), event);
record.headers().add("traceId", traceId.getBytes(StandardCharsets.UTF_8));
kafkaTemplate.send(record);

(B) Message<> 빌더 — 권장

Message<OrderCreated> message = MessageBuilder
        .withPayload(event)
        .setHeader(KafkaHeaders.TOPIC, "orders")
        .setHeader(KafkaHeaders.KEY, event.orderId())
        .setHeader("traceId", traceId)                 // ← String 그대로
        .build();
kafkaTemplate.send(message);

(B) 는 KafkaTemplateMessagingMessageConverterDefaultKafkaHeaderMapper 를 타므로, String 을 UTF-8 바이트로 바꾸고 spring_json_header_types 에 타입을 기록해 줍니다. 5-9 의 사고를 예방하는 것이 바로 이 부분입니다.

💡 (B) 는 KafkaHeaders.TOPIC/KEY 를 씁니다. RECEIVED_ 접두사가 붙은 것은 인바운드 전용입니다. 발행할 때 RECEIVED_KEY 를 쓰면 그냥 이름이 kafka_receivedKey 인 커스텀 헤더가 하나 생길 뿐, 레코드 키는 null 입니다. 조용한 실패입니다.

컨슈머 — MDC 에 넣고 로그 패턴에 출력

@KafkaListener(topics = "orders", groupId = "s05-inventory")
public void onOrder(@Payload OrderCreated event,
                    @Header(name = "traceId", required = false) String traceId,
                    @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                    @Header(KafkaHeaders.OFFSET) long offset) {
    MDC.put("traceId", traceId != null ? traceId : "-");
    try {
        log.info("재고 차감 시작 {} qty={} ({}@{})", event.orderId(), event.quantity(), partition, offset);
        inventory.decrease(event.sku(), event.quantity());
        log.info("재고 차감 완료 {}", event.orderId());
    } finally {
        MDC.remove("traceId");   // ← 반드시 finally
    }
}

application.yml 에 로그 패턴을 추가합니다.

logging:
  pattern:
    console: "%clr(%5p) %clr(${PID:- }){magenta} --- [%15.15t] %clr([%X{traceId:-        }]){yellow} %-40.40logger{39} : %m%n"

결과

 INFO 13405 --- [ntainer#0-1-C-1] [3f7a1c8e] c.e.o.step05.Practice$Inventory          : 재고 차감 시작 ORD-0001 qty=2 (1@0)
 INFO 13405 --- [ntainer#0-1-C-1] [3f7a1c8e] c.e.o.step05.Practice$Inventory          : 재고 차감 완료 ORD-0001
 INFO 13405 --- [ntainer#0-2-C-1] [        ] c.e.o.step05.Practice$Inventory          : 재고 차감 시작 ORD-0003 qty=4 (2@0)

마지막 줄은 traceId 없이 발행된 메시지입니다. 대괄호가 비어 있는 것으로 "추적 ID 를 안 붙인 프로듀서가 있다" 를 즉시 알 수 있습니다.

⚠️ 함정 — MDC.remove() 를 빼먹으면 다음 메시지에 앞 메시지의 traceId 가 찍힌다 리스너 컨테이너 스레드(ntainer#0-1-C-1)는 재사용됩니다. MDC 는 ThreadLocal 입니다. traceId 가 없는 메시지가 오면 MDC.put 이 호출되지 않아 직전 메시지의 값이 그대로 남습니다. 로그는 완벽해 보이지만 다른 요청의 로그가 한 트레이스로 뭉칩니다. 장애 분석에서 가장 나쁜 종류의 오염입니다. finally { MDC.remove(...) } 또는 MDC.clear() 를 습관화하세요. 더 나은 방법은 5-8 의 RecordInterceptor 화입니다.

재발행 시 헤더 이어받기

orders 를 소비해 payments 로 다시 발행할 때, traceId 가 끊기면 추적이 거기서 멈춥니다.

@KafkaListener(topics = "orders", groupId = "s05-payment")
public void relay(ConsumerRecord<String, OrderCreated> in) {
    var out = new ProducerRecord<String, OrderCreated>("payments", in.key(), in.value());
    // 인프라 헤더는 이어받고, 타입 헤더는 새로 붙게 두는 것이 안전하다
    for (Header h : in.headers()) {
        if (h.key().equals("traceId") || h.key().startsWith("x-")) {
            out.headers().add(h);          // byte[] 그대로 복사 — 변환 불필요
        }
    }
    kafkaTemplate.send(out);
}

결과

 INFO 13405 --- [ntainer#0-0-C-1] [3f7a1c8e] c.e.o.step05.Practice$Relay              : payments 로 중계 ORD-0001

⚠️ 헤더를 통째로 복사하면 안 됩니다. __TypeId__spring_json_header_types 까지 따라오면, payments 토픽의 값 타입이 달라졌을 때 컨슈머가 옛 타입으로 역직렬화하려다 실패합니다. 화이트리스트로 이어받으세요. 실무에서는 traceId, x- 접두사, correlationId 정도면 충분합니다.


5-9. ⚠️ 함정 — 헤더 값은 byte[]

Kafka 헤더의 타입은 String key + byte[] value 입니다. 그 이상은 없습니다. 문자열도, 숫자도, JSON 도, 전부 여러분이 인코딩과 디코딩을 책임져야 합니다.

byte[] raw = record.headers().lastHeader("traceId").value();   // byte[]
String traceId = new String(raw, StandardCharsets.UTF_8);      // 직접 변환

문제는 @Headers Map<String,Object> 로 받았을 때입니다.

@KafkaListener(topics = "orders", groupId = "s05-headers")
public void onOrder(@Payload OrderCreated event, @Headers Map<String, Object> headers) {
    log.info("traceId={} type={}", headers.get("traceId"),
             headers.get("traceId") == null ? "-" : headers.get("traceId").getClass().getSimpleName());
}

결과 — 프로듀서가 5-8 의 (A) 방식(ProducerRecord.headers().add(...))으로 발행한 경우

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$RawHeaders         : traceId=[B@6f2c0754 type=byte[]

결과 — 프로듀서가 (B) 방식(MessageBuilder.setHeader)으로 발행한 경우

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$RawHeaders         : traceId=3f7a1c8e type=String

같은 리스너 코드인데 결과가 다릅니다. 갈림길은 DefaultKafkaHeaderMapper 입니다.

  • (B) 는 아웃바운드에서 spring_json_header_types 라는 동반 헤더{"traceId":"java.lang.String"} 을 함께 기록합니다. 인바운드에서 매퍼가 이 힌트를 보고 String 으로 복원합니다.
  • (A) 는 그 힌트가 없습니다. 매퍼는 타입을 모르니 byte[] 를 그대로 맵에 넣습니다.

더 나쁜 경우 — @Header String 에서 조용히 깨진다

@Header(name = "traceId", required = false) String traceId

힌트가 없는 (A) 메시지에서는 byte[]String 파라미터에 맞춰야 합니다. DefaultFormattingConversionServicebyte[] → String 전용 변환기가 없어서, 최후의 수단인 Object → String(즉 toString())으로 떨어집니다.

결과

 INFO 13405 --- [ntainer#0-1-C-1] [[B@6f2c0754] c.e.o.step05.Practice$Inventory          : 재고 차감 시작 ORD-0001 qty=2 (1@0)
 INFO 13405 --- [ntainer#0-0-C-1] [[B@1d3a4b91] c.e.o.step05.Practice$Inventory          : 재고 차감 시작 ORD-0002 qty=3 (0@0)

예외가 나지 않습니다. 로그에는 [B@6f2c0754 가 찍히고, 이 값은 JVM 재시작마다 달라지므로 트레이스로 아무 쓸모가 없습니다. "traceId 는 잘 들어오는데 검색이 안 되네요"의 정체입니다. 다른 언어(Go/Python) 프로듀서가 붙는 순간 100% 발생합니다.

해결 — 셋 중 하나

DefaultKafkaHeaderMapper.setRawMappedHeaders 로 "이 헤더는 String 으로 읽어라" 라고 알려 준다 (권장)

@Bean
ConcurrentKafkaListenerContainerFactory<String, Object> rawHeaderFactory(ConsumerFactory<String, Object> cf) {
    var mapper = new DefaultKafkaHeaderMapper();
    mapper.setRawMappedHeaders(Map.of("traceId", true, "x-source", true));  // true = 인바운드에서 String 변환

    var converter = new MessagingMessageConverter();
    converter.setHeaderMapper(mapper);

    var factory = new ConcurrentKafkaListenerContainerFactory<String, Object>();
    factory.setConsumerFactory(cf);
    factory.setRecordMessageConverter(converter);
    return factory;
}

결과 (프로듀서가 (A) 방식이어도)

INFO 13405 --- [ntainer#0-0-C-1] c.e.o.step05.Practice$RawHeaders         : traceId=3f7a1c8e type=String

byte[] 로 받아 직접 변환한다@Header(name="traceId", required=false) byte[] bytes 로 받고 new String(bytes, UTF_8). 가장 정직하고 오해가 없습니다.

③ 프로듀서를 (B) 방식으로 통일한다 — 우리 팀 코드만 있을 때는 제일 깔끔합니다. 다만 다른 언어 프로듀서는 통제할 수 없습니다.

💡 실무 팁 — 헤더 값은 UTF-8 문자열로만 넣는 것을 팀 규칙으로 삼으세요. 숫자도 String.valueOf(n).getBytes(UTF_8) 로 넣습니다. ByteBuffer.allocate(4).putInt(n) 같은 바이너리 인코딩은 엔디안·언어별 차이 때문에 다른 팀이 소비하는 순간 깨집니다. kafka-console-consumer --property print.headers=true 로 눈으로 읽을 수 있다는 점만으로도 UTF-8 문자열의 가치는 충분합니다.


5-10. 헤더 크기와 개수의 비용

헤더는 레코드마다 반복 전송됩니다. 압축이 걸려 있어도 헤더 이름이 매번 실립니다. 레코드 하나의 헤더 비용은 대략 varint(키 길이) + 키 바이트 + varint(값 길이) + 값 바이트 입니다.

OrderCreated 하나를 기준으로 실측했습니다.

구성헤더레코드 크기
헤더 없음 (spring.json.add.type.headers: false)0개209 B
__TypeId__1개256 B
+ traceId + spring_json_header_types3개431 B

orders 토픽에 10만 건을 넣고 세그먼트 크기를 재면:

docker exec -it learn-kafka du -sh /var/lib/kafka/data/orders-0

결과

7.1M	/var/lib/kafka/data/orders-0      # 헤더 없음
14.7M	/var/lib/kafka/data/orders-0      # 헤더 3개

7.1MB → 14.7MB. 페이로드는 그대로인데 저장량이 2.07배가 됐습니다. 네트워크·디스크·리텐션 비용이 전부 2배입니다. 헤더는 공짜가 아닙니다.

한계설정기본값초과 시
프로듀서 요청max.request.size1,048,576RecordTooLargeException (즉시, 브로커 가기 전)
브로커 수신message.max.bytes1,048,588RecordTooLargeException (브로커가 거절)
토픽별max.message.bytes브로커값 상속위와 동일
컨슈머 fetchmax.partition.fetch.bytes1,048,576작으면 그 레코드를 못 읽고 멈춤

헤더 크기는 이 한계에 전부 포함됩니다. 페이로드가 900KB 인데 헤더에 스택트레이스를 200KB 넣으면 초과합니다.

ERROR 13405 --- [ad | producer-1] o.s.k.c.KafkaTemplate                    : Failed to send
org.apache.kafka.common.errors.RecordTooLargeException: The message is 1160244 bytes when serialized which is larger than 1048576, which is the value of the max.request.size configuration.

⚠️ 함정 — DLT 헤더가 메시지를 한계 너머로 밀어 올린다 DeadLetterPublishingRecoverer(Step 07)는 원본 메시지에 예외 스택트레이스를 헤더로 통째로 붙입니다 (kafka_dlt-exception-stacktrace). 스택이 깊으면 수십 KB 입니다. 원본이 이미 900KB 였다면 DLT 발행이 RecordTooLargeException 으로 실패하고, 그러면 그 메시지는 DLT 에도 못 가고 사라집니다. 에러 처리기가 에러를 내는 상황이라 로그도 어지럽습니다. DeadLetterPublishingRecoverer.setMaxStackTraceLength(...) 로 잘라 두세요.

💡 실무 팁 — 안 쓰는 헤더부터 끄세요. spring.json.add.type.headers: false__TypeId__(약 47B)를 없앱니다. 컨슈머가 spring.json.value.default.type 이나 MessageConverter 로 타입을 정한다면 이 헤더는 순수한 낭비입니다. Step 04 에서 본 "패키지 결합" 문제까지 함께 사라집니다.

💡 헤더가 Kafka 레코드 배치 안에서 실제로 어떻게 인코딩되는지(varint, 레코드 배치 오버헤드)는 Kafka 코스 Step 02 를 참고하세요. 이 코스는 클라이언트 쪽 비용만 다룹니다.


정리

개념핵심
리스너 파라미터페이로드·@Header·@Headers·ConsumerRecord·Message·Acknowledgment·Consumer 를 자유 조합
파라미터 해석MessagingMessageListenerAdapterDefaultMessageHandlerMethodFactory → 리졸버 체인
애노테이션 없는 파라미터PayloadMethodArgumentResolver 가 마지막에 전부 가져감. 2개 이상이면 @Payload 필수
KafkaHeaders 3.0 변경RECEIVED_MESSAGE_KEYRECEIVED_KEY, RECEIVED_PARTITION_IDRECEIVED_PARTITION. 문자열 값도 바뀜
헤더 이름 하드코딩상수를 쓰면 컴파일 에러(다행), 문자열을 쓰면 런타임 null
@Header 기본값required=true → 없으면 MessageHandlingException. 옵셔널은 required=false + 래퍼 타입
ConsumerRecord유일하게 serializedValueSize·중복 헤더 전체·timestampType 을 볼 수 있음
Serializer vs Converter① kafka-clients vs ② Spring Messaging. 둘 중 정확히 한 곳에서만 변환
Converter 의 장점타입을 메서드 시그니처가 결정 → 한 팩토리로 여러 타입
@KafkaHandler클래스 레벨 @KafkaListener + 타입별 메서드. isDefault=true반드시 로그를 남길 것
traceId 전파MessageBuilder.setHeader 로 발행 → @Header 수신 → MDC → %X{traceId}. finally 로 remove
재발행헤더 통째 복사 금지. traceId·x- 만 화이트리스트로 이어받기
헤더는 byte[]@Headers Map 은 raw byte. @Header String[B@1a2b3c조용히 깨짐
해결setRawMappedHeaders(Map.of("traceId", true)) 또는 byte[] 로 받아 직접 디코딩
헤더 비용레코드마다 반복. 3개로 저장량 2.07배. max.request.size·message.max.bytes 에 포함

연습문제

Exercise.java 에 7문제가 있습니다. 정답은 Solution.java.

  1. 토픽·파티션·오프셋·키·타임스탬프를 한 줄로 남기는 리스너를 @Header 만으로 작성하기
  2. MessageBuilder 로 traceId 헤더를 붙여 발행하고, 컨슈머에서 MDC 로 전파한 뒤 finally 로 정리하기
  3. ConsumerRecord모든 헤더를 이름=값(UTF-8) 형태로 덤프하기 (중복 헤더 포함)
  4. 클래스 레벨 @KafkaListener + @KafkaHandlerOrderCreated / OrderCancelled 분기하고, isDefault 로 나머지 처리
  5. required=falsedefaultValue 로 옵셔널 헤더 두 개(traceId, retryCount)를 안전하게 받기
  6. [B@1a2b3c 사고를 재현한 뒤 setRawMappedHeaders 로 고치기
  7. StringJsonMessageConverter 를 쓰는 컨테이너 팩토리를 만들어 하나의 팩토리로 두 타입 받기

다음 단계

여기까지는 "메시지를 어떻게 꺼내 읽는가"였습니다. 그런데 읽은 다음이 진짜 문제입니다. 언제 "다 읽었다"고 브로커에 알릴 것인가 — 이 타이밍 하나가 메시지 유실과 중복을 가릅니다. AckMode 6종을 전부 비교하고, 배치 리스너의 중간이 실패했을 때 어디까지 커밋되는지를 오프셋으로 직접 확인합니다.

Step 06 — 오프셋 커밋과 AckMode


실습 파일

세 파일을 순서대로 씁니다. 먼저 Practice.java 를 프로필 step05 로 띄워 5-1 ~ 5-10 의 로그를 교재와 대조하고, 특히 5-9 의 [B@... 출력이 여러분 화면에도 나오는지 확인하세요(안 나오면 프로듀서가 (B) 방식으로 돌고 있는 것입니다). 그다음 Exercise.java 의 7문제를 풀고, Solution.java 로 채점합니다. 세 파일 모두 com.example.order.step05 패키지이며 @Profile 로 격리되어 있어 서로 간섭하지 않습니다.

Practice.java

본문 전 절의 예제를 절 번호 주석과 함께 담은 단일 파일입니다.

  • 파일 상단 주석에 네 개의 보조 프로필이 정리돼 있습니다. step05-raw(5-9 의 (A) 방식 프로듀서), step05-legacy(5-3 의 2.x 문자열 헤더), step05-required(5-4 의 예외 재현), step05-conv(5-6/5-7 의 컨버터 방식). 기본 step05 만 켜면 정상 경로만 돕니다.
  • [5-1] AllParams[5-5] Audit 는 같은 토픽을 다른 그룹(s05-inventory, s05-audit)으로 읽습니다. 두 리스너의 로그가 번갈아 찍히는 것이 정상입니다.
  • [5-6] jsonConverterFactory 빈은 StringDeserializer 로 재정의한 별도 ConsumerFactory 를 씁니다. 기본 팩토리를 건드리지 않습니다. 5-6 의 함정("둘 다 JSON")을 피하기 위해서입니다.
  • [5-8] Inventory 는 MDC 를 try/finally 로 감쌌습니다. finally 를 지우고 다시 돌려 보면 함정 블록의 "앞 메시지 traceId 가 남는" 현상을 직접 볼 수 있습니다. 꼭 한 번 해 보세요.
  • [5-9] RawHeaders@Headers Map@Header String동시에 받아 한 줄에 찍습니다. step05-raw 프로필과 함께 켜야 [B@... 가 나옵니다.
  • [5-10] SizeProbeserializedKeySize + serializedValueSize + 헤더 바이트 합 을 계산해 레코드 실제 크기를 로그로 남깁니다. 표의 209 B / 431 B 를 여러분 환경에서 재현하는 코드입니다.
package com.example.order.step05;

/*
 * ============================================================================
 * Step 05 — 메시지 변환과 헤더 : Practice
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step05/Practice.java
 *
 * 사전 준비
 *   spring-kafka-lab/src/main/java/com/example/order/domain/OrderCancelled.java 를 먼저 만드세요 (본문 5-0).
 *
 *       package com.example.order.domain;
 *       import java.time.Instant;
 *       public record OrderCancelled(String orderId, String reason, Instant cancelledAt) {
 *           public static OrderCancelled of(int seq) {
 *               return new OrderCancelled("ORD-%04d".formatted(seq),
 *                       seq % 2 == 0 ? "OUT_OF_STOCK" : "USER_REQUEST",
 *                       Instant.parse("2025-01-01T00:00:00Z").plusSeconds(seq * 60L));
 *           }
 *       }
 *
 * 실행 (기본 — 정상 경로만 돕니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step05'
 *
 * 보조 프로필 (함정 재현용. 기본 프로필과 함께 켭니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step05,step05-raw'
 *       → [5-9] ProducerRecord.headers().add(byte[]) 방식으로 발행. @Headers Map 에 byte[] 가 들어옵니다.
 *   ./gradlew bootRun --args='--spring.profiles.active=step05,step05-legacy'
 *       → [5-3] 2.x 문자열 헤더 이름("kafka_receivedMessageKey")을 쓰면 key 가 영원히 null.
 *   ./gradlew bootRun --args='--spring.profiles.active=step05,step05-required'
 *       → [5-4] required=true 인 @Header 가 없을 때 MessageHandlingException 을 봅니다.
 *   ./gradlew bootRun --args='--spring.profiles.active=step05,step05-conv'
 *       → [5-6][5-7] StringJsonMessageConverter + @KafkaHandler 타입 분기.
 *   ./gradlew bootRun --args='--spring.profiles.active=step05,step05-raw,step05-fixed'
 *       → [5-9] setRawMappedHeaders 로 고친 팩토리로 같은 메시지를 다시 읽습니다.
 *
 * 로그 패턴 (traceId 를 보려면 application.yml 에 추가)
 *   logging:
 *     pattern:
 *       console: "%clr(%5p) %clr(${PID:- }){magenta} --- [%15.15t] %clr([%X{traceId:-        }]){yellow} %-40.40logger{39} : %m%n"
 *
 * 확인할 CLI
 *   kcc --topic orders --from-beginning --property print.key=true --property print.headers=true
 *   kcg --describe --group s05-inventory
 *   kcg --group s05-inventory --topic orders --reset-offsets --to-earliest --execute   (앱 종료 후)
 *   docker exec -it learn-kafka du -sh /var/lib/kafka/data/orders-0
 * ============================================================================
 */

import com.example.order.domain.OrderCancelled;
import com.example.order.domain.OrderCreated;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.KafkaHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.DefaultKafkaHeaderMapper;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.converter.MessagingMessageConverter;
import org.springframework.kafka.support.converter.StringJsonMessageConverter;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.Map;
import java.util.UUID;

/**
 * Step 05 의 모든 예제를 담은 단일 파일입니다.
 * 각 nested static class 는 본문 절 번호와 1:1 로 대응합니다.
 *
 * 주의: org.apache.kafka.common.header.Header 와
 *       org.springframework.messaging.handler.annotation.Header 는 이름이 같습니다.
 *       이 파일은 전자를 import 하고, 후자는 완전한 이름으로 씁니다.
 */
public final class Practice {

    private Practice() {
        // 유틸리티 홀더. 인스턴스화하지 않습니다.
    }

    // ========================================================================
    // [5-0] 공통 — 결정적인 traceId 생성기
    // ========================================================================
    //
    // 실행할 때마다 값이 바뀌면 교재의 로그와 대조할 수 없습니다.
    // seq 로부터 항상 같은 8자리 hex 를 만듭니다. (ORD-0001 → 3f7a1c8e)
    //
    static final class TraceIds {
        private TraceIds() {
        }

        static String of(int seq) {
            long h = UUID.nameUUIDFromBytes(("ORD-%04d".formatted(seq)).getBytes(StandardCharsets.UTF_8))
                    .getMostSignificantBits();
            return "%08x".formatted((int) (h >>> 32));
        }
    }

    // ========================================================================
    // [5-8] 프로듀서 (기본) — MessageBuilder 로 traceId 를 붙여 발행한다
    // ========================================================================
    //
    // MessageBuilder 로 만든 Message<?> 를 send 하면 KafkaTemplate 의
    // MessagingMessageConverter → DefaultKafkaHeaderMapper 를 탑니다.
    // 그 결과 traceId 가 UTF-8 바이트로 실리고, 동반 헤더
    // spring_json_header_types = {"traceId":"java.lang.String"} 가 함께 붙습니다.
    // 이 동반 헤더 덕분에 컨슈머가 byte[] 가 아니라 String 으로 복원합니다. ([5-9] 참고)
    //
    @Component
    @Profile("step05 & !step05-raw")
    public static class Publisher implements ApplicationRunner {

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

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            for (int seq = 1; seq <= 4; seq++) {
                String orderId = "ORD-%04d".formatted(seq);
                String traceId = TraceIds.of(seq);

                Object payload = (seq == 2) ? OrderCancelled.of(seq) : OrderCreated.of(seq);

                // seq == 3 은 일부러 traceId 를 붙이지 않습니다. ([5-4] 옵셔널 헤더 확인용)
                var builder = MessageBuilder
                        .withPayload(payload)
                        .setHeader(KafkaHeaders.TOPIC, "orders")
                        .setHeader(KafkaHeaders.KEY, orderId);
                if (seq != 3) {
                    builder.setHeader("traceId", traceId);
                }

                Message<?> message = builder.build();
                template.send(message);
                log.info("발행 {} traceId={}", orderId, seq != 3 ? traceId : "(없음)");
            }
            template.flush();
        }
    }

    // ========================================================================
    // [5-9] 프로듀서 (raw) — ProducerRecord 에 byte[] 헤더를 직접 붙인다
    // ========================================================================
    //
    // 이쪽 경로는 KafkaTemplate 의 메시지 컨버터를 타지 않습니다.
    // 따라서 spring_json_header_types 동반 헤더가 없고, 컨슈머는 타입을 모릅니다.
    // 다른 언어(Go/Python) 프로듀서가 붙었을 때와 동일한 상황입니다.
    //
    @Component
    @Profile("step05-raw")
    public static class RawPublisher implements ApplicationRunner {

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

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            for (int seq = 1; seq <= 4; seq++) {
                String orderId = "ORD-%04d".formatted(seq);
                Object payload = (seq == 2) ? OrderCancelled.of(seq) : OrderCreated.of(seq);

                var record = new ProducerRecord<String, Object>("orders", orderId, payload);
                record.headers().add("traceId",
                        TraceIds.of(seq).getBytes(StandardCharsets.UTF_8));   // ← 날 byte[]
                record.headers().add("x-source", "raw-publisher".getBytes(StandardCharsets.UTF_8));

                template.send(record);
                log.info("raw 발행 {} (동반 타입 헤더 없음)", orderId);
            }
            template.flush();
        }
    }

    // ========================================================================
    // [5-1] 리스너가 받을 수 있는 파라미터 전부
    // ========================================================================
    //
    // 파라미터가 2개 이상이므로 페이로드에 @Payload 를 명시합니다. ([5-2] 의 규칙)
    // 파라미터를 많이 받는다고 브로커에서 더 가져오지 않습니다. 헤더는 이미
    // MessageHeaders 맵으로 만들어진 뒤 리플렉션으로 꺼내 쓰는 것뿐입니다.
    //
    @Component
    @Profile("step05")
    public static class AllParams {

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

        @KafkaListener(topics = "orders", groupId = "s05-allparams")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_KEY) String key,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.OFFSET) long offset,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_TIMESTAMP) long timestamp,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.GROUP_ID) String groupId) {

            log.info("topic={} p={} off={} key={} ts={} group={} payload={}",
                    topic, partition, offset, key, timestamp, groupId,
                    event.getClass().getSimpleName());
        }
    }

    // ========================================================================
    // [5-3] 함정 — 2.x 문자열 헤더 이름을 그대로 복사한 경우
    // ========================================================================
    //
    // 3.x 는 kafka_receivedKey 라는 이름으로 헤더를 심습니다.
    // kafka_receivedMessageKey 라는 헤더는 존재하지 않으므로:
    //   · required = true  → MessageHandlingException: Missing header '...'
    //   · required = false → 영원히 null. 예외도 경고도 없습니다.
    // 상수(KafkaHeaders.RECEIVED_KEY)만 썼다면 컴파일 에러로 즉시 잡혔을 문제입니다.
    //
    // 마이그레이션 점검:  grep -rn '"kafka_' src/
    //
    @Component
    @Profile("step05-legacy")
    public static class LegacyHeader {

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

        @KafkaListener(topics = "orders", groupId = "s05-legacy")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "kafka_receivedMessageKey", required = false) String key,   // ← 2.x 이름
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.OFFSET) long offset) {

            log.info("key={} off={} payload={}", key, offset, event.getClass().getSimpleName());
        }
    }

    // ========================================================================
    // [5-4] required = true (기본값) — 헤더가 없으면 예외
    // ========================================================================
    //
    // seq == 3 메시지에는 traceId 가 없습니다. 그 메시지에서만 터집니다.
    // 재시도해도 헤더가 생기지 않으므로 백오프는 무의미하고,
    // DefaultErrorHandler 가 재시도를 소진한 뒤 skip 하면
    // "처리되지 않았는데 오프셋만 전진" 하는 조용한 유실이 됩니다.
    //
    @Component
    @Profile("step05-required")
    public static class RequiredHeader {

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

        @KafkaListener(topics = "orders", groupId = "s05-required")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header("traceId") String traceId) {

            log.info("traceId={} payload={}", traceId, event.getClass().getSimpleName());
        }
    }

    // ========================================================================
    // [5-4] required = false / defaultValue — 옵셔널 헤더의 올바른 선언
    // ========================================================================
    //
    // ⚠️ retryCount 를 int 로 선언하고 required=false 를 붙이면,
    //    헤더가 없을 때 null 이 전달되어 리플렉션 호출에서
    //    IllegalArgumentException: argument type mismatch 가 납니다.
    //    래퍼 타입(Integer) 또는 defaultValue 를 반드시 함께 쓰세요.
    //
    @Component
    @Profile("step05")
    public static class OptionalHeader {

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

        @KafkaListener(topics = "orders", groupId = "s05-optional")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String traceId,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "retryCount", defaultValue = "0") Integer retryCount,
                @org.springframework.messaging.handler.annotation.Header(
                        name = KafkaHeaders.RECEIVED_KEY, required = false) String key) {

            log.info("traceId={} retry={} key={}", traceId == null ? "-" : traceId, retryCount, key);
        }
    }

    // ========================================================================
    // [5-5] ConsumerRecord 를 직접 받기 — 메타데이터 전부
    // ========================================================================
    //
    // ConsumerRecord 로만 얻을 수 있는 것:
    //   · serializedKeySize / serializedValueSize   (크기 감시)
    //   · timestampType                             (CreateTime vs LogAppendTime)
    //   · headers() 순회                            (Kafka 헤더는 멀티맵. @Headers Map 은 중복을 잃는다)
    //
    // 로깅·감사·재처리 도구에 쓰세요. 비즈니스 리스너에는 쓰지 마세요.
    // 도메인 코드에 org.apache.kafka 임포트가 스며듭니다.
    //
    @Component
    @Profile("step05")
    public static class Audit {

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

        @KafkaListener(topics = "orders", groupId = "s05-audit")
        public void audit(ConsumerRecord<String, Object> record) {
            StringBuilder headers = new StringBuilder();
            for (Header h : record.headers()) {
                headers.append(h.key()).append('=')
                        .append(new String(h.value(), StandardCharsets.UTF_8))
                        .append(' ');
            }
            log.info("{}-{}@{} key={} ts={}({}) keySize={} valSize={} headers=[{}]",
                    record.topic(), record.partition(), record.offset(),
                    record.key(), record.timestamp(), record.timestampType(),
                    record.serializedKeySize(), record.serializedValueSize(),
                    headers.toString().trim());
        }
    }

    // ========================================================================
    // [5-6] MessageConverter 방식 — Serializer 를 String 으로 내리고 컨버터가 변환
    // ========================================================================
    //
    // ⚠️ 가장 흔한 실수: value-deserializer 를 JsonDeserializer 로 둔 채
    //    StringJsonMessageConverter 만 얹는 것. ①에서 이미 객체가 되어 있으므로
    //    ②는 "변환할 게 없네" 하고 통과합니다. 예외도 로그도 없습니다.
    //    그래서 [5-7] 의 @KafkaHandler 분기가 안 되는데 원인을 못 찾습니다.
    //
    // 그래서 여기서는 이 팩토리 전용 ConsumerFactory 를 새로 만들고,
    // key/value deserializer 를 StringDeserializer 로 명시적으로 내립니다.
    // 기본 팩토리(application.yml)는 건드리지 않습니다.
    //
    @Configuration
    @Profile("step05-conv")
    public static class ConverterConfig {

        @Bean
        public ConsumerFactory<String, String> stringConsumerFactory(KafkaProperties properties) {
            Map<String, Object> props = new HashMap<>(properties.buildConsumerProperties());
            props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            // ErrorHandlingDeserializer 의 위임 설정은 여기선 불필요하므로 제거합니다.
            props.remove("spring.deserializer.value.delegate.class");
            props.remove("spring.deserializer.key.delegate.class");
            props.remove("spring.json.value.default.type");
            return new DefaultKafkaConsumerFactory<>(props);
        }

        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, String> jsonConverterFactory(
                ConsumerFactory<String, String> stringConsumerFactory) {

            var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
            factory.setConsumerFactory(stringConsumerFactory);
            // ② Spring Messaging 레이어에서 JSON → 객체. 타입은 메서드 시그니처가 정한다.
            factory.setRecordMessageConverter(new StringJsonMessageConverter());
            factory.setConcurrency(3);
            return factory;
        }
    }

    @Component
    @Profile("step05-conv")
    public static class ConverterListener {

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

        // 파라미터 타입이 곧 변환 대상 타입입니다. JsonDeserializer 설정을 바꾸지 않아도
        // 같은 팩토리로 다른 리스너가 다른 타입을 받을 수 있습니다.
        @KafkaListener(topics = "orders", groupId = "s05-conv",
                containerFactory = "jsonConverterFactory")
        public void onOrder(OrderCreated event) {
            log.info("converted -> {} sku={}", event.orderId(), event.sku());
        }
    }

    // ========================================================================
    // [5-7] @KafkaHandler — 한 토픽에 여러 이벤트 타입이 섞여 올 때
    // ========================================================================
    //
    // @KafkaListener 를 "클래스"에 붙이고, 타입별 메서드에 @KafkaHandler 를 붙입니다.
    // instanceof 분기보다 낫습니다: 타입이 늘어도 기존 메서드를 건드리지 않습니다.
    //
    // ⚠️ isDefault=true 핸들러는 "모르는 타입 전부"를 받습니다.
    //    거기서 아무것도 안 하면 메시지가 조용히 사라집니다.
    //    반드시 WARN 이상으로 로그를 남기거나 DLT 로 보내세요.
    //    반대로 isDefault 핸들러가 아예 없으면
    //    KafkaException: No method found for class ... 로 시끄럽게 실패합니다.
    //
    @Component
    @Profile("step05-conv")
    @KafkaListener(topics = "orders", groupId = "s05-multi",
            containerFactory = "jsonConverterFactory")
    public static class MultiTypeListener {

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

        @KafkaHandler
        public void onCreated(OrderCreated event) {
            log.info("[created]   {} sku={}", event.orderId(), event.sku());
        }

        @KafkaHandler
        public void onCancelled(OrderCancelled event) {
            log.info("[cancelled] {} reason={}", event.orderId(), event.reason());
        }

        @KafkaHandler(isDefault = true)
        public void onUnknown(Object payload,
                              @org.springframework.messaging.handler.annotation.Header(
                                      name = KafkaHeaders.RECEIVED_KEY, required = false) String key) {
            log.warn("[unknown]   key={} type={} value={}", key, payload.getClass().getName(), payload);
        }
    }

    // ========================================================================
    // [5-8] traceId → MDC → 로그 패턴 %X{traceId}
    // ========================================================================
    //
    // ⚠️ finally 의 MDC.remove 를 빼면 어떻게 되는지 꼭 한 번 지워서 실행해 보세요.
    //    리스너 컨테이너 스레드(ntainer#0-1-C-1)는 재사용되고 MDC 는 ThreadLocal 입니다.
    //    traceId 가 없는 메시지가 오면 MDC.put 이 호출되지 않아
    //    직전 메시지의 traceId 가 그대로 남습니다.
    //    로그는 완벽해 보이는데 서로 다른 요청이 한 트레이스로 뭉칩니다.
    //
    @Component
    @Profile("step05")
    public static class Inventory {

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

        @KafkaListener(topics = "orders", groupId = "s05-inventory")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String traceId,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.OFFSET) long offset) {

            MDC.put("traceId", traceId != null ? traceId : "-");
            try {
                if (event instanceof OrderCreated created) {
                    log.info("재고 차감 시작 {} qty={} ({}@{})",
                            created.orderId(), created.quantity(), partition, offset);
                    log.info("재고 차감 완료 {}", created.orderId());
                } else {
                    log.info("재고 대상 아님 {} ({}@{})", event.getClass().getSimpleName(), partition, offset);
                }
            } finally {
                MDC.remove("traceId");   // ← 반드시 finally
            }
        }
    }

    // ========================================================================
    // [5-8] 재발행 시 헤더 이어받기
    // ========================================================================
    //
    // ⚠️ in.headers() 를 통째로 복사하면 __TypeId__ 와 spring_json_header_types 까지 따라갑니다.
    //    payments 토픽의 값 타입이 다르면 컨슈머가 옛 타입으로 역직렬화하려다 실패합니다.
    //    traceId / x- 접두사 / correlationId 정도만 화이트리스트로 이어받으세요.
    //    byte[] 를 그대로 복사하므로 인코딩 변환이 전혀 필요 없습니다.
    //
    @Component
    @Profile("step05")
    public static class Relay {

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

        private final KafkaTemplate<String, Object> template;

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

        @KafkaListener(topics = "orders", groupId = "s05-payment")
        public void relay(ConsumerRecord<String, Object> in) {
            var out = new ProducerRecord<String, Object>("payments", in.key(), in.value());
            for (Header h : in.headers()) {
                if ("traceId".equals(h.key()) || "correlationId".equals(h.key()) || h.key().startsWith("x-")) {
                    out.headers().add(h);
                }
            }
            template.send(out);

            Header t = in.headers().lastHeader("traceId");
            MDC.put("traceId", t == null ? "-" : new String(t.value(), StandardCharsets.UTF_8));
            try {
                log.info("payments 로 중계 {}", in.key());
            } finally {
                MDC.remove("traceId");
            }
        }
    }

    // ========================================================================
    // [5-9] 함정 — 헤더 값은 byte[] 다
    // ========================================================================
    //
    // step05-raw 프로필과 함께 켜면 아래가 찍힙니다.
    //   mapValue=[B@6f2c0754 mapType=byte[] asString=[B@6f2c0754
    //
    // 기본 프로필(MessageBuilder 발행)이면 이렇게 찍힙니다.
    //   mapValue=3f7a1c8e mapType=String asString=3f7a1c8e
    //
    // 차이는 spring_json_header_types 동반 헤더의 유무입니다.
    // @Header String 이 [B@... 로 찍히는 이유:
    //   DefaultFormattingConversionService 에 byte[] → String 전용 변환기가 없어서
    //   최후 수단인 Object → String (즉 toString()) 으로 떨어집니다. 예외가 안 납니다.
    //   그 값은 JVM 재시작마다 달라지므로 트레이스로 아무 쓸모가 없습니다.
    //
    @Component
    @Profile("step05")
    public static class RawHeaders {

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

        @KafkaListener(topics = "orders", groupId = "s05-headers")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Headers Map<String, Object> headers,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String asString,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) byte[] asBytes) {

            Object v = headers.get("traceId");
            log.info("mapValue={} mapType={} asString={} asBytes={}",
                    v, v == null ? "-" : v.getClass().getSimpleName(), asString,
                    asBytes == null ? "-" : new String(asBytes, StandardCharsets.UTF_8));
        }
    }

    // ========================================================================
    // [5-9] 해결 — DefaultKafkaHeaderMapper.setRawMappedHeaders
    // ========================================================================
    //
    // Map<String, Boolean> 의 value 가 true 면 "인바운드에서 UTF-8 String 으로 변환" 입니다.
    // false 로 두면 byte[] 그대로입니다.
    // 동반 타입 헤더가 없는 외부 프로듀서(다른 언어)를 상대할 때의 정석 해법입니다.
    //
    @Configuration
    @Profile("step05-fixed")
    public static class RawHeaderMapperConfig {

        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, Object> rawHeaderFactory(
                ConsumerFactory<String, Object> consumerFactory) {

            var mapper = new DefaultKafkaHeaderMapper();
            mapper.setRawMappedHeaders(Map.of(
                    "traceId", true,
                    "x-source", true));

            var converter = new MessagingMessageConverter();
            converter.setHeaderMapper(mapper);

            var factory = new ConcurrentKafkaListenerContainerFactory<String, Object>();
            factory.setConsumerFactory(consumerFactory);
            factory.setRecordMessageConverter(converter);
            factory.setConcurrency(3);
            return factory;
        }
    }

    @Component
    @Profile("step05-fixed")
    public static class FixedRawHeaders {

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

        @KafkaListener(topics = "orders", groupId = "s05-headers-fixed",
                containerFactory = "rawHeaderFactory")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Headers Map<String, Object> headers,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String traceId) {

            Object v = headers.get("traceId");
            log.info("[fixed] mapValue={} mapType={} asString={}",
                    v, v == null ? "-" : v.getClass().getSimpleName(), traceId);
        }
    }

    // ========================================================================
    // [5-10] 헤더 비용 실측 — 레코드 하나의 실제 크기
    // ========================================================================
    //
    // 헤더는 매 레코드마다 반복 전송됩니다. 이름까지 매번 실립니다.
    // 헤더 하나의 비용 ≈ varint(키 길이) + 키 바이트 + varint(값 길이) + 값 바이트
    //
    // 아래 계산은 varint 를 무시한 하한값이지만, 표의 209 B / 431 B 를 재현하기에 충분합니다.
    // 줄이는 법: spring.json.add.type.headers: false 로 __TypeId__(약 47 B)를 없앤다.
    //
    @Component
    @Profile("step05")
    public static class SizeProbe {

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

        @KafkaListener(topics = "orders", groupId = "s05-size")
        public void probe(ConsumerRecord<String, Object> record) {
            int headerBytes = 0;
            var breakdown = new LinkedHashMap<String, Integer>();
            for (Header h : record.headers()) {
                int size = h.key().getBytes(StandardCharsets.UTF_8).length
                        + (h.value() == null ? 0 : h.value().length);
                headerBytes += size;
                breakdown.put(h.key(), size);
            }
            int keySize = Math.max(record.serializedKeySize(), 0);
            int valSize = Math.max(record.serializedValueSize(), 0);

            log.info("{}@{} key={}B value={}B headers={}B({}개) total≈{}B {}",
                    record.partition(), record.offset(), keySize, valSize,
                    headerBytes, breakdown.size(), keySize + valSize + headerBytes, breakdown);
        }
    }
}

Exercise.java

7문제의 문제지입니다. 각 문제는 컴파일되는 뼈대와 // 여기에 작성: 자리로 되어 있습니다.

  • 문제 1·3·5 는 리스너 시그니처를 채우는 문제, 문제 2·6·7 은 프로듀서와 설정까지 함께 손대야 하는 문제, 문제 4 는 클래스 구조 자체를 바꾸는 문제입니다.
  • 문제 6 은 일부러 깨진 코드를 먼저 실행해 보라고 요구합니다. [B@ 로 시작하는 로그가 나오는 것을 확인한 뒤에 고쳐야 학습이 됩니다. 바로 정답 코드를 쓰면 절반만 배웁니다.
  • 각 문제마다 기대 로그를 주석으로 붙여 두었습니다. 여러분 출력과 글자 단위로 맞을 필요는 없지만, traceId= 뒤에 [B@ 가 붙는지 아닌지 같은 판정 지점은 정확히 일치해야 합니다.
  • 문제 4·7 은 OrderCancelled 를 씁니다. 5-0 에서 만든 domain/OrderCancelled.java 가 없으면 컴파일되지 않습니다.
  • 컨슈머 그룹은 문제별로 s05-ex1 ~ s05-ex7 로 분리했습니다. 한 문제를 다시 풀 때 그 그룹만 오프셋을 리셋하면 됩니다.
package com.example.order.step05;

/*
 * ============================================================================
 * Step 05 — 메시지 변환과 헤더 : Exercise (7문제)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step05/Exercise.java
 *
 * 실행
 *   ./gradlew bootRun --args='--spring.profiles.active=step05-ex'
 *
 * 문제 6 만 따로 (일부러 깨진 코드를 먼저 보기 위한 프로필)
 *   ./gradlew bootRun --args='--spring.profiles.active=step05-ex,step05-ex6-broken'
 *   ./gradlew bootRun --args='--spring.profiles.active=step05-ex,step05-ex6-fixed'
 *
 * 사전 준비
 *   domain/OrderCancelled.java (본문 5-0) 가 있어야 문제 4·7 이 컴파일됩니다.
 *
 * 다시 풀 때 (앱 종료 후 해당 그룹만 리셋)
 *   kcg --group s05-ex1 --topic orders --reset-offsets --to-earliest --execute
 *
 * 로그 패턴 (문제 2 에 필요)
 *   logging.pattern.console: "%clr(%5p) ... [%X{traceId:-        }] ..."
 * ============================================================================
 */

import com.example.order.domain.OrderCancelled;
import com.example.order.domain.OrderCreated;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.messaging.handler.annotation.Headers;
import org.springframework.messaging.handler.annotation.Payload;
import org.springframework.stereotype.Component;

import java.util.Map;

public final class Exercise {

    private Exercise() {
    }

    // ========================================================================
    // 공통 — 문제들이 소비할 메시지를 발행합니다. 수정하지 마세요.
    // ========================================================================
    //
    // ORD-0001 : OrderCreated,   traceId 있음, retryCount 없음
    // ORD-0002 : OrderCancelled, traceId 있음
    // ORD-0003 : OrderCreated,   traceId 없음          ← 문제 5 의 판정 대상
    // ORD-0004 : OrderCreated,   traceId 있음, retryCount=2
    //
    @Component
    @Profile("step05-ex")
    public static class Fixture implements ApplicationRunner {

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            for (int seq = 1; seq <= 4; seq++) {
                String orderId = "ORD-%04d".formatted(seq);
                Object payload = (seq == 2) ? OrderCancelled.of(seq) : OrderCreated.of(seq);

                var builder = org.springframework.messaging.support.MessageBuilder
                        .withPayload(payload)
                        .setHeader(KafkaHeaders.TOPIC, "orders")
                        .setHeader(KafkaHeaders.KEY, orderId);
                if (seq != 3) {
                    builder.setHeader("traceId", "trace-%04d".formatted(seq));
                }
                if (seq == 4) {
                    builder.setHeader("retryCount", "2");
                }
                template.send(builder.build());
            }
            template.flush();
        }
    }

    // ========================================================================
    // 문제 1. 메타데이터를 한 줄로 남기는 리스너
    // ========================================================================
    //
    // 요구사항
    //   · groupId = "s05-ex1", topics = "orders"
    //   · ConsumerRecord 를 쓰지 말고 @Header 만으로 다음을 전부 받으세요:
    //       토픽 / 파티션 / 오프셋 / 키 / 타임스탬프
    //   · 페이로드는 Object 로 받고, 파라미터가 2개 이상이므로 @Payload 를 명시할 것
    //   · KafkaHeaders 상수를 쓸 것. 문자열("kafka_receivedPartitionId" 등) 금지
    //
    // 기대 로그
    //   topic=orders p=1 off=0 key=ORD-0001 ts=1735689660000 payload=OrderCreated
    //
    @Component
    @Profile("step05-ex")
    public static class Q1 {

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

        // 여기에 작성:
        // @KafkaListener(...)
        // public void onOrder(...) { ... }
    }

    // ========================================================================
    // 문제 2. traceId 를 MDC 로 전파하기
    // ========================================================================
    //
    // 요구사항
    //   · groupId = "s05-ex2"
    //   · traceId 헤더를 옵셔널로 받아 MDC 에 넣고, 처리 로그 2줄을 남긴 뒤 반드시 정리할 것
    //   · traceId 가 없으면 MDC 에 "-" 를 넣을 것
    //   · 정리를 어디에서 해야 하는지 그 이유를 주석으로 한 줄 적을 것
    //
    // 기대 로그 (logging.pattern.console 에 %X{traceId} 를 넣은 뒤)
    //    INFO ... [trace-0001] ...Q2 : 처리 시작 ORD-0001
    //    INFO ... [trace-0001] ...Q2 : 처리 완료 ORD-0001
    //    INFO ... [        ]   ...Q2 : 처리 시작 ORD-0003     ← 대괄호가 비어야 정답
    //
    // ⚠️ ORD-0003 줄에 trace-0002 가 찍힌다면 정리를 안 한 것입니다.
    //
    @Component
    @Profile("step05-ex")
    public static class Q2 {

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

        // 여기에 작성:
    }

    // ========================================================================
    // 문제 3. 모든 헤더를 덤프하기
    // ========================================================================
    //
    // 요구사항
    //   · groupId = "s05-ex3"
    //   · ConsumerRecord 를 받아 헤더 전체를 "이름=값(UTF-8)" 형태로 한 줄에 출력할 것
    //   · @Headers Map 을 쓰지 말 것. 왜 안 되는지 주석으로 한 줄 적을 것
    //   · 값이 null 인 헤더도 안전하게 처리할 것
    //
    // 기대 로그
    //   ORD-0001 headers=[__TypeId__=com.example.order.domain.OrderCreated, traceId=trace-0001,
    //                     spring_json_header_types={"traceId":"java.lang.String"}]
    //
    @Component
    @Profile("step05-ex")
    public static class Q3 {

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

        // 여기에 작성:
    }

    // ========================================================================
    // 문제 4. @KafkaHandler 로 타입 분기
    // ========================================================================
    //
    // 요구사항
    //   · 이 클래스에 클래스 레벨 @KafkaListener 를 붙이세요 (topics="orders", groupId="s05-ex4",
    //     containerFactory = "exJsonFactory" — 문제 7 에서 만들 팩토리입니다)
    //   · OrderCreated 용 @KafkaHandler, OrderCancelled 용 @KafkaHandler 를 각각 작성
    //   · @KafkaHandler(isDefault = true) 로 나머지를 받되, WARN 으로 타입과 키를 남길 것
    //     (조용히 삼키면 안 되는 이유를 주석으로 한 줄)
    //
    // 기대 로그
    //   [created]   ORD-0001 sku=SKU-002
    //   [cancelled] ORD-0002 reason=OUT_OF_STOCK
    //
    @Component
    @Profile("step05-ex")
    // 여기에 작성: 클래스 레벨 @KafkaListener
    public static class Q4 {

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

        // 여기에 작성: @KafkaHandler 3개
    }

    // ========================================================================
    // 문제 5. 옵셔널 헤더 두 개를 안전하게 받기
    // ========================================================================
    //
    // 요구사항
    //   · groupId = "s05-ex5"
    //   · traceId  : 없을 수 있음 → 없으면 null
    //   · retryCount : 없을 수 있음 → 없으면 0 (defaultValue 사용)
    //   · ⚠️ retryCount 를 int 로 선언하면 안 되는 이유를 주석으로 적을 것
    //
    // 기대 로그
    //   ORD-0001 traceId=trace-0001 retry=0
    //   ORD-0003 traceId=- retry=0
    //   ORD-0004 traceId=trace-0004 retry=2
    //
    @Component
    @Profile("step05-ex")
    public static class Q5 {

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

        // 여기에 작성:
    }

    // ========================================================================
    // 문제 6. [B@1a2b3c 사고를 재현하고 고치기
    // ========================================================================
    //
    // (6-a) 먼저 아래 BrokenPublisher/BrokenListener 를 그대로 실행해 보세요.
    //       ./gradlew bootRun --args='--spring.profiles.active=step05-ex,step05-ex6-broken'
    //       traceId 가 [B@ 로 시작하는 값으로 찍히는 것을 눈으로 확인해야 합니다.
    //       ⚠️ 확인하지 않고 정답부터 쓰면 절반만 배웁니다.
    //
    // (6-b) 그다음 FixConfig 에 컨테이너 팩토리를 만들어 고치세요.
    //       힌트: DefaultKafkaHeaderMapper 의 setRawMappedHeaders(Map<String, Boolean>)
    //             MessagingMessageConverter 의 setHeaderMapper(...)
    //             factory.setRecordMessageConverter(...)
    //       고친 뒤 FixedListener 가 traceId=trace-raw-1 처럼 찍혀야 정답입니다.
    //
    @Component
    @Profile("step05-ex6-broken | step05-ex6-fixed")
    public static class BrokenPublisher implements ApplicationRunner {

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            for (int seq = 1; seq <= 2; seq++) {
                var record = new org.apache.kafka.clients.producer.ProducerRecord<String, Object>(
                        "orders", "ORD-%04d".formatted(seq), OrderCreated.of(seq));
                // 날 byte[] 로 넣습니다. 동반 타입 헤더(spring_json_header_types)가 없습니다.
                record.headers().add("traceId",
                        ("trace-raw-" + seq).getBytes(java.nio.charset.StandardCharsets.UTF_8));
                template.send(record);
            }
            template.flush();
        }
    }

    @Component
    @Profile("step05-ex6-broken")
    public static class BrokenListener {

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

        @org.springframework.kafka.annotation.KafkaListener(
                topics = "orders", groupId = "s05-ex6-broken")
        public void onOrder(@Payload Object event,
                            @Headers Map<String, Object> headers,
                            @Header(name = "traceId", required = false) String traceId) {
            Object v = headers.get("traceId");
            log.info("mapType={} traceId={}", v == null ? "-" : v.getClass().getSimpleName(), traceId);
        }
    }

    @Configuration
    @Profile("step05-ex6-fixed")
    public static class FixConfig {

        // 여기에 작성: @Bean ConcurrentKafkaListenerContainerFactory<String, Object> exFixedFactory(...)
    }

    @Component
    @Profile("step05-ex6-fixed")
    public static class FixedListener {

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

        // 여기에 작성: containerFactory = "exFixedFactory" 로 같은 메시지를 다시 읽어
        //             traceId 가 String 으로 찍히는지 확인하는 리스너
    }

    // ========================================================================
    // 문제 7. StringJsonMessageConverter 로 하나의 팩토리에서 두 타입 받기
    // ========================================================================
    //
    // 요구사항
    //   · @Bean ConsumerFactory<String, String> exStringConsumerFactory(KafkaProperties)
    //       - KafkaProperties.buildConsumerProperties() 로 시작해서
    //         key/value deserializer 를 StringDeserializer 로 내릴 것
    //       - ⚠️ 이걸 안 하면 무슨 일이 생기는지 주석으로 한 줄 적을 것 (본문 5-6 의 함정)
    //   · @Bean ConcurrentKafkaListenerContainerFactory<String, String> exJsonFactory(...)
    //       - setRecordMessageConverter(new StringJsonMessageConverter())
    //   · 이 팩토리를 문제 4 의 Q4 가 씁니다.
    //
    @Configuration
    @Profile("step05-ex")
    public static class Q7Config {

        // 여기에 작성:
        // @Bean
        // public ConsumerFactory<String, String> exStringConsumerFactory(KafkaProperties properties) { ... }
        //
        // @Bean
        // public ConcurrentKafkaListenerContainerFactory<String, String> exJsonFactory(
        //         ConsumerFactory<String, String> exStringConsumerFactory) { ... }
    }
}

Solution.java

정답과 "왜 그 답인가"를 설명하는 긴 주석이 문제마다 붙어 있습니다.

  • 정답 1@Payload 를 굳이 붙입니다. 파라미터가 6개이므로 5-2 의 규칙상 명시가 안전하고, 나중에 누가 파라미터를 하나 더 넣어도 안 깨지기 때문입니다.
  • 정답 2 의 핵심은 MessageBuilder 쪽이 아니라 finally 입니다. 주석에서 컨테이너 스레드 재사용과 ThreadLocal 의 관계를 다시 설명하고, RecordInterceptor 로 옮기는 확장안을 덧붙였습니다.
  • 정답 3@Headers Map 이 아니라 ConsumerRecord.headers() 를 순회합니다. Kafka 헤더가 멀티맵이라 Map 으로 받으면 중복 헤더가 사라지기 때문이고, 이것이 문제 3 이 요구한 판정 지점입니다.
  • 정답 4isDefault=true 핸들러에서 log.warn + 헤더 덤프를 남깁니다. 조용히 삼키면 안 된다는 5-7 의 함정을 코드로 지킨 것입니다.
  • 정답 5retryCountint 가 아니라 Integer + defaultValue = "0" 으로 받습니다. 원시 타입 + required=falseargument type mismatch 를 피하는 유일한 방법이고, 주석에 그 예외의 스택 위치까지 적어 두었습니다.
  • 정답 6setRawMappedHeaders(Map.of("traceId", true)) 를 씁니다. true 가 "인바운드에서 String 으로 변환"이라는 뜻이라는 점, false 로 두면 byte 그대로라는 점을 대조해 설명합니다.
  • 정답 7StringDeserializer 로 별도 ConsumerFactory 를 만드는 부분이 진짜 답입니다. 팩토리에 컨버터만 얹고 deserializer 를 안 내리면 5-6 의 함정대로 아무 일도 안 일어납니다.
package com.example.order.step05;

/*
 * ============================================================================
 * Step 05 — 메시지 변환과 헤더 : Solution (7문제 정답 + 해설)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step05/Solution.java
 *
 * 실행
 *   ./gradlew bootRun --args='--spring.profiles.active=step05-sol'
 *   ./gradlew bootRun --args='--spring.profiles.active=step05-sol,step05-sol6'
 *
 * Exercise.java 와 컨슈머 그룹이 겹치지 않도록 s05-sol* 접두사를 씁니다.
 * 문제를 직접 풀어 본 "뒤에" 여세요.
 * ============================================================================
 */

import com.example.order.domain.OrderCancelled;
import com.example.order.domain.OrderCreated;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.header.Header;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.slf4j.MDC;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.boot.autoconfigure.kafka.KafkaProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.KafkaHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.DefaultKafkaHeaderMapper;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.kafka.support.converter.MessagingMessageConverter;
import org.springframework.kafka.support.converter.StringJsonMessageConverter;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;
import java.util.StringJoiner;

public final class Solution {

    private Solution() {
    }

    // ========================================================================
    // 공통 픽스처 (Exercise 와 동일)
    // ========================================================================
    @Component
    @Profile("step05-sol")
    public static class Fixture implements ApplicationRunner {

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            for (int seq = 1; seq <= 4; seq++) {
                String orderId = "ORD-%04d".formatted(seq);
                Object payload = (seq == 2) ? OrderCancelled.of(seq) : OrderCreated.of(seq);

                var builder = MessageBuilder.withPayload(payload)
                        .setHeader(KafkaHeaders.TOPIC, "orders")
                        .setHeader(KafkaHeaders.KEY, orderId);
                if (seq != 3) {
                    builder.setHeader("traceId", "trace-%04d".formatted(seq));
                }
                if (seq == 4) {
                    builder.setHeader("retryCount", "2");
                }
                template.send(builder.build());
            }
            template.flush();
        }
    }

    // ========================================================================
    // 정답 1 — 메타데이터를 @Header 만으로 전부 받기
    // ========================================================================
    /*
     * 왜 이 답인가
     *
     * 1) 상수를 씁니다. RECEIVED_PARTITION / RECEIVED_KEY 는 Spring Kafka 3.0 에서
     *    RECEIVED_PARTITION_ID / RECEIVED_MESSAGE_KEY 로부터 이름이 바뀐 것들입니다.
     *    상수를 쓰면 2.x 예제를 복사했을 때 "컴파일 에러"로 즉시 잡힙니다.
     *    문자열("kafka_receivedPartitionId")을 쓰면 3.x 에는 그런 헤더가 없으므로
     *    required=true 면 런타임 MessageHandlingException, required=false 면 "영원히 null" 입니다.
     *    후자는 예외도 경고도 없어서 가장 오래 걸립니다.
     *
     * 2) @Payload 를 굳이 붙였습니다. 파라미터가 6개이기 때문입니다.
     *    PayloadMethodArgumentResolver 는 "앞의 리졸버가 못 가져간 파라미터 전부"를 페이로드로
     *    간주합니다(useDefaultResolution=true). 지금은 나머지가 전부 @Header 라 문제가 없지만,
     *    나중에 누가 애노테이션 없는 파라미터를 하나 더 넣는 순간 "페이로드가 둘"이 되어
     *    MessageConversionException 이 납니다. @Payload 명시는 그 사고에 대한 보험입니다.
     *
     * 3) partition 은 int, offset/timestamp 는 long 입니다. 이 세 헤더는 항상 존재하므로
     *    원시 타입으로 받아도 안전합니다. (반면 traceId 나 RECEIVED_KEY 는 없을 수 있습니다 → 정답 5)
     */
    @Component
    @Profile("step05-sol")
    public static class S1 {

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

        @KafkaListener(topics = "orders", groupId = "s05-sol1")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.OFFSET) long offset,
                @org.springframework.messaging.handler.annotation.Header(
                        name = KafkaHeaders.RECEIVED_KEY, required = false) String key,
                @org.springframework.messaging.handler.annotation.Header(KafkaHeaders.RECEIVED_TIMESTAMP) long ts) {

            log.info("topic={} p={} off={} key={} ts={} payload={}",
                    topic, partition, offset, key, ts, event.getClass().getSimpleName());
        }
    }

    // ========================================================================
    // 정답 2 — traceId → MDC → %X{traceId}
    // ========================================================================
    /*
     * 왜 이 답인가
     *
     * 핵심은 MessageBuilder 쪽이 아니라 finally 입니다.
     *
     * 리스너 컨테이너의 컨슈머 스레드(ntainer#0-1-C-1)는 메시지마다 새로 만들어지지 않고
     * 계속 재사용됩니다. MDC 는 ThreadLocal 입니다. 두 사실을 곱하면 이렇게 됩니다:
     *
     *   ORD-0002 처리 → MDC.put("traceId", "trace-0002")  → 로그
     *   ORD-0003 처리 → traceId 헤더가 없다 → MDC 를 건드리지 않음
     *                 → 로그에 trace-0002 가 그대로 찍힌다   ← 오염
     *
     * 예외도 경고도 없고, 로그는 오히려 "완벽해 보입니다". 그래서 장애 분석 때
     * 서로 다른 요청이 한 트레이스로 뭉쳐 있는 것을 발견하고서야 알게 됩니다.
     *
     * 대책은 둘 중 하나입니다.
     *   (a) traceId 가 없어도 항상 MDC.put 으로 덮어쓴다 (아래 정답이 이 방식, "-" 를 넣음)
     *   (b) finally 에서 MDC.remove 로 반드시 지운다      (아래 정답이 함께 씀)
     * 둘 다 하는 것이 가장 안전합니다. 리스너가 예외로 빠져나가도 (b)가 지켜 줍니다.
     *
     * 확장: 리스너마다 이 코드를 복붙하는 대신 RecordInterceptor 로 올리면
     * intercept() 에서 put, success()/failure() 에서 remove 를 한 곳에 모을 수 있습니다.
     * (RecordInterceptor 는 Step 10 에서 다룹니다.)
     */
    @Component
    @Profile("step05-sol")
    public static class S2 {

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

        @KafkaListener(topics = "orders", groupId = "s05-sol2")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String traceId,
                @org.springframework.messaging.handler.annotation.Header(
                        name = KafkaHeaders.RECEIVED_KEY, required = false) String key) {

            MDC.put("traceId", traceId != null ? traceId : "-");   // (a) 항상 덮어쓴다
            try {
                log.info("처리 시작 {}", key);
                log.info("처리 완료 {}", key);
            } finally {
                MDC.remove("traceId");                              // (b) 항상 지운다
            }
        }
    }

    // ========================================================================
    // 정답 3 — ConsumerRecord.headers() 순회로 전체 덤프
    // ========================================================================
    /*
     * 왜 @Headers Map 이 아니라 ConsumerRecord 인가
     *
     * Kafka 헤더는 Map 이 아니라 "멀티맵" 입니다. 같은 이름의 헤더를 여러 번 넣을 수 있고,
     * 실제로 그렇게 쓰는 컴포넌트가 있습니다(@RetryableTopic 의 재시도 이력, Step 08).
     * @Headers Map<String,Object> 로 받으면 같은 이름의 헤더 중 하나만 남습니다.
     * "덤프" 가 목적이라면 잃어버린 헤더가 있다는 사실조차 모르게 됩니다.
     *
     * 그리고 값 디코딩은 우리가 합니다. Kafka 헤더 값의 타입은 byte[] 하나뿐입니다.
     * new String(h.value(), UTF_8) 이 없으면 [B@1a2b3c 가 찍힙니다(본문 5-9).
     * h.value() 는 null 일 수 있으므로(값 없는 헤더는 합법입니다) null 검사도 필요합니다.
     *
     * 부수적으로 serializedValueSize 도 함께 찍었습니다. 헤더 비용을 눈으로 보라는 뜻입니다
     * (본문 5-10: 헤더 3개면 저장량이 2배가 됩니다).
     */
    @Component
    @Profile("step05-sol")
    public static class S3 {

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

        @KafkaListener(topics = "orders", groupId = "s05-sol3")
        public void dump(ConsumerRecord<String, Object> record) {
            StringJoiner joiner = new StringJoiner(", ", "[", "]");
            for (Header h : record.headers()) {
                String value = (h.value() == null) ? "(null)" : new String(h.value(), StandardCharsets.UTF_8);
                joiner.add(h.key() + "=" + value);
            }
            log.info("{} valSize={}B headers={}", record.key(), record.serializedValueSize(), joiner);
        }
    }

    // ========================================================================
    // 정답 4 — 클래스 레벨 @KafkaListener + @KafkaHandler
    // ========================================================================
    /*
     * 왜 이 답인가
     *
     * 1) @KafkaListener 를 클래스에 붙이고 타입별 메서드에 @KafkaHandler 를 붙입니다.
     *    if (payload instanceof ...) 사슬보다 나은 이유는 확장 방향입니다.
     *    새 이벤트 타입이 생기면 메서드를 하나 추가할 뿐, 기존 메서드를 건드리지 않습니다.
     *
     * 2) containerFactory 는 정답 7 의 exJsonFactory(StringJsonMessageConverter) 입니다.
     *    타입 분기가 성립하려면 "타입이 파라미터 시그니처로 결정" 되어야 하는데,
     *    그건 MessageConverter 방식에서만 그렇습니다.
     *    value-deserializer 가 JsonDeserializer 인 채로 두면 ①단계에서 이미 타입이 고정되어
     *    분기가 의도대로 동작하지 않습니다(본문 5-6 의 함정).
     *
     * 3) isDefault=true 핸들러에서 반드시 로그를 남깁니다.
     *    이 핸들러가 없으면 KafkaException: No method found for class ... 로 시끄럽게 실패합니다.
     *    있는데 아무것도 안 하면 모르는 타입이 전부 조용히 사라집니다.
     *    "시끄러운 실패" 를 "조용한 유실" 로 바꾸는 것이 최악이므로, 기본 핸들러는
     *    WARN 이상 로그 + (실무라면) DLT 발행까지 하는 것이 원칙입니다.
     *    __TypeId__ 가 없는 레코드는 컨버터가 타입을 못 정해 LinkedHashMap 으로 넘어옵니다.
     */
    @Component
    @Profile("step05-sol")
    @KafkaListener(topics = "orders", groupId = "s05-sol4", containerFactory = "solJsonFactory")
    public static class S4 {

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

        @KafkaHandler
        public void onCreated(OrderCreated event) {
            log.info("[created]   {} sku={}", event.orderId(), event.sku());
        }

        @KafkaHandler
        public void onCancelled(OrderCancelled event) {
            log.info("[cancelled] {} reason={}", event.orderId(), event.reason());
        }

        @KafkaHandler(isDefault = true)
        public void onUnknown(Object payload,
                              @org.springframework.messaging.handler.annotation.Header(
                                      name = KafkaHeaders.RECEIVED_KEY, required = false) String key) {
            // 조용히 삼키면 안 됩니다. 여기 들어온 메시지는 아무도 처리하지 않았는데
            // 오프셋만 전진합니다 — 에러 없는 유실입니다.
            log.warn("[unknown]   key={} type={} value={}", key, payload.getClass().getName(), payload);
        }
    }

    // ========================================================================
    // 정답 5 — 옵셔널 헤더는 래퍼 타입 + defaultValue
    // ========================================================================
    /*
     * 왜 Integer 이고 왜 defaultValue 인가
     *
     * @Header(name = "retryCount", required = false) int retryCount   ← 이렇게 쓰면 안 됩니다.
     *
     * 헤더가 없을 때 HeaderMethodArgumentResolver 는 null 을 돌려줍니다.
     * 그 null 이 int 파라미터로 들어가는 순간, 리플렉션 호출(InvocableHandlerMethod.doInvoke →
     * Method.invoke) 에서 이렇게 터집니다.
     *
     *   Caused by: java.lang.IllegalArgumentException: argument type mismatch
     *       at java.base/jdk.internal.reflect.DirectMethodHandleAccessor.invoke(...)
     *       at org.springframework.messaging.handler.invocation.InvocableHandlerMethod.doInvoke(...)
     *
     * 메시지에 "retryCount" 도 "header" 도 없어서 원인을 찾기가 아주 어렵습니다.
     * 그래서 옵셔널 헤더는 (a) 래퍼 타입 Integer/Long 으로 받거나 (b) defaultValue 를 주거나,
     * 둘 다 하는 것이 안전합니다. 아래 정답은 둘 다 했습니다.
     *
     * defaultValue 는 문자열 "0" 입니다. 리졸버가 ConversionService 로 Integer 로 변환합니다.
     * 헤더 값도 문자열 "2" 로 들어오지만 같은 경로로 Integer 2 가 됩니다.
     *
     * traceId 는 String 이라 null 이 들어가도 예외가 안 나므로 required=false 만으로 충분합니다.
     * 다만 로그에서 null 을 "-" 로 바꿔 주는 편이 읽기 좋습니다.
     */
    @Component
    @Profile("step05-sol")
    public static class S5 {

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

        @KafkaListener(topics = "orders", groupId = "s05-sol5")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Header(
                        name = KafkaHeaders.RECEIVED_KEY, required = false) String key,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String traceId,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "retryCount", defaultValue = "0") Integer retryCount) {

            log.info("{} traceId={} retry={}", key, traceId == null ? "-" : traceId, retryCount);
        }
    }

    // ========================================================================
    // 정답 6 — setRawMappedHeaders 로 byte[] 헤더 고치기
    // ========================================================================
    /*
     * 왜 이 답인가
     *
     * 증상부터 정리합니다. 프로듀서가 ProducerRecord.headers().add("traceId", bytes) 로
     * 날 byte[] 를 넣으면, Spring 의 DefaultKafkaHeaderMapper 가 아웃바운드에서 붙여 주는
     * 동반 헤더 spring_json_header_types = {"traceId":"java.lang.String"} 가 없습니다.
     * 인바운드 매퍼는 타입 힌트가 없으니 그 헤더를 byte[] 그대로 MessageHeaders 에 넣습니다.
     *
     *   · @Headers Map 으로 받으면      → 값의 실제 타입이 byte[]
     *   · @Header String 으로 받으면    → byte[] → String 전용 변환기가 없어
     *                                     최후 수단인 Object→String(toString())으로 떨어져
     *                                     "[B@6f2c0754" 가 들어옵니다. 예외가 안 납니다.
     *
     * 이 값은 JVM 재시작마다 달라지므로 로그 검색에 아무 쓸모가 없습니다.
     * 다른 언어(Go/Python) 프로듀서가 붙는 순간 반드시 겪게 됩니다.
     *
     * 해결은 인바운드 매퍼에게 "이 헤더는 String 으로 읽어라" 라고 알려 주는 것입니다.
     *
     *   mapper.setRawMappedHeaders(Map.of("traceId", true, "x-source", true));
     *
     * Map 의 value 가 Boolean 이라는 점이 헷갈리기 쉽습니다.
     *   true  = 인바운드에서 UTF-8 String 으로 변환해서 넘긴다
     *   false = byte[] 그대로 넘긴다 (동반 타입 헤더도 만들지 않는다)
     * 즉 true 가 우리가 원하는 값입니다.
     *
     * 매퍼는 컨버터에 물리고, 컨버터는 컨테이너 팩토리에 물립니다.
     *   DefaultKafkaHeaderMapper → MessagingMessageConverter → factory.setRecordMessageConverter
     * MessagingMessageConverter 는 페이로드를 변환하지 않고 그대로 통과시키므로,
     * value-deserializer 설정(JsonDeserializer)은 그대로 두어도 됩니다.
     * 이 점이 정답 7 의 StringJsonMessageConverter 와 다릅니다.
     *
     * 대안 두 가지도 알아 두세요.
     *   · @Header(name="traceId", required=false) byte[] 로 받아 직접 new String(..., UTF_8)
     *     → 가장 정직하고 오해가 없습니다. 리스너가 몇 개 없으면 이쪽이 낫습니다.
     *   · 프로듀서를 MessageBuilder.setHeader 방식으로 통일
     *     → 우리 팀 코드만 있을 때는 제일 깔끔하지만, 외부 프로듀서는 통제할 수 없습니다.
     */
    @Component
    @Profile("step05-sol6")
    public static class S6Publisher implements ApplicationRunner {

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            for (int seq = 1; seq <= 2; seq++) {
                var record = new ProducerRecord<String, Object>(
                        "orders", "ORD-%04d".formatted(seq), OrderCreated.of(seq));
                record.headers().add("traceId",
                        ("trace-raw-" + seq).getBytes(StandardCharsets.UTF_8));
                template.send(record);
            }
            template.flush();
        }
    }

    @Configuration
    @Profile("step05-sol6")
    public static class S6Config {

        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, Object> solRawHeaderFactory(
                ConsumerFactory<String, Object> consumerFactory) {

            var mapper = new DefaultKafkaHeaderMapper();
            mapper.setRawMappedHeaders(Map.of(
                    "traceId", true,      // true = 인바운드에서 UTF-8 String 으로 변환
                    "x-source", true));

            var converter = new MessagingMessageConverter();
            converter.setHeaderMapper(mapper);

            var factory = new ConcurrentKafkaListenerContainerFactory<String, Object>();
            factory.setConsumerFactory(consumerFactory);
            factory.setRecordMessageConverter(converter);
            factory.setConcurrency(3);
            return factory;
        }
    }

    @Component
    @Profile("step05-sol6")
    public static class S6Listener {

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

        @KafkaListener(topics = "orders", groupId = "s05-sol6",
                containerFactory = "solRawHeaderFactory")
        public void onOrder(
                @org.springframework.messaging.handler.annotation.Payload Object event,
                @org.springframework.messaging.handler.annotation.Headers Map<String, Object> headers,
                @org.springframework.messaging.handler.annotation.Header(
                        name = "traceId", required = false) String traceId) {

            Object v = headers.get("traceId");
            log.info("[fixed] mapType={} traceId={}",
                    v == null ? "-" : v.getClass().getSimpleName(), traceId);
        }
    }

    // ========================================================================
    // 정답 7 — StringJsonMessageConverter 를 쓰는 컨테이너 팩토리
    // ========================================================================
    /*
     * 왜 ConsumerFactory 를 새로 만드는가 — 이게 이 문제의 진짜 답입니다.
     *
     * 컨버터만 얹고 value-deserializer 를 그대로(JsonDeserializer) 두면,
     * ① kafka-clients 레이어에서 이미 OrderCreated 객체가 만들어져 있으므로
     * ② StringJsonMessageConverter 는 "변환할 String 이 없네" 하고 그냥 통과시킵니다.
     * 예외도 로그도 없습니다. 겉으로는 잘 도는 것처럼 보이는데
     * 정답 4 의 @KafkaHandler 타입 분기가 의도대로 동작하지 않고, 원인을 찾을 단서가 없습니다.
     *
     * 그래서 이 팩토리 전용 ConsumerFactory 를 만들고 key/value 를 StringDeserializer 로 내립니다.
     * application.yml 의 기본 팩토리는 그대로 두므로 다른 리스너에 영향이 없습니다.
     *
     * KafkaProperties.buildConsumerProperties() 로 시작하는 이유는
     * bootstrap-servers / auto-offset-reset / enable-auto-commit 같은 공통 설정을
     * 손으로 다시 쓰지 않기 위해서입니다. 그 위에 필요한 것만 덮어씁니다.
     *
     * ErrorHandlingDeserializer 관련 위임 프로퍼티와 spring.json.value.default.type 은
     * StringDeserializer 에게 아무 의미가 없으므로 지웁니다. 남겨도 동작은 하지만
     * 기동 로그에 "isn't a known config" 성격의 잡음이 남습니다.
     *
     * 이 방식의 이점(본문 5-6 표):
     *   · 타입을 리스너 메서드 시그니처가 결정 → 한 팩토리로 여러 타입
     *   · 역직렬화 실패가 poll 이 아니라 리스너 호출 시점에 나므로
     *     DefaultErrorHandler 가 정상적으로 처리 → 포이즌 필이 파티션을 막지 않음
     */
    @Configuration
    @Profile("step05-sol")
    public static class S7Config {

        @Bean
        public ConsumerFactory<String, String> solStringConsumerFactory(KafkaProperties properties) {
            Map<String, Object> props = new HashMap<>(properties.buildConsumerProperties());
            props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
            props.remove("spring.deserializer.key.delegate.class");
            props.remove("spring.deserializer.value.delegate.class");
            props.remove("spring.json.value.default.type");
            return new DefaultKafkaConsumerFactory<>(props);
        }

        @Bean
        public ConcurrentKafkaListenerContainerFactory<String, String> solJsonFactory(
                ConsumerFactory<String, String> solStringConsumerFactory) {

            var factory = new ConcurrentKafkaListenerContainerFactory<String, String>();
            factory.setConsumerFactory(solStringConsumerFactory);
            factory.setRecordMessageConverter(new StringJsonMessageConverter());
            factory.setConcurrency(3);
            return factory;
        }
    }
}