Step 08 — @RetryableTopic 논블로킹 재시도

학습 목표

  • 블로킹 재시도가 파티션을 막는 문제를 재시도 전용 토픽으로 푸는 구조를 설명한다
  • @RetryableTopic 한 줄로 자동 생성되는 토픽 목록을 kt --list 로 확인한다
  • 실패 메시지 한 건이 ordersorders-retry-0-retry-1-retry-2orders.DLT 로 이동하는 전 과정을 로그로 추적한다
  • 같은 키 3건으로 순서가 깨지는 것을 재현하고 타임라인 표로 기록한다
  • @DltHandlerKafkaHeaders.ORIGINAL_* 헤더로 실패 원인을 복원한다
  • 백오프가 길면 retry 토픽 안에서 다시 블로킹된다는 것을 파티션 pause 로그로 확인한다
  • 블로킹/논블로킹 선택 기준표를 만들고, 두 가지를 조합하는 설정을 작성한다

선행 스텝: Step 07 — 에러 처리와 재시도 예상 소요: 90분


8-0. 실습 준비

이 스텝은 컨슈머 그룹 s08-inventory 를 씁니다. 토픽은 orders(파티션 3) 하나로 시작하고, 나머지는 Spring 이 만듭니다.

# 앱을 끈 상태에서
kcg --group s08-inventory --topic orders --reset-offsets --to-earliest --execute
kt --list

결과 — 아직은 깨끗합니다

orders
orders.DLT
payments

orders-retry-* 가 하나도 없는 것을 확인하세요. 8-2 에서 앱을 한 번 띄우는 것만으로 네 개가 생깁니다.


8-1. 블로킹 재시도의 한계 — 그리고 발상의 전환

Step 07 의 결론은 이랬습니다. DefaultErrorHandler 의 백오프 대기는 컨슈머 스레드 안의 Thread.sleep 이고, 그 스레드는 poll() 도 겸하기 때문에 재시도 중에는 그 파티션의 모든 일이 멈춥니다. FixedBackOff(10000L, 5L) 하나로 orders-1 이 50초 정지했고 LAG 이 47까지 올랐습니다.

여기서 나오는 질문은 하나입니다. "실패한 메시지를 붙잡고 있지 말고, 어딘가에 치워 두면 안 되나?"

그게 논블로킹 재시도입니다. 실패하면 그 레코드를 별도의 retry 토픽으로 발행하고, 원본 토픽의 오프셋은 즉시 커밋합니다. 본선 컨슈머는 바로 다음 레코드로 넘어갑니다. 치워 둔 레코드는 retry 토픽을 구독하는 다른 컨슈머가 백오프 시간이 지난 뒤 다시 처리합니다.

                    ┌─────────────────────────────────────────────┐
                    │  topic: orders  (partitions=3)              │
   프로듀서 ───────▶│  [98][99][100][101][102] ...                │
                    └───────────────┬─────────────────────────────┘
                                    │ group=s08-inventory

                          ┌───────────────────┐
                          │ InventoryListener │  ① 100 처리 → 예외
                          └─────────┬─────────┘  ② 100 을 retry-0 으로 발행
                                    │            ③ 오프셋 101 로 즉시 커밋
                                    │            ④ 101 처리 시작 ← 안 막힘 ★

        ┌──────────────────┐  실패  ┌──────────────────┐  실패  ┌──────────────────┐
        │ orders-retry-0   │──────▶│ orders-retry-1   │──────▶│ orders-retry-2   │
        │ (1초 뒤 처리)     │       │ (2초 뒤 처리)     │       │ (4초 뒤 처리)     │
        └──────────────────┘       └──────────────────┘       └────────┬─────────┘
                                                                       │ 재시도 소진

                                                     ┌──────────────────────────────┐
                                                     │ orders.DLT  @DltHandler 수신 │
                                                     └──────────────────────────────┘

본선 파티션은 ③에서 이미 커밋되었습니다. 재시도가 10분이 걸리든 1시간이 걸리든 orders-1 의 처리량에는 아무 영향이 없습니다. max.poll.interval.ms 초과도 없습니다.

대가는 하나입니다. 100 번 레코드가 101, 102 보다 늦게 처리됩니다. 이것이 8-6 의 주제이고, 이 스텝에서 딱 하나만 기억해야 한다면 그 절입니다.


8-2. @RetryableTopic — 애노테이션 한 줄로 시작

@RetryableTopic(
        attempts = "4",
        backoff = @Backoff(delay = 1000, multiplier = 2.0),
        kafkaTemplate = "kafkaTemplate")
@KafkaListener(topics = "orders", groupId = "s08-inventory")
public void onMessage(ConsumerRecord<String, OrderCreated> record) {
    if (record.partition() == 1 && record.offset() == 42) {
        throw new RemoteApiException("재고 API 타임아웃: " + record.value().orderId());
    }
    log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
}

attempts = "4"최초 시도를 포함한 총 호출 횟수 4번입니다. Step 07 의 FixedBackOff(interval, maxAttempts) 가 "재시도 횟수"였던 것과 세는 방식이 다릅니다. 4번 시도 = 최초 1 + 재시도 3 = retry 토픽 3개.

값이 전부 문자열인 것에도 이유가 있습니다. attempts = "${app.retry.attempts:4}" 처럼 프로퍼티 플레이스홀더를 쓸 수 있게 하기 위해서입니다.

앱을 띄우고 다시 목록을 봅니다.

kt --list

결과

orders
orders-retry-1000
orders-retry-2000
orders-retry-4000
orders-dlt
orders.DLT
payments

토픽 4개가 생겼습니다. 그런데 이름이 예상과 다릅니다.

토픽 명명 규칙 — 기본값은 "지연 시간 접미사"

@RetryableTopictopicSuffixingStrategy 기본값은 SUFFIX_WITH_DELAY_VALUE 입니다. retry 토픽 이름 뒤에 인덱스가 아니라 백오프 지연값(ms) 이 붙습니다. 그래서 orders-retry-1000, -2000, -4000 이 됩니다.

속성기본값결과
retryTopicSuffix-retryorders-retry...
dltTopicSuffix-dltorders-dlt
topicSuffixingStrategySUFFIX_WITH_DELAY_VALUEorders-retry-1000
topicSuffixingStrategySUFFIX_WITH_INDEX_VALUEorders-retry-0

이 코스의 규약은 orders-retry-0orders.DLT 입니다(코스 개요의 토픽 표). 기본값 그대로 두면 Step 07 에서 만든 orders.DLT 와 새로 생긴 orders-dlt따로 놀게 됩니다. 맞춰 줍니다.

@RetryableTopic(
        attempts = "4",
        backoff = @Backoff(delay = 1000, multiplier = 2.0),
        topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
        retryTopicSuffix = "-retry",
        dltTopicSuffix = ".DLT",
        kafkaTemplate = "kafkaTemplate")

결과

orders
orders-retry-0
orders-retry-1
orders-retry-2
orders.DLT
payments

kt --describe --topic orders-retry-0 로 파티션 수도 봅니다.

결과

Topic: orders-retry-0	TopicId: qL8vN1RhSem3xT4dPzKcAg	PartitionCount: 3	ReplicationFactor: 1	Configs: segment.bytes=1073741824
	Topic: orders-retry-0	Partition: 0	Leader: 1	Replicas: 1	Isr: 1
	Topic: orders-retry-0	Partition: 1	Leader: 1	Replicas: 1	Isr: 1
	Topic: orders-retry-0	Partition: 2	Leader: 1	Replicas: 1	Isr: 1

파티션 3, 복제 1. @RetryableTopic 은 원본 토픽의 파티션 수를 읽어 오지 않습니다. numPartitions 기본값이 -1 이라 KafkaAdmin 의 기본(브로커 num.partitions, 실습 환경은 3)을 따랐을 뿐입니다. 원본이 12 파티션이면 retry 토픽만 3 파티션이 되어 병렬도가 조용히 1/4 로 떨어집니다. numPartitions = "12" 로 명시하세요.

⚠️ 함정 — 기본 DLT 접미사는 -dlt 이고, Step 07 의 .DLT 와 다른 토픽이다 Step 07 에서 DeadLetterPublishingRecoverer 가 쓰던 목적지는 <topic>.DLT 였습니다. @RetryableTopic 의 기본값은 <topic>-dlt 입니다. 점이냐 하이픈이냐 하나 차이로 완전히 다른 토픽입니다. 증상: 기존 DLT 모니터링·알림은 orders.DLT 를 보고 있는데 실패 메시지는 전부 orders-dlt 로 쌓입니다. 에러는 없습니다. 대시보드가 조용해서 "장애가 없어졌다"고 착각합니다. 해결: dltTopicSuffix = ".DLT" 를 명시하거나, 팀 전체 규약을 -dlt 로 통일하세요. 둘 중 무엇이든 한 프로젝트에 두 규칙이 공존하면 안 됩니다.


8-3. 실패 메시지 한 건의 여정을 끝까지 따라가기

orders-1@42 하나만 항상 실패하도록 두고 전체 흐름을 봅니다. attempts = "4", delay = 1000, multiplier = 2.0 입니다.

결과 — 기동 로그

INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-0, orders-1, orders-2]
INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-retry-0-0, orders-retry-0-1, orders-retry-0-2]
INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-retry-1-0, orders-retry-1-1, orders-retry-1-2]
INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-retry-2-0, orders-retry-2-1, orders-retry-2-2]
INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders.DLT-0, orders.DLT-1, orders.DLT-2]

리스너 하나를 선언했는데 컨테이너가 5개 떴습니다. 본선 1 + retry 3 + DLT 1. 전부 s08-inventory 같은 컨슈머 그룹입니다.

결과 — 실패 레코드의 이동

INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$JourneyDemo        : 처리 완료 orders-1@41
WARN  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$JourneyDemo        : ★ 실패 topic=orders attempt=1 ORD-0043
INFO  13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-0-1
INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$JourneyDemo        : 처리 완료 orders-1@43   ← ★ 0.4ms 뒤. 본선은 안 막혔다
INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$JourneyDemo        : 처리 완료 orders-1@44

WARN  13309 --- [ntainer#1-1-C-1] c.e.o.step08.Practice$JourneyDemo        : ★ 실패 topic=orders-retry-0 attempt=2 ORD-0043   (t=1006ms)
INFO  13309 --- [ntainer#1-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-1-1

WARN  13309 --- [ntainer#2-1-C-1] c.e.o.step08.Practice$JourneyDemo        : ★ 실패 topic=orders-retry-1 attempt=3 ORD-0043   (t=3012ms)
INFO  13309 --- [ntainer#2-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-2-1

WARN  13309 --- [ntainer#3-1-C-1] c.e.o.step08.Practice$JourneyDemo        : ★ 실패 topic=orders-retry-2 attempt=4 ORD-0043   (t=7021ms)
INFO  13309 --- [ntainer#3-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders.DLT-1
ERROR 13309 --- [ntainer#4-1-C-1] c.e.o.step08.Practice$JourneyDemo        : DLT 수신 ORD-0043 origin=orders-1@42 ex=RemoteApiException

읽어야 할 것이 넷입니다.

  1. 스레드 이름이 매번 바뀝니다. [ntainer#0-1-C-1]#1#2#3#4. 컨테이너가 다르다는 증거입니다. Step 07 에서는 재시도가 전부 같은 스레드였습니다.
  2. 처리 완료 orders-1@43 이 실패 직후에 찍혔습니다. 본선은 0.4ms 만에 다음 레코드로 넘어갔습니다.
  3. 시각이 1006ms → 3012ms → 7021ms. 대기 간격이 1초 → 2초 → 4초로 두 배씩 늘었습니다.
  4. 마지막은 orders.DLT-1. 원본 파티션 번호 1 이 그대로 유지됩니다.
kcg --describe --group s08-inventory

결과 — 재시도가 도는 도중(t=3s)에 찍은 것

GROUP          TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CLIENT-ID
s08-inventory  orders          0          100             100             0    consumer-s08-inventory-1
s08-inventory  orders          1          100             100             0    consumer-s08-inventory-2
s08-inventory  orders          2          100             100             0    consumer-s08-inventory-3
s08-inventory  orders-retry-0  1          1               1               0    consumer-s08-inventory-5
s08-inventory  orders-retry-1  1          1               1               0    consumer-s08-inventory-8
s08-inventory  orders-retry-2  1          0               1               1    consumer-s08-inventory-11

orders-1 의 LAG 이 0 입니다. Step 07 의 같은 상황에서는 47이었습니다. 밀린 것은 orders-retry-2 의 1건뿐입니다.

구간본선 평균 지연본선 최대 지연실패 레코드의 총 지연
Step 07 블로킹 (FixedBackOff(10s, 5))24,300ms50,180ms50,180ms
Step 08 논블로킹 (attempts=4, 1s×2배)8ms43ms7,021ms

본선 지연 24.3초 → 8ms. 3,000배 빨라졌습니다. 그런데 이 숫자만 보고 "논블로킹이 우월하다"고 결론 내리면 8-6 에서 크게 다칩니다.


8-4. 백오프 전략과 토픽 개수의 관계

@Backoff 는 Spring Retry 의 애노테이션입니다. 네 가지 속성이 의미가 있습니다.

@Backoff(delay = 1000, multiplier = 2.0, maxDelay = 10000, random = true)
속성의미기본값
delay첫 대기(ms)1000
multiplier배수. 1.0 이면 고정 간격0 (= 고정)
maxDelay상한(ms). 0 이면 무제한0
random지터. 실제 대기를 [delay, delay*multiplier) 에서 무작위 선택false

random = true동시에 실패한 수천 건이 정확히 같은 순간에 재시도하는 것(thundering herd) 을 막습니다. 외부 API 장애 복구 직후가 정확히 그런 상황입니다.

시도 횟수와 토픽 개수

retry 토픽 개수는 attempts - 1 이 기본입니다. 각 재시도마다 대기 시간이 다르니 토픽도 따로 있어야 하기 때문입니다.

설정대기 시퀀스retry 토픽전체 토픽 수
attempts=3, delay=1000, multiplier=2.01s, 2s-retry-0, -retry-14
attempts=4, delay=1000, multiplier=2.01s, 2s, 4s-retry-0 ~ -25
attempts=5, delay=1000, multiplier=2.0, maxDelay=30001s, 2s, 3s, 3s-retry-0 ~ -36
attempts=5, delay=1000 (multiplier 없음)1s, 1s, 1s, 1s-retry 1개3

마지막 줄이 중요합니다. multiplier 를 안 주면 모든 재시도의 대기가 같으므로 토픽을 나눌 이유가 없습니다. Spring 이 알아서 orders-retry 하나로 합치고 4번 순환시킵니다. 접미사에 숫자도 안 붙습니다.

SameIntervalTopicReuseStrategy — 같은 간격 구간을 합칠지

maxDelay 로 잘린 뒷부분은 대기가 전부 같습니다. 위 세 번째 줄의 3s, 3s 가 그렇습니다. 이걸 토픽 하나로 합칠지 정하는 것이 sameIntervalTopicReuseStrategy 입니다.

@RetryableTopic(
        attempts = "6",
        backoff = @Backoff(delay = 1000, multiplier = 2.0, maxDelay = 4000),
        sameIntervalTopicReuseStrategy = SameIntervalTopicReuseStrategy.SINGLE_TOPIC)
전략대기 시퀀스 1s,2s,4s,4s,4s생성 토픽
MULTIPLE_TOPICS (기본)각각 따로-retry-0, -1, -2, -3, -4 (5개)
SINGLE_TOPIC뒤 3개를 합침-retry-0, -1, -retry (3개)

💡 Spring Kafka 3.0 에서 fixedDelayTopicStrategy(타입 FixedDelayStrategy)가 deprecated 되고 sameIntervalTopicReuseStrategy(SameIntervalTopicReuseStrategy)로 대체됐습니다. 값 이름(SINGLE_TOPIC / MULTIPLE_TOPICS)은 같습니다. 3.1.4 에서 옛 속성을 쓰면 컴파일은 되지만 deprecation 경고가 뜨고, 두 속성을 같이 주면 새 쪽이 이깁니다.

⚠️ 함정 — attempts 를 크게 잡으면 토픽이 그만큼 늘어난다 attempts = "10" 은 retry 토픽 9개를 만듭니다. 파티션 3이면 파티션 27개, 컨슈머 컨테이너 9개, 컨슈머 스레드 27개가 리스너 하나에서 생깁니다. 리스너가 5개면 토픽 50개, 파티션 150개입니다. 증상: 기동이 눈에 띄게 느려지고(토픽 생성이 순차적입니다), 브로커의 파티션 수가 통제를 벗어나며, kt --list 가 읽을 수 없는 상태가 됩니다. 리밸런스 시간도 파티션 수에 비례해 늘어납니다. 해결: attempts4~5를 넘기지 마세요. 더 오래 기다려야 하면 횟수가 아니라 maxDelay 를 키우고 SameIntervalTopicReuseStrategy.SINGLE_TOPIC 으로 토픽을 합치세요. 다만 그 경우 8-8 의 함정에 걸립니다.


8-5. @DltHandler — 종착지 처리

Step 07 에서는 DLT 리스너를 @KafkaListener(topics = "orders.DLT") 로 따로 만들었습니다. @RetryableTopic 을 쓰면 같은 클래스 안에 @DltHandler 메서드 하나면 됩니다.

@DltHandler
public void onDlt(OrderCreated event,
                  @Header(KafkaHeaders.ORIGINAL_TOPIC)     String topic,
                  @Header(KafkaHeaders.ORIGINAL_PARTITION) int    partition,
                  @Header(KafkaHeaders.ORIGINAL_OFFSET)    long   offset,
                  @Header(KafkaHeaders.EXCEPTION_FQCN)     String exFqcn,
                  @Header(KafkaHeaders.EXCEPTION_MESSAGE)  String exMessage) {
    log.error("DLT 수신 order={} origin={}-{}@{} ex={} msg={}",
            event.orderId(), topic, partition, offset, exFqcn, exMessage);
}

헤더 이름이 Step 07 과 다릅니다. @RetryableTopic 인프라는 kafka_dlt-original-topic 이 아니라 kafka_original-topic 을 씁니다.

상수실제 헤더내용
KafkaHeaders.ORIGINAL_TOPICkafka_original-topicorders (retry 토픽이 아니라 최초 출발지)
KafkaHeaders.ORIGINAL_PARTITIONkafka_original-partition4바이트 int
KafkaHeaders.ORIGINAL_OFFSETkafka_original-offset8바이트 long
KafkaHeaders.EXCEPTION_FQCNkafka_exception-fqcn예외 FQCN
KafkaHeaders.EXCEPTION_MESSAGEkafka_exception-message예외 메시지
KafkaHeaders.EXCEPTION_STACKTRACEkafka_exception-stacktrace스택트레이스 전문
RetryTopicHeaders.DEFAULT_HEADER_ATTEMPTSretry_topic-attempts현재 시도 횟수
RetryTopicHeaders.DEFAULT_HEADER_BACKOFF_TIMESTAMPretry_topic-backoff-timestamp이 시각 전에는 처리하지 말라
RetryTopicHeaders.DEFAULT_HEADER_ORIGINAL_TIMESTAMPretry_topic-original-timestamp최초 발행 시각
kcc --topic orders-retry-1 --from-beginning --property print.headers=true --property print.key=true

결과 (읽기 쉽게 줄바꿈했습니다)

kafka_original-topic:orders,
kafka_original-partition:\x00\x00\x00\x01,
kafka_original-offset:\x00\x00\x00\x00\x00\x00\x00\x2A,
kafka_exception-fqcn:com.example.order.step08.RemoteApiException,
kafka_exception-message:재고 API 타임아웃: ORD-0043,
retry_topic-attempts:\x00\x00\x00\x03,
retry_topic-backoff-timestamp:1735689603012,
retry_topic-original-timestamp:1735689600006,
__TypeId__:com.example.order.domain.OrderCreated
	ORD-0043	{"orderId":"ORD-0043","customerId":1013,"sku":"SKU-002","quantity":2,"amount":11000,...}

retry_topic-backoff-timestamp 가 이 스텝의 숨은 주인공입니다. 8-8 에서 다시 봅니다. 그리고 retry_topic-attempts바이너리 int, retry_topic-backoff-timestamp문자열로 찍힌 epoch ms 라 형식이 섞여 있습니다. 직접 파싱하지 말고 @Header(name = ..., required = false) byte[] 로 받아 확인하세요.

DltStrategy — DLT 처리마저 실패하면

@RetryableTopic(attempts = "4", dltStrategy = DltStrategy.FAIL_ON_ERROR)
동작
FAIL_ON_ERROR (기본)@DltHandler 가 예외를 던지면 그 DLT 파티션이 막힙니다. 오프셋이 안 넘어갑니다
ALWAYS_RETRY_ON_ERRORDLT 처리 실패 시 DLT 토픽으로 다시 발행. 원인이 안 고쳐지면 무한 루프
NO_DLTDLT 자체를 안 만듭니다. 소진되면 버립니다

⚠️ 함정 — NO_DLT 는 Step 07 의 "조용히 버린다" 로 되돌아가는 스위치다 "DLT 토픽이 늘어나는 게 부담스럽다"는 이유로 dltStrategy = DltStrategy.NO_DLT 를 켜는 경우가 있습니다. 이 순간 재시도를 소진한 메시지는 어디에도 남지 않고 사라집니다. LAG 은 0, 에러 로그는 KafkaBackoffException 도 아닌 평범한 WARN 한 줄입니다. 증상: retry 토픽에는 메시지가 흘러간 흔적이 있는데 최종 목적지가 없습니다. kt --listorders.DLT 가 아예 없으니 "원래 그런 설계인가" 하고 넘어가게 됩니다. 해결: NO_DLT재시도 자체가 부가 기능인 경우(예: 조회수 집계)에만 쓰세요. 그 경우에도 소진 시점에 카운터 지표를 올려 두어야 합니다. 기본값 FAIL_ON_ERROR 를 유지하는 것이 원칙입니다.


8-6. ⚠️ 핵심 함정 — 논블로킹 재시도는 메시지 순서를 깨뜨린다

이 절이 Step 08 의 전부입니다.

주문 ORD-0001 에 대해 세 이벤트를 순서대로 발행합니다. 키가 같으므로 셋 다 같은 파티션(orders-1) 에 들어가고, 오프셋은 100, 101, 102 입니다.

template.send("orders", "ORD-0001", OrderEvent.of("ORD-0001", OrderState.CREATED));    // offset 100
template.send("orders", "ORD-0001", OrderEvent.of("ORD-0001", OrderState.UPDATED));    // offset 101
template.send("orders", "ORD-0001", OrderEvent.of("ORD-0001", OrderState.CANCELLED));  // offset 102

컨슈머는 상태를 맵에 반영합니다. 그리고 첫 번째(CREATED)만 한 번 실패하도록 만듭니다. 두 번째 시도에서는 성공합니다 — 실제 일시적 장애가 딱 이렇게 동작합니다.

@RetryableTopic(attempts = "4", backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                dltTopicSuffix = ".DLT", kafkaTemplate = "kafkaTemplate")
@KafkaListener(topics = "orders", groupId = "s08-order-state")
public void onMessage(ConsumerRecord<String, OrderEvent> record) {
    OrderEvent e = record.value();
    if (e.status() == OrderState.CREATED && FIRST_TRY.getAndSet(false)) {
        throw new RemoteApiException("재고 API 타임아웃");   // 최초 1회만 실패
    }
    STATE.put(e.orderId(), e.status());
    log.info("상태 반영 {} → {}  (topic={} t={}ms)", e.orderId(), e.status(), record.topic(), elapsed());
}

결과

INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : 발행 ORD-0001 CREATED → orders-1@100  (t=0ms)
INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : 발행 ORD-0001 UPDATED → orders-1@101  (t=2ms)
INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : 발행 ORD-0001 CANCELLED → orders-1@102  (t=3ms)
WARN  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : ★ 실패 ORD-0001 CREATED (topic=orders t=8ms)
INFO  13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-0-1
INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : 상태 반영 ORD-0001 → UPDATED    (topic=orders t=12ms)
INFO  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : 상태 반영 ORD-0001 → CANCELLED  (topic=orders t=15ms)
INFO  13309 --- [ntainer#1-1-C-1] c.e.o.step08.Practice$OrderStateDemo     : 상태 반영 ORD-0001 → CREATED    (topic=orders-retry-0 t=1012ms)
INFO  13309 --- [           main] c.e.o.step08.Practice$OrderStateDemo     : 최종 상태 ORD-0001 = CREATED     ← ★★ 취소된 주문이 살아났다

타임라인

시각토픽이벤트결과STATE["ORD-0001"]
t=0msCREATED 발행orders-1@100(없음)
t=2msUPDATED 발행orders-1@101(없음)
t=3msCANCELLED 발행orders-1@102(없음)
t=8msordersCREATED 처리실패orders-retry-0 발행, 오프셋 101 커밋(없음)
t=12msordersUPDATED 처리성공UPDATED
t=15msordersCANCELLED 처리성공CANCELLED
t=1012msorders-retry-0CREATED 재처리성공CREATED ← 오염

취소된 주문이 1초 뒤에 되살아났습니다. 예외는 하나도 없고, kcg --describe 의 LAG 은 전부 0 이며, 로그에도 ERROR 가 없습니다. @DltHandler 는 호출되지 않았습니다 — 재시도가 성공했으니까요. 이 코스가 계속 말하는 "에러 없이 조용히 잘못 동작하는" 상태의 교과서적 사례입니다.

블로킹 재시도였다면 이런 일이 없습니다. orders-1@100 이 성공하거나 DLT 로 갈 때까지 101, 102 는 아예 poll 되지 않기 때문입니다. Step 07 의 "파티션이 멈춘다" 는 단점은, 동시에 순서 보장이라는 장점의 다른 이름입니다.

⚠️ 함정 — 논블로킹 재시도는 같은 키의 순서 보장을 전면 포기한다 Kafka 의 순서 보장은 "하나의 파티션 안" 에서만 성립합니다. @RetryableTopic 은 실패한 레코드를 다른 토픽으로 옮기므로 그 순간 보장 범위 밖으로 나갑니다. 같은 키라도 마찬가지입니다. 증상 3종: ① 상태 전이가 역행합니다(취소 → 생성, 환불 → 결제). ② 카운터가 어긋납니다(재고 차감/복원 순서 뒤바뀜). ③ 재현이 안 됩니다. 실패가 없으면 순서가 지켜지므로, 외부 API 가 흔들린 몇 분 동안만 데이터가 오염되고 로그에는 흔적이 없습니다. 해결 3가지: ① 순서가 중요하면 논블로킹을 쓰지 마세요. 블로킹 + 짧은 백오프 + 즉시 DLT 가 정답입니다. ② 이벤트에 버전/시퀀스 번호를 넣고 컨슈머가 if (incoming.version() <= current.version()) return; 으로 역행을 거부하게 만드세요(Step 13 의 멱등 컨슈머). ③ 이벤트를 상태 전이(delta)가 아니라 전체 상태(snapshot) 로 설계하면 늦게 온 오래된 스냅샷을 버릴 수 있습니다.

💡 판단 문장 하나로 줄이면 이렇습니다. "이 메시지가 1초 늦게 처리돼도 결과가 같은가?" 같으면 논블로킹, 다르면 블로킹입니다.


8-7. retry 토픽의 컨슈머 그룹과 파티션

8-3 에서 컨테이너가 5개 떴습니다. 컨슈머 그룹은 몇 개일까요.

kcg --list

결과

s08-inventory

하나입니다. retry 토픽과 DLT 는 본선과 같은 그룹 ID 를 씁니다. kcg --describe 를 보면 한 그룹이 5개 토픽을 구독하고 있습니다.

kcg --describe --group s08-inventory

결과

GROUP          TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG  CLIENT-ID
s08-inventory  orders          0          300             300             0    consumer-s08-inventory-1
s08-inventory  orders          1          300             300             0    consumer-s08-inventory-2
s08-inventory  orders          2          300             300             0    consumer-s08-inventory-3
s08-inventory  orders-retry-0  1          3               3               0    consumer-s08-inventory-5
s08-inventory  orders-retry-1  1          3               3               0    consumer-s08-inventory-8
s08-inventory  orders-retry-2  1          3               3               0    consumer-s08-inventory-11
s08-inventory  orders.DLT      1          3               3               0    consumer-s08-inventory-14

(LAG 0 인 나머지 파티션 행은 생략했습니다)

concurrency: 3 이 본선뿐 아니라 retry 토픽에도 그대로 적용되어 컨슈머가 3+3+3+3+3 = 15개 만들어졌습니다. retry 토픽에는 보통 메시지가 거의 없는데도 스레드 12개가 놀고 있습니다.

@RetryableTopic(attempts = "4", concurrency = "1", ...)

concurrency = "1" 로 retry 컨테이너만 1 스레드로 줄입니다. 본선의 concurrency 는 그대로입니다.

결과concurrency = "1" 적용 후

INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-0, orders-1, orders-2]
INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-retry-0-0, orders-retry-0-1, orders-retry-0-2]

두 번째 줄의 파티션 3개를 한 스레드가 전부 맡습니다. 스레드 15개 → 6개. 다만 8-8 의 이유로 retry 토픽의 concurrency 를 1로 줄이면 대기가 직렬화되므로, 실패율이 높은 리스너에서는 신중해야 합니다.

💡 실무 팁 — 컨슈머 그룹이 하나라 오프셋 리셋이 편합니다. kcg --group s08-inventory --topic orders-retry-0 --reset-offsets --to-earliest --execute 처럼 토픽 단위로 각각 리셋해야 하고, 실습을 되돌리려면 retry 토픽 3개 + DLT 까지 4번 리셋해야 합니다. --all-topics 를 쓰면 한 번에 됩니다.


8-8. ⚠️ 함정 — retry 토픽의 지연은 "최소" 보장일 뿐이고, 길면 다시 블로킹된다

@Backoff(delay = 60000) 이면 "60초 뒤에 처리된다"고 읽기 쉽습니다. 틀렸습니다. "60초 전에는 처리되지 않는다" 가 맞습니다.

Spring Kafka 의 논블로킹 재시도에는 스케줄러가 없습니다. orders-retry-0 컨슈머는 평범한 컨슈머라서 메시지를 즉시 poll() 해 옵니다. 그리고 KafkaBackoffAwareMessageListenerAdapterretry_topic-backoff-timestamp 헤더를 보고 판단합니다.

retry-0 컨슈머 스레드 [ntainer#1-0-C-1]
  poll() → 레코드 (backoff-timestamp = now + 1000ms)

  KafkaBackoffAwareMessageListenerAdapter
     아직 시간이 안 됐다 → KafkaBackoffException 을 던진다

  KafkaConsumerBackoffManager
     ① consumer.pause(orders-retry-0-1)      ← ★ 그 파티션을 통째로 멈춘다
     ② consumer.seek(orders-retry-0-1, 그 오프셋)

  1000ms 뒤 → resume → 다시 poll → 이번엔 시간이 됐다 → 리스너 호출

pause 가 핵심입니다. 대기 중에는 그 retry 토픽 파티션 전체가 멈춥니다. Thread.sleep 이 아니라 pause 라서 poll() 은 계속 돌고 max.poll.interval.ms 초과는 나지 않습니다 — 이 점이 Step 07 과 결정적으로 다릅니다. 하지만 그 파티션의 뒤 레코드가 대기한다는 사실 자체는 동일합니다.

결과delay = 30000 으로 두고 짧은 간격으로 3건이 실패했을 때

INFO  13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-0-1
INFO  13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-0-1
INFO  13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-0-1
DEBUG 13309 --- [ntainer#1-0-C-1] o.s.k.l.KafkaConsumerBackoffManager      : Backing off partition orders-retry-0-1 until 2025-01-01T00:00:30.014Z (dueTimestamp=1735689630014)
INFO  13309 --- [ntainer#1-0-C-1] o.s.k.l.KafkaMessageListenerContainer    : Paused consumption from: [orders-retry-0-1]
INFO  13309 --- [ntainer#1-0-C-1] c.e.o.step08.Practice$LongBackoffDemo    : 재처리 ORD-0043 (topic=orders-retry-0 대기 30,012ms)
DEBUG 13309 --- [ntainer#1-0-C-1] o.s.k.l.KafkaConsumerBackoffManager      : Backing off partition orders-retry-0-1 until 2025-01-01T00:00:30.019Z (dueTimestamp=1735689630019)
INFO  13309 --- [ntainer#1-0-C-1] c.e.o.step08.Practice$LongBackoffDemo    : 재처리 ORD-0044 (topic=orders-retry-0 대기 30,018ms)
INFO  13309 --- [ntainer#1-0-C-1] c.e.o.step08.Practice$LongBackoffDemo    : 재처리 ORD-0045 (topic=orders-retry-0 대기 30,021ms)

여기까지는 괜찮습니다. 세 건의 due time 이 거의 같아서 한 번의 pause 로 셋이 함께 풀렸습니다. 문제는 due time 이 흩어져 있을 때입니다. @Backoff(random = true) 를 켜거나, maxDelay 가 다른 두 리스너가 같은 retry 토픽을 공유하면 파티션 안에서 due time 이 뒤죽박죽이 됩니다. delay = 300000(5분) + random = true 로 재현했습니다.

retry-0 파티션에 쌓인 순서due time실제 처리 시각초과 지연
ORD-Et=900st=900s0s
ORD-Ft=310s (이미 지남)t=900s590s
ORD-Gt=320s (이미 지남)t=900s580s

ORD-F 와 ORD-G 는 진작에 처리됐어야 하는데 앞의 ORD-E 가 파티션을 900초까지 붙잡고 있어서 함께 기다렸습니다. 백오프 5분이 실제로는 15분이 됐습니다.

⚠️ 함정 — 긴 백오프는 retry 토픽 안에서 다시 블로킹이다 @RetryableTopic 의 장점은 "본선 파티션을 안 막는다"입니다. retry 토픽 파티션은 그대로 막힙니다. maxDelay 를 시간 단위로 크게 잡는 순간, retry 토픽이 Step 07 의 본선 토픽과 똑같은 병목이 됩니다. 증상: kcg --describe 에서 orders 의 LAG 은 0인데 orders-retry-N 의 LAG 만 계속 쌓입니다. 처리는 되지만 백오프 설정값보다 훨씬 늦게 됩니다. 지연 지표를 안 보면 발견이 안 됩니다. 해결: ① retry 토픽의 파티션 수와 concurrency 를 늘려 대기를 병렬화하세요. 파티션 하나가 막혀도 다른 파티션은 돕니다. ② @Backoff(random = true) 로 due time 을 분산시키면 한 파티션에 몰리는 것을 줄일 수 있지만, 위 표처럼 역전도 함께 만듭니다. 트레이드오프입니다. ③ 백오프가 분 단위를 넘어가면 재시도 토픽이 아니라 스케줄러/배치로 설계하세요. DLT 에 넣고 별도 잡이 주기적으로 재처리하는 편이 예측 가능합니다.


8-9. 블로킹 vs 논블로킹 — 선택 기준표

블로킹 (DefaultErrorHandler)논블로킹 (@RetryableTopic)
순서 보장보존됨 (seek 으로 제자리 반복)깨짐 (다른 토픽으로 이동)
파티션 블로킹본선 파티션이 백오프 총합만큼 정지본선은 즉시 통과. retry 토픽은 막힘
처리 지연(정상 메시지)실패 1건에 뒤 전부가 대기 (실측 24.3초)영향 없음 (실측 8ms)
최대 백오프max.poll.interval.ms 의 1/3 이하 (~100초)사실상 무제한 (pause 방식이라 poll timeout 없음)
토픽 개수원본 + DLT = 2원본 + retry attempts-1 + DLT = 최대 6~10
컨슈머 스레드concurrency 그대로concurrency × (attempts+1) 까지 증가
운영 복잡도낮음. 토픽 2개, 그룹 1개높음. 토픽/파티션 폭증, 리밸런스 시간 증가
재시도 가시성로그와 LAG 뿐. 밖에서 안 보임retry 토픽에 실물이 쌓여 kcc 로 눈으로 확인
재시작 시 재시도 상태소실. 처음부터 다시보존. retry 토픽에 남아 있음
적합한 실패 유형짧은 순간 장애(DB 락 경합, 커넥션 재수립)긴 외부 장애(외부 API 다운, 서드파티 점검)

두 줄로 요약하면 이렇습니다. 블로킹은 "정확성을 지키고 처리량을 포기" 하고, 논블로킹은 "처리량을 지키고 정확성을 포기" 합니다. 그리고 재시작 시 재시도 상태 보존은 논블로킹의 잘 안 알려진 강점입니다. 블로킹 재시도 중에 앱을 배포하면 재시도 진행 상황이 통째로 사라지고 처음부터 다시 시작하지만, retry 토픽의 메시지는 브로커에 남아 있어 새 인스턴스가 이어받습니다.

결정 트리

실패한 메시지가 5초 뒤에 똑같이 하면 성공할 가능성이 있는가?

├─ 없다 (검증 오류, 스키마 불일치, 잔액 부족)
│    → 재시도 금지. addNotRetryableExceptions / @RetryableTopic(exclude=...)
│      → 즉시 DLT

└─ 있다

     └─ 이 메시지가 뒤 메시지보다 늦게 처리돼도 결과가 같은가?

          ├─ 아니다 (주문 상태 전이, 계좌 잔액, 재고 증감)
          │    → 블로킹. DefaultErrorHandler
          │      + 짧은 백오프 (총합 ≤ max.poll.interval.ms / 3)
          │      + DeadLetterPublishingRecoverer

          └─ 그렇다 (알림 발송, 검색 색인, 통계 집계, 캐시 갱신)

               └─ 얼마나 기다려야 하나?
                    ├─ 수 초         → 블로킹도 충분. 굳이 토픽을 늘리지 말 것
                    ├─ 수십 초 ~ 수 분 → 논블로킹 @RetryableTopic ★
                    └─ 수십 분 이상   → 즉시 DLT + 별도 배치 재처리 잡
상황선택설정
주문 상태 전이(OrderStatus)블로킹ExponentialBackOffWithMaxRetries(4) 0.5s/2배/상한 5s + DLT
계좌 잔액 증감블로킹FixedBackOff(500L, 3L) + DLT. 재시도보다 멱등성이 우선
결제 알림 SMS 발송논블로킹attempts=4, delay=5000, multiplier=3.0
검색 색인 갱신논블로킹attempts=5, delay=10000, maxDelay=60000, SINGLE_TOPIC
스키마 검증 실패재시도 없음exclude = InvalidOrderException.class → 즉시 DLT
역직렬화 실패재시도 없음ErrorHandlingDeserializer + 기본 비재시도 목록 (Step 04·07)

8-10. 둘을 같이 쓰기 — 짧은 블로킹 뒤에 논블로킹

가장 실용적인 구성은 둘 중 하나를 고르는 게 아니라 겹쳐 쓰는 것입니다. 대부분의 순간 장애는 500ms 뒤 재시도로 끝납니다. 그런 실패까지 retry 토픽을 거치면 토픽만 지저분해지고 순서도 깨집니다.

"컨테이너에서 짧게 2회 블로킹 재시도 → 그래도 실패하면 retry 토픽으로" 가 정답에 가깝습니다.

@Configuration
public class RetryTopicGlobalConfig extends RetryTopicConfigurationSupport {

    /** ★ 논블로킹으로 넘기기 전에 컨테이너에서 짧게 블로킹 재시도한다. */
    @Override
    protected void configureBlockingRetries(BlockingRetriesConfigurer blockingRetries) {
        blockingRetries
                .retryOn(RemoteApiException.class, SocketTimeoutException.class)
                .backOff(new FixedBackOff(500L, 2L));   // 0.5초 × 2회 = 1초
    }

    /** retry/DLT 발행기를 손보는 훅. 반환 타입이 팩토리라는 점에 주의하세요. */
    @Override
    protected Consumer<DeadLetterPublishingRecovererFactory> configureDeadLetterPublishingContainerFactory() {
        return factory -> factory.setDeadLetterPublishingRecovererCustomizer(
                recoverer -> recoverer.setStripPreviousExceptionHeaders(true));
    }
}

configureBlockingRetries 에 등록한 예외만 먼저 블로킹으로 재시도합니다. 여기서 성공하면 retry 토픽에 아무것도 안 남고 순서도 유지됩니다. 실패하면 그때 @RetryableTopic 의 논블로킹 경로로 넘어갑니다.

결과RemoteApiException 이 5번째 시도에서 성공하는 경우

WARN  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$HybridDemo         : 실패 attempt=1 topic=orders       t=0ms
WARN  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$HybridDemo         : 실패 attempt=2 topic=orders       t=504ms    ← 블로킹
WARN  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$HybridDemo         : 실패 attempt=3 topic=orders       t=1008ms   ← 블로킹
INFO  13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Republishing failed record to orders-retry-0-1
WARN  13309 --- [ntainer#1-0-C-1] c.e.o.step08.Practice$HybridDemo         : 실패 attempt=4 topic=orders-retry-0 t=2014ms  ← 논블로킹
INFO  13309 --- [ntainer#2-0-C-1] c.e.o.step08.Practice$HybridDemo         : 성공 attempt=5 topic=orders-retry-1 t=4023ms

본선에서 3번(0ms, 504ms, 1008ms) 시도한 뒤에야 retry 토픽으로 넘어갑니다. 1초 안에 회복되는 장애는 순서를 지키며 처리되고, 그보다 긴 장애만 순서를 포기합니다.

⚠️ 함정 — 총 시도 횟수는 곱이 아니라 합이지만, 대기 시간은 곱해질 수 있다 블로킹 3회 × 논블로킹 4회를 "총 12번"으로 착각하기 쉽습니다. 실제로는 블로킹 재시도가 각 논블로킹 단계마다 반복됩니다. retry-0 컨슈머에서 실패해도 거기서 또 블로킹 2회를 돌고 나서 retry-1 로 갑니다. 증상: 예상 총 시간이 1+2+4=7초인데 실측이 11초입니다. 각 단계에서 1초씩 블로킹이 추가된 것입니다. retry 토픽 파티션도 그만큼 오래 막힙니다(8-8). 해결: 블로킹 백오프는 1초 이내 총합으로 아주 짧게 두세요. FixedBackOff(500L, 2L) 정도가 상한입니다.

애노테이션 대신 빈으로 — RetryTopicConfiguration

리스너가 여러 개면 애노테이션을 복붙하게 됩니다. 빈 하나로 토픽 단위 전역 설정을 할 수 있습니다.

@Bean
public RetryTopicConfiguration ordersRetryTopicConfig(KafkaTemplate<String, Object> template) {
    return RetryTopicConfigurationBuilder
            .newInstance()
            .includeTopic("orders")                       // 이 토픽을 구독하는 모든 리스너에 적용
            .maxAttempts(4)
            .exponentialBackoff(1000, 2.0, 10000)
            .suffixTopicsWithIndexValues()                // orders-retry-0, -1, -2
            .retryTopicSuffix("-retry")
            .dltSuffix(".DLT")
            .concurrency(1)
            .autoCreateTopics(true, 3, (short) 1)         // 파티션 3, 복제 1
            .notRetryOn(InvalidOrderException.class)      // 검증 실패는 즉시 DLT
            .doNotRetryOnDltFailure()                     // = DltStrategy.FAIL_ON_ERROR
            .create(template);
}

@RetryableTopic 애노테이션을 하나도 안 붙여도 orders 를 구독하는 모든 @KafkaListener 에 적용됩니다. 애노테이션과 빈이 둘 다 있으면 애노테이션이 이깁니다.

방식적용 범위언제
@RetryableTopic그 리스너 하나리스너마다 정책이 다를 때
RetryTopicConfiguration 빈 + includeTopic그 토픽을 구독하는 전부팀 표준 정책
RetryTopicConfigurationSupport 상속애플리케이션 전역 인프라블로킹 재시도 조합, 커스터마이저

8-11. ⚠️ 함정 — 토픽 자동 생성이 꺼진 운영 환경

@RetryableTopicautoCreateTopics 기본값은 "true" 이고, 기동 시 KafkaAdminNewTopic 빈으로 retry 토픽과 DLT 를 만듭니다. 그런데 운영 클러스터는 보통 브로커의 auto.create.topics.enable=false 이고, 애플리케이션 계정에 토픽 생성 ACL 이 없습니다(CREATE 없이 READ/WRITE 만). autoCreateTopics = "true" 그대로 두고 권한만 없는 상태로 띄워 봅니다.

결과 — 기동 로그

INFO  13309 --- [           main] o.s.k.core.KafkaAdmin                    : Could not configure topics
org.springframework.kafka.KafkaException: Failed to create topics; nested exception is
	java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TopicAuthorizationException:
	Authorization failed for topics [orders-retry-0, orders-retry-1, orders-retry-2, orders.DLT]
INFO  13309 --- [           main] o.s.k.l.KafkaMessageListenerContainer    : s08-inventory: partitions assigned: [orders-0, orders-1, orders-2]
INFO  13309 --- [           main] c.e.o.OrderServiceApplication            : Started OrderServiceApplication in 3.104 seconds

KafkaAdmin 이 실패했는데 애플리케이션은 정상 기동합니다. fatalIfBrokerNotAvailable 기본값이 false 이기 때문입니다. partitions assignedStarted 도 정상이고, 본선 처리는 몇 시간이고 잘 돌아갑니다 — 최초의 실패가 발생하는 순간까지는.

결과 — 첫 실패 발생 시

WARN  13309 --- [ntainer#0-1-C-1] c.e.o.step08.Practice$JourneyDemo        : ★ 실패 topic=orders attempt=1 ORD-0043
ERROR 13309 --- [ntainer#0-1-C-1] o.s.k.l.DeadLetterPublishingRecoverer    : Dead-letter publication to orders-retry-0 failed for: orders-1@42
org.apache.kafka.common.errors.TimeoutException: Topic orders-retry-0 not present in metadata after 60000 ms.
ERROR 13309 --- [ntainer#0-1-C-1] o.s.k.l.KafkaMessageListenerContainer    : Error handler threw an exception
org.springframework.kafka.KafkaException: Dead-letter publication failed; nested exception is ...

60초 침묵 뒤 실패하고, 그 60초 동안 orders-1 파티션은 완전히 멈춥니다. 논블로킹을 쓰는 이유였던 "파티션을 안 막는다"가 정반대로 뒤집혔습니다. 게다가 발행이 실패했으므로 그 레코드는 retry 토픽에도 DLT 에도 못 갑니다. Step 07 의 "조용히 버림" 과 같은 결말입니다.

해결은 토픽을 미리 만들고 자동 생성을 끄는 것입니다.

for i in 0 1 2; do
  kt --create --topic orders-retry-$i --partitions 3 --replication-factor 1
done
kt --create --topic orders.DLT --partitions 3 --replication-factor 1
@RetryableTopic(
        attempts = "4",
        backoff = @Backoff(delay = 1000, multiplier = 2.0),
        topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
        dltTopicSuffix = ".DLT",
        autoCreateTopics = "false",        // ★ 운영 필수
        kafkaTemplate = "kafkaTemplate")

⚠️ 함정 — autoCreateTopics = "false" 는 "토픽이 없어도 기동된다" 는 뜻이다 false 로 두면 Spring 은 토픽 생성을 시도하지 않을 뿐, 존재 여부를 검증하지도 않습니다. 토픽 이름에 오타가 있어도(dltTopicSuffix = ".DTL") 기동은 성공하고, 첫 실패가 날 때까지 아무도 모릅니다. 증상: 배포 후 며칠 뒤 첫 장애에서 Topic ... not present in metadata after 60000 ms 와 함께 파티션이 60초 단위로 멈춥니다. 하필 장애 대응 중에 두 번째 장애가 터지는 셈입니다. 해결: ① 사전 생성한 토픽 목록을 kt --list배포 체크리스트에 넣으세요. ② 애플리케이션 기동 시 AdminClient.describeTopics(...) 로 존재를 검증하고 없으면 기동을 실패시키세요. 조용히 뜨는 것보다 안 뜨는 게 낫습니다. ③ 스테이징에서 일부러 실패 메시지를 하나 흘려 전체 경로를 한 번 태워 보세요. 이것이 유일하게 확실한 검증입니다.

💡 토픽 생성 권한·auto.create.topics.enable·retention 같은 브로커 쪽 설정은 Kafka 코스 Step 09Step 14 를 참고하세요. retry 토픽의 retention.ms백오프 최대값보다 넉넉히 잡아야 합니다. 백오프 7일인데 retention 이 3일이면 재처리 전에 메시지가 삭제됩니다.


정리

개념핵심
논블로킹의 아이디어실패 레코드를 retry 토픽으로 옮기고 원본은 즉시 커밋. 본선 파티션 안 막힘
attempts최초 시도 포함 총 호출 횟수. attempts=4 → retry 토픽 3개
토픽 명명 기본값<topic>-retry-<지연ms>, <topic>-dlt. SUFFIX_WITH_DELAY_VALUE 가 기본
코스 규약 맞추기topicSuffixingStrategy=SUFFIX_WITH_INDEX_VALUE, dltTopicSuffix=".DLT"
컨테이너 수본선 1 + retry attempts-1 + DLT 1. 컨슈머 그룹은 하나
실측본선 지연 24.3초(블로킹) → 8ms(논블로킹)
핵심 함정순서가 깨진다. 재시도된 CREATED 가 CANCELLED 뒤에 도착해 취소된 주문이 부활
순서 깨짐의 특징예외 없음, LAG 0, ERROR 로그 없음, 재현 불가
@DltHandler같은 클래스에 두면 DLT 처리. 헤더는 kafka_original-* (Step 07 의 kafka_dlt-* 와 다름)
DltStrategyFAIL_ON_ERROR(기본) / ALWAYS_RETRY_ON_ERROR / NO_DLT(= 조용히 버림)
백오프의 실체스케줄러 없음. backoff-timestamp 헤더 + 파티션 pause
긴 백오프retry 토픽 안에서 다시 블로킹. 분 단위를 넘으면 배치로 설계할 것
지연 보장"정확히 N초 뒤"가 아니라 "N초 전에는 아님"
조합configureBlockingRetries 로 짧게 블로킹 → 실패 시 논블로킹
전역 설정RetryTopicConfiguration 빈 + includeTopic. 애노테이션이 우선
운영 함정autoCreateTopics 실패해도 기동은 성공. 첫 실패 때 60초 타임아웃
선택 기준 한 줄"이 메시지가 1초 늦게 처리돼도 결과가 같은가?"

연습문제

Exercise.java 에 6문제가 있습니다. 정답은 Solution.java. 반드시 직접 실행해 kt --list 와 로그를 확인하세요.

  1. @RetryableTopic(attempts = "4", backoff = @Backoff(delay = 1000, multiplier = 2.0)) 를 붙이고 기동한 뒤, kt --list 로 생성된 토픽 이름을 그대로 적기. 기본 접미사 전략이 무엇인지 이름에서 역추론하기
  2. 코스 규약(orders-retry-0 ~ -2, orders.DLT)에 맞도록 접미사와 접미사 전략을 지정하고, 지수 백오프 1s → 2s → 4s 의 시도 시각을 실측하기
  3. @DltHandler 를 구현해 KafkaHeaders.ORIGINAL_TOPIC / ORIGINAL_PARTITION / ORIGINAL_OFFSET / EXCEPTION_FQCN 을 로그로 남기기. 원본 토픽이 orders-retry-2 가 아니라 orders 로 찍히는 것을 확인하기
  4. 같은 키 ORD-0001 의 3건(생성 → 수정 → 취소)으로 순서 깨짐을 재현하고, 본문 8-6 형태의 타임라인 표를 채우기. 최종 상태가 CREATED 임을 확인하기
  5. InvalidOrderExceptionexclude 에 등록해 재시도 없이 즉시 DLT 로 보내기. retry 토픽에 아무것도 안 쌓이는 것을 kcc 로 확인하기
  6. 애노테이션을 전부 지우고 RetryTopicConfiguration 빈 하나로 같은 설정을 재현하기. includeTopic("orders") 로 두 리스너에 동시 적용되는 것을 확인하기

다음 단계

논블로킹 재시도는 파티션 블로킹 문제를 확실히 해결하지만, 순서 보장이라는 값을 치릅니다. 두 방식 중 무엇을 쓸지는 성능이 아니라 도메인의 성격이 정합니다.

그런데 8-6 의 오염을 근본적으로 막으려면 재시도 설정만으로는 부족합니다. "재고를 차감하고 이벤트를 발행하는 것"이 하나의 원자적 단위여야 합니다. 다음 스텝에서 Kafka 트랜잭션을 다룹니다. 그리고 거기서 이 코스의 가장 흔한 오해 하나를 깹니다 — @Transactional 하나로 DB 와 Kafka 를 묶을 수 있다는 믿음입니다.

Step 09 — 트랜잭션


실습 파일

이 스텝도 Java 파일 세 개로 진행합니다. Practice.java 를 보조 프로필로 하나씩 켜 가며 8-2 ~ 8-11 의 로그를 재현하고, 특히 step08-order 프로필로 8-6 의 순서 깨짐을 눈으로 확인하세요. 그다음 Exercise.java 의 6문제를 풀고 Solution.java 로 대조합니다. 세 파일 모두 com.example.order.step08 패키지에 둡니다. ⚠️ 프로필을 바꿀 때마다 retry 토픽을 지우세요. 접미사 전략이 다른 프로필끼리 토픽이 섞이면 로그를 읽을 수 없게 됩니다.

Practice.java

본문 8-2 ~ 8-11 의 모든 예제를 절 번호 주석과 함께 nested static class 로 담은 실행 파일입니다.

  • 프로필은 7개입니다. step08-default(기본 접미사 확인) / step08-journey(여정 추적) / step08-order(★ 순서 깨짐) / step08-longbackoff(retry 토픽 재블로킹) / step08-hybrid(블로킹+논블로킹) / step08-builder(빈 설정) / step08-noautocreate(운영 함정). 파일 상단 주석에 각각이 재현하는 절과 application.yml 에 추가할 두 줄이 정리돼 있습니다.
  • step08-default 는 반드시 제일 먼저 한 번만 실행하세요. 기본 접미사로 orders-retry-1000, orders-dlt 를 만드는 프로필입니다. 확인 후 kt --delete 로 지워야 이후 프로필의 kt --list 결과가 교재와 일치합니다.
  • [8-6] OrderStateDemo 가 이 파일의 핵심입니다. ORD-0001 로 3건을 발행하고 AtomicBoolean FIRST_TRY첫 CREATED 만 1회 실패시킵니다. STATE 맵의 최종값을 ApplicationRunner 가 3초 뒤에 찍으므로, 마지막 줄이 CANCELLED 가 아니라 CREATED 인지 반드시 확인하세요.
  • [8-8] LongBackoffDemodelay = 30000 으로 두고 3건을 연속 실패시킵니다. logging.level.org.springframework.kafka.listener.KafkaConsumerBackoffManager=DEBUG 를 켜야 Backing off partition ... 로그가 보입니다. 파일 주석에 그 설정 줄이 적혀 있습니다.
  • [8-11] NoAutoCreateDemoautoCreateTopics = "false" 로 두고 토픽을 일부러 만들지 않은 상태로 실행합니다. 60초 타임아웃을 실제로 겪어 보는 프로필이라 인내가 필요합니다. 확인 후 사전 생성 스크립트를 돌리고 다시 실행해 정상 동작을 대조하세요.
  • RemoteApiException / InvalidOrderException 은 Step 07 과 같은 이름이지만 step08 패키지에 다시 선언했습니다. 두 스텝의 파일을 같은 프로젝트에 두면 패키지가 달라 충돌하지 않습니다.
package com.example.order.step08;

/*
 * ============================================================================
 * Step 08 — @RetryableTopic 논블로킹 재시도 : Practice
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step08/Practice.java
 *
 * ★ 사전 준비 — application.yml 에 두 줄을 추가하세요.
 *
 *   spring.kafka.consumer.properties.spring.json.trusted.packages:
 *       "com.example.order.domain,com.example.order.step08"
 *       → [8-6] 에서 쓰는 OrderEvent 가 step08 패키지에 있기 때문입니다.
 *         빼먹으면 "The class is not in the trusted packages" 로 전부 실패합니다(Step 04).
 *
 *   logging.level.org.springframework.kafka.listener.KafkaConsumerBackoffManager: DEBUG
 *       → [8-8] 의 "Backing off partition ..." 로그를 보려면 필요합니다.
 *
 * 실행 (보조 프로필을 하나만 함께 켭니다. 예제끼리 서로 간섭합니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-default'
 *       → [8-2] 기본 접미사. orders-retry-1000 / -2000 / -4000 / orders-dlt 가 생깁니다.
 *         ★ 확인 후 반드시 kt --delete 로 지우세요. 안 지우면 이후 kt --list 가 교재와 다릅니다.
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-journey'
 *       → [8-3][8-5][8-7] 실패 한 건의 여정 추적 + @DltHandler + 컨테이너 5개 확인
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-order'
 *       → [8-6] ★★핵심★★ 같은 키 3건으로 순서가 깨지는 것을 재현
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-longbackoff'
 *       → [8-8] 긴 백오프가 retry 토픽 파티션을 pause 로 막는 것
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-hybrid'
 *       → [8-10] 짧은 블로킹 재시도 뒤에 논블로킹으로 넘기기
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-builder'
 *       → [8-10] 애노테이션 대신 RetryTopicConfiguration 빈으로 전역 설정
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,step08-noautocreate'
 *       → [8-11] autoCreateTopics=false 인데 토픽이 없을 때. 60초 타임아웃을 겪습니다.
 *
 * ★ 매 실행 전 오프셋 리셋 (앱을 먼저 종료할 것)
 *   kcg --group s08-inventory   --topic orders --reset-offsets --to-earliest --execute
 *   kcg --group s08-order-state --all-topics   --reset-offsets --to-earliest --execute
 *
 * ★ 프로필을 바꿀 때는 retry 토픽을 지우세요. 접미사 전략이 다르면 이름이 섞입니다.
 *   kt --list | grep -E 'orders-retry|orders-dlt' | xargs -I{} \
 *     docker exec -i learn-kafka /opt/kafka/bin/kafka-topics.sh \
 *       --bootstrap-server localhost:9092 --delete --topic {}
 *
 * 실행 중에 확인할 CLI
 *   kt  --list                                   ← ★ 8-2 의 핵심 확인
 *   kt  --describe --topic orders-retry-0        ← 파티션 수가 원본과 같은지
 *   kcg --list                                   ← ★ 그룹이 하나뿐인 것 (8-7)
 *   kcg --describe --group s08-inventory         ← 5개 토픽이 한 그룹에
 *   kcc --topic orders-retry-1 --from-beginning --property print.headers=true
 * ============================================================================
 */

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.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.DltHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.RetryableTopic;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.retrytopic.DeadLetterPublishingRecovererFactory;
import org.springframework.kafka.retrytopic.DltStrategy;
import org.springframework.kafka.retrytopic.RetryTopicConfiguration;
import org.springframework.kafka.retrytopic.RetryTopicConfigurationBuilder;
import org.springframework.kafka.retrytopic.RetryTopicConfigurationSupport;
import org.springframework.kafka.retrytopic.SameIntervalTopicReuseStrategy;
import org.springframework.kafka.retrytopic.TopicSuffixingStrategy;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.retry.annotation.Backoff;
import org.springframework.stereotype.Component;
import org.springframework.util.backoff.FixedBackOff;

import java.math.BigDecimal;
import java.time.Instant;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Consumer;

/**
 * Step 08 의 모든 예제를 담은 단일 파일입니다.
 * 각 nested static class 는 본문 절 번호와 1:1 로 대응합니다.
 */
public final class Practice {

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

    // ========================================================================
    // [공통] 실패 조건과 시도 시각 기록
    // ========================================================================
    //
    // Step 07 과 같은 orders-1@42 를 실패 지점으로 씁니다.
    // ⚠️ 단, 논블로킹에서는 재시도가 "다른 토픽" 에서 일어나므로 offset 이 달라집니다.
    //    그래서 카운터의 키를 offset 이 아니라 orderId 로 잡아야 attempt 가 이어집니다.
    static final class FailurePolicy {

        static final int  FAIL_PARTITION = 1;
        static final long FAIL_OFFSET    = 42L;

        private static final Map<String, AtomicInteger> ATTEMPTS = new ConcurrentHashMap<>();
        private static final Map<String, Long>          FIRST_AT = new ConcurrentHashMap<>();

        /** 본선(orders) 에서만 판정합니다. retry 토픽에서는 무조건 실패시켜야 여정을 끝까지 봅니다. */
        static boolean shouldFail(ConsumerRecord<?, ?> record) {
            if (!"orders".equals(record.topic())) {
                return true;                    // retry 토픽으로 넘어온 것은 계속 실패
            }
            return record.partition() == FAIL_PARTITION && record.offset() == FAIL_OFFSET;
        }

        /** "attempt=2 t=1006ms topic=orders-retry-0 ORD-0043" 형태의 도장. */
        static String stamp(ConsumerRecord<?, ?> record, String orderId) {
            int attempt = ATTEMPTS.computeIfAbsent(orderId, k -> new AtomicInteger()).incrementAndGet();
            long first  = FIRST_AT.computeIfAbsent(orderId, k -> System.currentTimeMillis());
            long elapsed = System.currentTimeMillis() - first;
            return "attempt=%d t=%dms topic=%s %s".formatted(attempt, elapsed, record.topic(), orderId);
        }
    }

    // ========================================================================
    // [8-2] 기본 접미사 — 아무것도 지정하지 않으면 어떤 토픽이 생기는가
    // ========================================================================
    //
    // topicSuffixingStrategy 의 기본값은 SUFFIX_WITH_DELAY_VALUE 입니다.
    // 인덱스가 아니라 "백오프 지연값(ms)" 이 접미사로 붙습니다.
    //
    // 기대 결과 (kt --list)
    //   orders
    //   orders-retry-1000
    //   orders-retry-2000
    //   orders-retry-4000
    //   orders-dlt              ← ★ Step 07 의 orders.DLT 와 다른 토픽입니다
    //
    // ★ 확인 후 반드시 이 네 개를 kt --delete 로 지우세요.
    @Component
    @Profile("step08-default")
    public static class DefaultSuffixDemo {

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

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (FailurePolicy.shouldFail(record)) {
                throw new RemoteApiException("재고 API 타임아웃: " + record.value().orderId());
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event) {
            log.error("DLT 수신 {}", event.orderId());
        }
    }

    // ========================================================================
    // [8-3][8-5][8-7] 실패 한 건의 여정 추적
    // ========================================================================
    //
    // 코스 규약(orders-retry-0/-1/-2, orders.DLT)에 맞춘 설정입니다.
    //
    // ★ 로그에서 확인할 것
    //   ① 스레드 이름이 [ntainer#0-1-C-1] → #1 → #2 → #3 → #4 로 바뀐다 (컨테이너가 다르다)
    //   ② "처리 완료 orders-1@43" 이 실패 직후 곧바로 찍힌다 (본선이 안 막혔다)
    //   ③ t=1006ms → 3012ms → 7021ms (1초 → 2초 → 4초)
    //   ④ 마지막이 orders.DLT-1 (원본 파티션 번호 1 유지)
    //
    // ★ CLI 로 확인할 것
    //   kcg --list                            → s08-inventory 하나뿐
    //   kcg --describe --group s08-inventory  → 5개 토픽이 한 그룹에
    @Component
    @Profile("step08-journey")
    public static class JourneyDemo {

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

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",              // 원본과 같게. 기본값 -1 은 브로커 기본을 따릅니다
                replicationFactor = "1",
                concurrency = "1",                // [8-7] retry 컨테이너만 1 스레드로
                dltStrategy = DltStrategy.FAIL_ON_ERROR,
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            OrderCreated event = record.value();
            if (FailurePolicy.shouldFail(record)) {
                log.warn("★ 실패 {}", FailurePolicy.stamp(record, event.orderId()));
                throw new RemoteApiException("재고 API 타임아웃: " + event.orderId());
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        /**
         * [8-5] 같은 클래스 안에 @DltHandler 하나면 DLT 리스너가 됩니다.
         *
         * ⚠️ 헤더 상수가 Step 07 과 다릅니다.
         *    Step 07 : KafkaHeaders.DLT_ORIGINAL_TOPIC  → "kafka_dlt-original-topic"
         *    Step 08 : KafkaHeaders.ORIGINAL_TOPIC      → "kafka_original-topic"
         *
         * ★ topic 이 "orders-retry-2" 가 아니라 "orders" 로 찍힙니다.
         *   헤더는 최초 발행 시 한 번만 붙고 이후 단계에서 덮어쓰이지 않기 때문입니다.
         */
        @DltHandler
        public void onDlt(OrderCreated event,
                          @Header(KafkaHeaders.ORIGINAL_TOPIC)     String topic,
                          @Header(KafkaHeaders.ORIGINAL_PARTITION) int    partition,
                          @Header(KafkaHeaders.ORIGINAL_OFFSET)    long   offset,
                          @Header(KafkaHeaders.EXCEPTION_FQCN)     String exFqcn,
                          @Header(KafkaHeaders.EXCEPTION_MESSAGE)  String exMessage) {

            log.error("DLT 수신 {} origin={}-{}@{} ex={} msg={}",
                    event.orderId(), topic, partition, offset,
                    exFqcn.substring(exFqcn.lastIndexOf('.') + 1), exMessage);
        }
    }

    // ========================================================================
    // [8-6] ★★★ 핵심 함정 ★★★ 논블로킹 재시도가 메시지 순서를 깨뜨린다
    // ========================================================================
    //
    // 같은 키 ORD-0001 로 3건을 순서대로 발행합니다.
    //   CREATED(offset 100) → UPDATED(101) → CANCELLED(102)
    // 키가 같으니 셋 다 같은 파티션에 들어가고, Kafka 는 파티션 안의 순서를 보장합니다.
    //
    // 그런데 첫 번째(CREATED)만 1회 실패하게 두면:
    //   t=8ms    CREATED 실패    → orders-retry-0 으로 발행, 오프셋 101 커밋
    //   t=12ms   UPDATED 성공    → STATE = UPDATED
    //   t=15ms   CANCELLED 성공  → STATE = CANCELLED
    //   t=1012ms CREATED 재처리 성공 → STATE = CREATED     ← ★★ 취소된 주문이 부활
    //
    // 예외 없음. LAG 0. ERROR 로그 없음. @DltHandler 도 호출 안 됨.
    // "재시도가 성공했기 때문에" 벌어진 오염입니다.
    //
    // ★ ApplicationRunner 가 3초 뒤 최종 상태를 찍습니다.
    //   마지막 줄이 "최종 상태 ORD-0001 = CREATED" 인지 반드시 확인하세요.
    @Component
    @Profile("step08-order")
    public static class OrderStateDemo {

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

        /** 주문 상태 저장소. 실무의 DB 라고 생각하세요. */
        static final Map<String, OrderState> STATE = new ConcurrentHashMap<>();

        /** 최초 CREATED 한 번만 실패시킵니다. 일시적 장애를 흉내 냅니다. */
        static final AtomicBoolean FIRST_TRY = new AtomicBoolean(true);

        static final long T0 = System.currentTimeMillis();

        static long elapsed() {
            return System.currentTimeMillis() - T0;
        }

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-order-state")
        public void onMessage(ConsumerRecord<String, OrderEvent> record) {

            OrderEvent event = record.value();

            // 최초 CREATED 만 1회 실패. 두 번째 시도(retry-0)에서는 성공합니다.
            if (event.status() == OrderState.CREATED && FIRST_TRY.getAndSet(false)) {
                log.warn("★ 실패 {} {} (topic={} t={}ms)",
                        event.orderId(), event.status(), record.topic(), elapsed());
                throw new RemoteApiException("재고 API 타임아웃");
            }

            STATE.put(event.orderId(), event.status());
            log.info("상태 반영 {} → {}  (topic={} t={}ms)",
                    event.orderId(), event.status(), record.topic(), elapsed());
        }

        @DltHandler
        public void onDlt(OrderEvent event) {
            log.error("DLT 수신 {} {}", event.orderId(), event.status());
        }
    }

    /** ORD-0001 의 세 이벤트를 순서대로 발행하고, 3초 뒤 최종 상태를 찍습니다. */
    @Component
    @Profile("step08-order")
    public static class OrderStatePublisher implements ApplicationRunner {

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

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            Thread t = new Thread(() -> {
                sleep(2_000L);   // 리스너 5개가 전부 partitions assigned 될 때까지 기다립니다

                for (OrderState status : new OrderState[]{
                        OrderState.CREATED, OrderState.UPDATED, OrderState.CANCELLED}) {

                    OrderEvent event = OrderEvent.of("ORD-0001", status);
                    template.send("orders", event.orderId(), event);
                    log.info("발행 {} {}  (t={}ms)",
                            event.orderId(), status, OrderStateDemo.elapsed());
                }

                sleep(3_000L);   // 백오프 1초 + 여유
                log.info("최종 상태 {} = {}   ← ★ CANCELLED 가 아니면 순서가 깨진 것입니다",
                        "ORD-0001", OrderStateDemo.STATE.get("ORD-0001"));
            }, "order-state-publisher");
            t.setDaemon(true);
            t.start();
        }

        private static void sleep(long ms) {
            try {
                Thread.sleep(ms);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }

    // ========================================================================
    // [8-8] ⚠️ 함정 — 긴 백오프는 retry 토픽 안에서 다시 블로킹이다
    // ========================================================================
    //
    // 백오프는 스케줄러가 아닙니다. retry 토픽 컨슈머는 메시지를 즉시 poll 해 오고,
    // KafkaBackoffAwareMessageListenerAdapter 가 retry_topic-backoff-timestamp 헤더를
    // 보고 "아직 시간이 안 됐다" 면 KafkaBackoffException 을 던집니다.
    // 그러면 KafkaConsumerBackoffManager 가 그 파티션을 pause + seek 합니다.
    //
    // ★ pause 이므로 poll() 자체는 계속 돕니다 → max.poll.interval.ms 초과는 안 납니다.
    //   하지만 그 파티션의 뒤 레코드는 그대로 기다립니다. Step 07 과 같은 병목입니다.
    //
    // ★ 로그를 보려면 application.yml 에 이 줄이 필요합니다.
    //   logging.level.org.springframework.kafka.listener.KafkaConsumerBackoffManager: DEBUG
    //
    // 기대 로그
    //   DEBUG ... KafkaConsumerBackoffManager : Backing off partition orders-retry-0-1 until ...
    //   INFO  ... KafkaMessageListenerContainer : Paused consumption from: [orders-retry-0-1]
    @Component
    @Profile("step08-longbackoff")
    public static class LongBackoffDemo {

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

        private static final Map<String, Long> SENT_AT = new ConcurrentHashMap<>();

        @RetryableTopic(
                attempts = "3",
                // 30초. random=true 로 두면 due time 이 흩어져 "역전" 이 생깁니다(본문 8-8 표).
                backoff = @Backoff(delay = 30_000, multiplier = 1.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                // multiplier=1.0 이라 모든 단계의 대기가 같습니다.
                // SINGLE_TOPIC 으로 두면 orders-retry 하나로 합쳐집니다.
                sameIntervalTopicReuseStrategy = SameIntervalTopicReuseStrategy.SINGLE_TOPIC,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {

            String orderId = record.value().orderId();

            if ("orders".equals(record.topic())) {
                // 본선에서 연속 3건(ORD-0043, 0044, 0045)을 실패시켜 retry 토픽에 몰아넣습니다.
                if (record.partition() == 1 && record.offset() >= 42 && record.offset() <= 44) {
                    SENT_AT.put(orderId, System.currentTimeMillis());
                    log.warn("★ 실패 {} → retry 토픽으로", orderId);
                    throw new RemoteApiException("재고 API 타임아웃");
                }
                log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
                return;
            }

            // retry 토픽에서는 성공시키고, 실제로 몇 ms 기다렸는지 찍습니다.
            long waited = System.currentTimeMillis() - SENT_AT.getOrDefault(orderId, System.currentTimeMillis());
            log.info("재처리 {} (topic={} 대기 {}ms)   ← 설정값 30,000ms 와 비교하세요",
                    orderId, record.topic(), waited);
        }

        @DltHandler
        public void onDlt(OrderCreated event) {
            log.error("DLT 수신 {}", event.orderId());
        }
    }

    // ========================================================================
    // [8-10] 블로킹 + 논블로킹 조합
    // ========================================================================
    //
    // 대부분의 순간 장애는 500ms 뒤 재시도로 끝납니다.
    // 그런 실패까지 retry 토픽을 거치면 토픽만 지저분해지고 순서도 깨집니다.
    //
    // configureBlockingRetries 에 등록한 예외는 "먼저 컨테이너에서 블로킹으로" 재시도합니다.
    // 여기서 성공하면 retry 토픽에 아무것도 안 남고 순서도 유지됩니다.
    //
    // ⚠️ 총 시도는 합이지만 "대기 시간" 은 곱해집니다.
    //    retry-0 컨슈머에서도 블로킹 2회를 다시 돌고 나서 retry-1 로 갑니다.
    //    그래서 블로킹 백오프는 총합 1초 이내로 아주 짧게 두어야 합니다.
    @Configuration
    @Profile("step08-hybrid")
    public static class HybridRetryConfig extends RetryTopicConfigurationSupport {

        /** ★ 논블로킹으로 넘기기 전에 컨테이너에서 0.5초 × 2회만 블로킹 재시도. */
        @Override
        protected void configureBlockingRetries(BlockingRetriesConfigurer blockingRetries) {
            blockingRetries
                    .retryOn(RemoteApiException.class)
                    .backOff(new FixedBackOff(500L, 2L));
        }

        /**
         * DLT/retry 발행기를 커스터마이징하는 훅입니다.
         * 반환 타입이 Consumer&lt;DeadLetterPublishingRecovererFactory&gt; 인 것에 주의하세요.
         * 실제 recoverer 를 손대려면 팩토리의 customizer 를 한 겹 더 등록해야 합니다.
         */
        @Override
        protected Consumer<DeadLetterPublishingRecovererFactory> configureDeadLetterPublishingContainerFactory() {
            return factory -> factory.setDeadLetterPublishingRecovererCustomizer(
                    recoverer -> recoverer.setStripPreviousExceptionHeaders(true));
        }
    }

    @Component
    @Profile("step08-hybrid")
    public static class HybridDemo {

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

        /** 5번째 시도에서 성공하도록 만듭니다. 블로킹 3회 → 논블로킹 2회에 걸칩니다. */
        private static final AtomicInteger ATTEMPT = new AtomicInteger();

        @RetryableTopic(
                attempts = "3",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {

            if (!(record.partition() == 1 && record.offset() == 42) && "orders".equals(record.topic())) {
                log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
                return;
            }

            String orderId = record.value().orderId();
            int attempt = ATTEMPT.incrementAndGet();

            if (attempt < 5) {
                log.warn("실패 attempt={} topic={} {}", attempt, record.topic(), orderId);
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("성공 attempt={} topic={} {}", attempt, record.topic(), orderId);
        }

        @DltHandler
        public void onDlt(OrderCreated event) {
            log.error("DLT 수신 {}", event.orderId());
        }
    }

    // ========================================================================
    // [8-10] 애노테이션 대신 RetryTopicConfiguration 빈으로 전역 설정
    // ========================================================================
    //
    // includeTopic("orders") 하나로, orders 를 구독하는 "모든" @KafkaListener 에
    // 같은 재시도 정책이 적용됩니다. 애노테이션을 하나도 안 붙여도 됩니다.
    //
    // ⚠️ 애노테이션과 빈이 둘 다 있으면 애노테이션이 이깁니다.
    @Configuration
    @Profile("step08-builder")
    public static class BuilderConfig {

        @Bean
        public RetryTopicConfiguration ordersRetryTopicConfig(KafkaTemplate<String, Object> template) {
            return RetryTopicConfigurationBuilder
                    .newInstance()
                    .includeTopic("orders")
                    .maxAttempts(4)                          // = attempts = "4"
                    .exponentialBackoff(1000, 2.0, 10_000)   // = @Backoff(1000, 2.0, maxDelay=10000)
                    .suffixTopicsWithIndexValues()           // = SUFFIX_WITH_INDEX_VALUE
                    .retryTopicSuffix("-retry")
                    .dltSuffix(".DLT")
                    .concurrency(1)
                    .autoCreateTopics(true, 3, (short) 1)    // = numPartitions/replicationFactor
                    .notRetryOn(InvalidOrderException.class) // = exclude
                    .doNotRetryOnDltFailure()                // = DltStrategy.FAIL_ON_ERROR
                    .create(template);
        }
    }

    /** 애노테이션이 하나도 없습니다. 위 빈이 대신 적용됩니다. */
    @Component
    @Profile("step08-builder")
    public static class BuilderInventoryListener {

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

        @KafkaListener(topics = "orders", groupId = "s08-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (FailurePolicy.shouldFail(record)) {
                log.warn("★ 실패 {}", FailurePolicy.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("재고 처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }
    }

    /** 같은 토픽을 구독하는 두 번째 리스너. 여기에도 같은 정책이 적용됩니다. */
    @Component
    @Profile("step08-builder")
    public static class BuilderNotificationListener {

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

        @KafkaListener(topics = "orders", groupId = "s08-notification")
        public void notifyCustomer(ConsumerRecord<String, OrderCreated> record) {
            log.info("알림 발송 완료 {}", record.key());
        }
    }

    // ========================================================================
    // [8-11] ⚠️ 함정 — 토픽 자동 생성이 꺼진 운영 환경
    // ========================================================================
    //
    // autoCreateTopics = "false" 는 "토픽을 만들지 않는다" 일 뿐,
    // "토픽이 있는지 검증한다" 가 아닙니다. 없어도 기동은 성공합니다.
    //
    // ★ 재현 순서
    //   1) retry/DLT 토픽을 전부 지웁니다.
    //   2) 이 프로필로 기동합니다 → 정상 기동. partitions assigned 도 찍힙니다.
    //   3) 첫 실패가 나면 60초 침묵 뒤:
    //      ERROR ... DeadLetterPublishingRecoverer : Dead-letter publication to
    //              orders-retry-0 failed for: orders-1@42
    //      org.apache.kafka.common.errors.TimeoutException:
    //              Topic orders-retry-0 not present in metadata after 60000 ms.
    //   4) 그 60초 동안 orders-1 파티션은 완전히 멈춥니다. 논블로킹의 의미가 사라집니다.
    //   5) 아래 스크립트로 토픽을 만들고 다시 실행해 정상 동작과 대조하세요.
    //
    //      for i in 0 1 2; do
    //        kt --create --topic orders-retry-$i --partitions 3 --replication-factor 1
    //      done
    //      kt --create --topic orders.DLT --partitions 3 --replication-factor 1
    @Component
    @Profile("step08-noautocreate")
    public static class NoAutoCreateDemo {

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

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                autoCreateTopics = "false",        // ★ 운영 필수. 대신 사전 생성이 전제입니다
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (FailurePolicy.shouldFail(record)) {
                log.warn("★ 실패 {} — 여기서 60초를 기다리게 됩니다",
                        FailurePolicy.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event) {
            log.error("DLT 수신 {}", event.orderId());
        }
    }
}

// ============================================================================
// 이 스텝에서 쓰는 보조 타입들
// ============================================================================

/** 외부 API 일시 장애. "5초 뒤에 하면 될 수도 있다" → 재시도 대상. */
class RemoteApiException extends RuntimeException {
    RemoteApiException(String message) {
        super(message);
    }
}

/**
 * 비즈니스 검증 실패. 100번 해도 100번 실패합니다.
 * @RetryableTopic(exclude = InvalidOrderException.class) 로 즉시 DLT 로 보냅니다.
 */
class InvalidOrderException extends RuntimeException {
    InvalidOrderException(String message) {
        super(message);
    }
}

/** 주문 상태. [8-6] 의 순서 깨짐을 보여 주기 위한 최소 집합입니다. */
enum OrderState {
    CREATED, UPDATED, CANCELLED
}

/**
 * [8-6] 전용 이벤트. OrderCreated 에는 상태 필드가 없어서 따로 둡니다.
 *
 * ⚠️ 이 타입은 com.example.order.step08 패키지에 있으므로 application.yml 의
 *    spring.json.trusted.packages 에 "com.example.order.step08" 을 추가해야 합니다.
 *    빼먹으면 "The class ... is not in the trusted packages" 로 전부 실패합니다(Step 04).
 */
record OrderEvent(String orderId, OrderState status, BigDecimal amount, Instant at) {

    static OrderEvent of(String orderId, OrderState status) {
        return new OrderEvent(orderId, status, new BigDecimal("11000"), Instant.now());
    }
}

Exercise.java

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

  • 문제 1·2 는 토픽 이름을 관찰하는 문제이고, 문제 3·5 는 애노테이션 속성을 설정하는 문제, 문제 4 는 순서 깨짐을 재현하는 문제, 문제 6 은 빈으로 재구성하는 문제입니다. 순서대로 푸세요.
  • ⚠️ 문제 1 과 문제 2 는 접미사 전략이 달라 서로 다른 토픽을 만듭니다. 문제 1 을 확인한 뒤 kt --delete --topic orders-retry-1000 식으로 반드시 지우고 문제 2 로 넘어가세요. 안 지우면 kt --list 에 7개가 섞여 보입니다.
  • 문제 4 는 코드를 거의 안 씁니다. // 관측 기록: 주석의 타임라인 표를 채우는 문제입니다. 각 행의 "실제 처리 시각"과 "STATE 값"을 로그에서 옮겨 적고, 마지막 줄에 최종 상태를 적으세요.
  • 각 문제 끝의 // 확인: 주석에 기대 로그 한 줄이 적혀 있습니다. 예를 들어 문제 5 의 확인 줄은 Republishing failed record to orders.DLT-1 이고, orders-retry-0 이 로그에 한 번도 나오면 안 됩니다.
  • 문제 6 은 @RetryableTopic 애노테이션을 전부 지운 상태에서 풀어야 합니다. 애노테이션이 남아 있으면 그쪽이 이겨서 빈이 무시되고, 답이 틀렸는지 맞았는지 구분이 안 됩니다.
package com.example.order.step08;

/*
 * ============================================================================
 * Step 08 — @RetryableTopic 논블로킹 재시도 : Exercise (6문제)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step08/Exercise.java
 *
 * ★ 사전 준비 — application.yml
 *   spring.kafka.consumer.properties.spring.json.trusted.packages:
 *       "com.example.order.domain,com.example.order.step08"
 *
 * 실행 (문제마다 보조 프로필을 하나만 켭니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,ex08-q1'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,ex08-q2'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,ex08-q3'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,ex08-q4'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,ex08-q5'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,ex08-q6'
 *
 * ★ 매 실행 전 오프셋 리셋 (앱을 먼저 종료할 것)
 *   kcg --group s08-ex-inventory   --all-topics --reset-offsets --to-earliest --execute
 *   kcg --group s08-ex-order-state --all-topics --reset-offsets --to-earliest --execute
 *
 * ⚠️ 문제 1 과 문제 2 는 접미사 전략이 달라 "서로 다른 토픽" 을 만듭니다.
 *    문제 1 을 확인한 뒤 반드시 지우고 문제 2 로 넘어가세요.
 *      kt --delete --topic orders-retry-1000
 *      kt --delete --topic orders-retry-2000
 *      kt --delete --topic orders-retry-4000
 *      kt --delete --topic orders-dlt
 *    안 지우면 kt --list 에 7개가 섞여 보여 답을 확인할 수 없습니다.
 *
 * ⚠️ 문제 6 은 @RetryableTopic 애노테이션을 "전부 지운 상태" 에서 풀어야 합니다.
 *    애노테이션이 남아 있으면 그쪽이 이겨서 빈이 무시되고, 정답 여부를 구분할 수 없습니다.
 * ============================================================================
 */

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.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.DltHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.RetryableTopic;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.retrytopic.RetryTopicConfiguration;
import org.springframework.kafka.retrytopic.RetryTopicConfigurationBuilder;
import org.springframework.kafka.retrytopic.TopicSuffixingStrategy;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.retry.annotation.Backoff;
import org.springframework.stereotype.Component;

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;

public final class Exercise {

    private Exercise() {
    }

    // ------------------------------------------------------------------
    // 공통 — 실패 조건과 시도 시각 (수정하지 마세요)
    // ------------------------------------------------------------------
    static final class Fail {

        static final int  PARTITION = 1;
        static final long OFFSET    = 42L;

        private static final Map<String, AtomicInteger> ATTEMPTS = new ConcurrentHashMap<>();
        private static final Map<String, Long>          FIRST_AT = new ConcurrentHashMap<>();

        /** 본선에서는 orders-1@42 만, retry 토픽에서는 전부 실패시킵니다. */
        static boolean shouldFail(ConsumerRecord<?, ?> r) {
            if (!"orders".equals(r.topic())) {
                return true;
            }
            return r.partition() == PARTITION && r.offset() == OFFSET;
        }

        static String stamp(ConsumerRecord<?, ?> r, String orderId) {
            int attempt = ATTEMPTS.computeIfAbsent(orderId, k -> new AtomicInteger()).incrementAndGet();
            long first  = FIRST_AT.computeIfAbsent(orderId, k -> System.currentTimeMillis());
            return "attempt=%d t=%dms topic=%s %s".formatted(
                    attempt, System.currentTimeMillis() - first, r.topic(), orderId);
        }
    }

    // ==================================================================
    // 문제 1. @RetryableTopic 을 붙이고 생성된 토픽 이름을 관찰하기
    // ==================================================================
    //
    // 요구사항
    //   - attempts = "4"
    //   - backoff  = @Backoff(delay = 1000, multiplier = 2.0)
    //   - kafkaTemplate = "kafkaTemplate"
    //   - 접미사나 접미사 전략은 "아무것도 지정하지 말 것" (기본값을 보는 것이 목적)
    //
    // 기동 후 `kt --list` 결과를 아래에 그대로 옮겨 적으세요.
    //
    //   // 관측 기록 (kt --list):
    //   //   orders
    //   //   ______________________
    //   //   ______________________
    //   //   ______________________
    //   //   ______________________
    //   //
    //   // Q. 접미사에 붙은 숫자는 무엇인가?  →  ______________________
    //   // Q. 그렇다면 topicSuffixingStrategy 의 기본값은?  →  ______________________
    //   // Q. 운영 중에 delay 를 1000 → 1500 으로 바꾸면 어떻게 되는가? → ______________________
    //
    // 확인: kt --list 에 "orders-dlt" 가 보여야 합니다 (orders.DLT 가 아닙니다)
    @Component
    @Profile("ex08-q1")
    public static class Q1DefaultSuffix {

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

        // 여기에 작성: @RetryableTopic 을 붙이세요

        @KafkaListener(topics = "orders", groupId = "s08-ex-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                log.warn("★ 실패 {}", Fail.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }
    }

    // ==================================================================
    // 문제 2. 코스 규약에 맞추고 백오프 시각을 실측하기
    // ==================================================================
    //
    // 요구사항
    //   - retry 토픽 이름이 orders-retry-0, orders-retry-1, orders-retry-2 가 되도록
    //   - DLT 이름이 orders.DLT 가 되도록 (Step 07 과 같은 토픽을 재사용합니다)
    //   - 파티션 3, 복제 1 로 만들어지도록 명시
    //   - 백오프 1s → 2s → 4s
    //
    // 로그의 t=Nms 를 아래에 적으세요.
    //   // 관측 기록:
    //   //   attempt=1 topic=orders          t=_______ms
    //   //   attempt=2 topic=orders-retry-0  t=_______ms
    //   //   attempt=3 topic=orders-retry-1  t=_______ms
    //   //   attempt=4 topic=orders-retry-2  t=_______ms
    //
    // 확인: kt --list 에 orders-retry-0/-1/-2 와 orders.DLT 가 있고 orders-dlt 는 없어야 합니다
    @Component
    @Profile("ex08-q2")
    public static class Q2CourseConvention {

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

        // 여기에 작성: @RetryableTopic (topicSuffixingStrategy, retryTopicSuffix,
        //              dltTopicSuffix, numPartitions, replicationFactor 를 명시)

        @KafkaListener(topics = "orders", groupId = "s08-ex-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                log.warn("★ 실패 {}", Fail.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }
    }

    // ==================================================================
    // 문제 3. @DltHandler 로 실패 원인 복원하기
    // ==================================================================
    //
    // 요구사항
    //   - 같은 클래스 안에 @DltHandler 메서드를 만들 것
    //   - @Header 로 다음 넷을 받아 로그로 남길 것
    //       KafkaHeaders.ORIGINAL_TOPIC      (String)
    //       KafkaHeaders.ORIGINAL_PARTITION  (int)
    //       KafkaHeaders.ORIGINAL_OFFSET     (long)
    //       KafkaHeaders.EXCEPTION_FQCN      (String)
    //
    //   // Q. ORIGINAL_TOPIC 에 무엇이 찍혔는가?  →  ______________________
    //   //    "orders-retry-2" 를 예상했다면 왜 틀렸는지 한 줄로 설명하세요.
    //   //    → ______________________________________________________
    //
    // 확인: "DLT 수신 ORD-0043 origin=orders-1@42" 가 찍혀야 합니다
    @Component
    @Profile("ex08-q3")
    public static class Q3DltHandler {

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

        @RetryableTopic(
                attempts = "3",
                backoff = @Backoff(delay = 500, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-ex-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                throw new RemoteApiException("재고 API 타임아웃: " + record.value().orderId());
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        // 여기에 작성: @DltHandler 메서드
    }

    // ==================================================================
    // 문제 4. ★핵심★ 같은 키 3건으로 순서 깨짐을 재현하기
    // ==================================================================
    //
    // 요구사항
    //   - 리스너에 @RetryableTopic 을 붙일 것 (attempts=4, delay=1000, multiplier=2.0,
    //     인덱스 접미사, dltTopicSuffix=".DLT")
    //   - 아래 onMessage 는 이미 "첫 CREATED 만 1회 실패" 하도록 되어 있습니다. 수정하지 마세요.
    //   - Q4Publisher 가 ORD-0002 의 3건을 순서대로 발행합니다.
    //
    // 로그를 보고 타임라인 표를 채우세요. 코드는 거의 안 씁니다.
    //
    //   // 관측 기록:
    //   // | 시각      | 토픽            | 이벤트    | 결과            | STATE          |
    //   // |----------|----------------|----------|----------------|----------------|
    //   // | t=____ms | orders         | CREATED  | ______________ | ______________ |
    //   // | t=____ms | orders         | UPDATED  | ______________ | ______________ |
    //   // | t=____ms | orders         | CANCELLED| ______________ | ______________ |
    //   // | t=____ms | ______________ | CREATED  | ______________ | ______________ |
    //   //
    //   // 최종 상태 ORD-0002 = ______________
    //   //
    //   // Q. ERROR 로그가 몇 줄 찍혔는가?  →  ______
    //   // Q. @DltHandler 는 호출됐는가?    →  ______
    //   // Q. kcg --describe 의 LAG 은?     →  ______
    //
    // 확인: 마지막 줄이 "최종 상태 ORD-0002 = CREATED" 여야 합니다 (CANCELLED 가 아닙니다)
    @Component
    @Profile("ex08-q4")
    public static class Q4OrderState {

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

        static final Map<String, OrderState> STATE = new ConcurrentHashMap<>();
        static final AtomicBoolean FIRST_TRY = new AtomicBoolean(true);
        static final long T0 = System.currentTimeMillis();

        static long elapsed() {
            return System.currentTimeMillis() - T0;
        }

        // 여기에 작성: @RetryableTopic

        @KafkaListener(topics = "orders", groupId = "s08-ex-order-state")
        public void onMessage(ConsumerRecord<String, OrderEvent> record) {
            OrderEvent event = record.value();
            if (event.status() == OrderState.CREATED && FIRST_TRY.getAndSet(false)) {
                log.warn("★ 실패 {} {} (topic={} t={}ms)",
                        event.orderId(), event.status(), record.topic(), elapsed());
                throw new RemoteApiException("재고 API 타임아웃");
            }
            STATE.put(event.orderId(), event.status());
            log.info("상태 반영 {} → {}  (topic={} t={}ms)",
                    event.orderId(), event.status(), record.topic(), elapsed());
        }

        @DltHandler
        public void onDlt(OrderEvent event) {
            log.error("DLT 수신 {} {}", event.orderId(), event.status());
        }
    }

    @Component
    @Profile("ex08-q4")
    public static class Q4Publisher implements ApplicationRunner {

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

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            Thread t = new Thread(() -> {
                sleep(2_000L);
                for (OrderState s : new OrderState[]{
                        OrderState.CREATED, OrderState.UPDATED, OrderState.CANCELLED}) {
                    OrderEvent e = OrderEvent.of("ORD-0002", s);
                    template.send("orders", e.orderId(), e);
                    log.info("발행 {} {} (t={}ms)", e.orderId(), s, Q4OrderState.elapsed());
                }
                sleep(3_000L);
                log.info("최종 상태 ORD-0002 = {}", Q4OrderState.STATE.get("ORD-0002"));
            }, "q4-publisher");
            t.setDaemon(true);
            t.start();
        }

        private static void sleep(long ms) {
            try {
                Thread.sleep(ms);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }

    // ==================================================================
    // 문제 5. 특정 예외는 재시도 없이 즉시 DLT
    // ==================================================================
    //
    // 요구사항
    //   - InvalidOrderException 은 재시도하지 말고 곧바로 orders.DLT 로 보낼 것
    //   - RemoteApiException 은 정상적으로 3번 재시도할 것
    //   - 힌트 1: @RetryableTopic 의 exclude 속성
    //   - 힌트 2: 리스너 예외는 ListenerExecutionFailedException 으로 "감싸여" 옵니다.
    //             원인 체인을 따라가게 하는 속성이 하나 더 필요합니다.
    //
    //   // Q. 힌트 2 의 속성 이름은?  →  ______________________
    //
    // 확인: orders-1@42(InvalidOrderException) 에 대해
    //         "Republishing failed record to orders.DLT-1" 이 곧바로 찍히고
    //       ★ "orders-retry-0" 이 로그에 단 한 번도 나오면 안 됩니다.
    //       kcc --topic orders-retry-0 --from-beginning 결과가 비어 있어야 합니다.
    @Component
    @Profile("ex08-q5")
    public static class Q5Exclude {

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

        // 여기에 작성: @RetryableTopic (exclude 와 원인 체인 탐색 속성 포함)

        @KafkaListener(topics = "orders", groupId = "s08-ex-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            OrderCreated event = record.value();

            // (A) 검증 실패 → 재시도 없이 즉시 DLT 여야 합니다
            if (record.partition() == 1 && record.offset() == 42) {
                log.warn("★ 검증 실패 {} (topic={})", event.orderId(), record.topic());
                throw new InvalidOrderException("quantity 는 1 이상이어야 합니다: " + event.orderId());
            }

            // (B) 외부 API 실패 → 정상적으로 재시도돼야 합니다
            if (record.partition() == 2 && record.offset() == 30) {
                log.warn("★ API 실패 {} (topic={})", event.orderId(), record.topic());
                throw new RemoteApiException("재고 API 타임아웃");
            }

            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event,
                          @Header(KafkaHeaders.EXCEPTION_FQCN) String exFqcn) {
            log.error("DLT 수신 {} ex={}", event.orderId(),
                    exFqcn.substring(exFqcn.lastIndexOf('.') + 1));
        }
    }

    // ==================================================================
    // 문제 6. RetryTopicConfiguration 빈으로 전역 설정하기
    // ==================================================================
    //
    // 요구사항
    //   - 아래 두 리스너에 @RetryableTopic 이 하나도 없습니다. 그대로 두세요.
    //   - RetryTopicConfiguration 빈 하나로 문제 2 와 같은 정책을 적용할 것
    //       includeTopic("orders") / maxAttempts 4 / 지수 1s·2배·상한 10s
    //       인덱스 접미사 / retry 접미사 "-retry" / DLT 접미사 ".DLT"
    //       concurrency 1 / 토픽 자동 생성 파티션 3 복제 1
    //       InvalidOrderException 은 재시도 제외
    //
    //   // Q. 두 리스너(s08-ex-inventory, s08-ex-notification) 모두에 적용됐는가? → ______
    //   // Q. 컨슈머 그룹은 몇 개인가? (kcg --list)                              → ______
    //
    // 확인: kt --list 에 orders-retry-0/-1/-2, orders.DLT 가 생겨야 합니다
    @Configuration
    @Profile("ex08-q6")
    public static class Q6BuilderConfig {

        @Bean
        public RetryTopicConfiguration ordersRetryTopicConfig(KafkaTemplate<String, Object> template) {
            // 여기에 작성: RetryTopicConfigurationBuilder 체인
            return RetryTopicConfigurationBuilder.newInstance().create(template);
        }
    }

    @Component
    @Profile("ex08-q6")
    public static class Q6Inventory {

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

        @KafkaListener(topics = "orders", groupId = "s08-ex-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                log.warn("★ 실패 {}", Fail.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("재고 처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }
    }

    @Component
    @Profile("ex08-q6")
    public static class Q6Notification {

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

        @KafkaListener(topics = "orders", groupId = "s08-ex-notification")
        public void notifyCustomer(ConsumerRecord<String, OrderCreated> record) {
            log.info("알림 발송 완료 {} (topic={})", record.key(), record.topic());
        }
    }
}

Solution.java

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

  • 정답 1 은 코드보다 관찰이 핵심입니다. orders-retry-1000/2000/4000 이라는 이름에서 접미사가 인덱스가 아니라 지연값이라는 것, 따라서 백오프를 바꾸면 토픽 이름 자체가 바뀐다는 결론을 이끌어 냅니다. 운영 중 백오프 변경이 왜 위험한지도 여기서 설명합니다.
  • 정답 2SUFFIX_WITH_INDEX_VALUE + dltTopicSuffix = ".DLT" 조합입니다. retryTopicSuffix 는 기본값과 같아 생략해도 되지만 명시하는 편이 낫다는 이유(팀원이 기본값을 외우고 있지 않다)를 주석에 적었습니다.
  • 정답 3 의 포인트는 ORIGINAL_TOPICorders-retry-2 가 아니라 orders 로 찍힌다는 점입니다. 헤더는 최초 발행 시 한 번만 붙고 이후 단계에서 덮어쓰이지 않기 때문입니다. Step 07 의 DLT_ORIGINAL_TOPIC 과 상수 이름이 다른 이유도 함께 설명합니다.
  • 정답 4 는 실측 타임라인표와 결론입니다. 최종 상태가 CREATED 인 것, 그리고 이 문제를 실패로 만드는 것은 예외가 아니라 성공이라는 역설을 짚습니다. 해결책 세 가지(블로킹 전환 / 버전 필드 / 스냅샷 이벤트)를 코드 스케치로 붙였습니다.
  • 정답 5exclude = { InvalidOrderException.class } 하나면 충분하다는 것입니다. traversingCauses = "true" 가 왜 필요한지 — 리스너 예외는 ListenerExecutionFailedException 으로 감싸여 오므로 원인 체인을 따라가야 분류가 맞는다는 점 — 을 자세히 적었습니다. Step 07 의 addNotRetryableExceptions 와 같은 함정입니다.
  • 정답 6RetryTopicConfigurationBuilder 체인입니다. .suffixTopicsWithIndexValues(), .autoCreateTopics(true, 3, (short) 1), .notRetryOn(...) 이 각각 애노테이션의 어느 속성에 대응하는지 1:1 대조표를 주석 표로 넣었습니다. 그리고 애노테이션이 빈보다 우선한다는 규칙도 마지막에 적었습니다.
package com.example.order.step08;

/*
 * ============================================================================
 * Step 08 — @RetryableTopic 논블로킹 재시도 : Solution (6문제 정답 + 해설)
 * ============================================================================
 *
 * 배치 위치
 *   spring-kafka-lab/src/main/java/com/example/order/step08/Solution.java
 *
 * 실행 (Exercise.java 와 프로필 이름이 다릅니다. 둘 다 두어도 충돌하지 않습니다)
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,sol08-q1'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,sol08-q2'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,sol08-q3'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,sol08-q4'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,sol08-q5'
 *   ./gradlew bootRun --args='--spring.profiles.active=step08,sol08-q6'
 *
 * ★ 매 실행 전 오프셋 리셋 + retry 토픽 정리
 *   kcg --group s08-sol-inventory   --all-topics --reset-offsets --to-earliest --execute
 *   kcg --group s08-sol-order-state --all-topics --reset-offsets --to-earliest --execute
 * ============================================================================
 */

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.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Profile;
import org.springframework.kafka.annotation.DltHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.RetryableTopic;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.retrytopic.RetryTopicConfiguration;
import org.springframework.kafka.retrytopic.RetryTopicConfigurationBuilder;
import org.springframework.kafka.retrytopic.TopicSuffixingStrategy;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.retry.annotation.Backoff;
import org.springframework.stereotype.Component;

import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;

public final class Solution {

    private Solution() {
    }

    static final class Fail {

        static final int  PARTITION = 1;
        static final long OFFSET    = 42L;

        private static final Map<String, AtomicInteger> ATTEMPTS = new ConcurrentHashMap<>();
        private static final Map<String, Long>          FIRST_AT = new ConcurrentHashMap<>();

        static boolean shouldFail(ConsumerRecord<?, ?> r) {
            if (!"orders".equals(r.topic())) {
                return true;
            }
            return r.partition() == PARTITION && r.offset() == OFFSET;
        }

        static String stamp(ConsumerRecord<?, ?> r, String orderId) {
            int attempt = ATTEMPTS.computeIfAbsent(orderId, k -> new AtomicInteger()).incrementAndGet();
            long first  = FIRST_AT.computeIfAbsent(orderId, k -> System.currentTimeMillis());
            return "attempt=%d t=%dms topic=%s %s".formatted(
                    attempt, System.currentTimeMillis() - first, r.topic(), orderId);
        }
    }

    // ==================================================================
    // 정답 1. 기본 접미사 관찰
    // ==================================================================
    /*
     * 왜 이 답인가
     *
     * 코드 자체는 애노테이션 세 줄이 전부입니다. 이 문제의 목적은 "코드" 가 아니라
     * "결과를 눈으로 확인하는 것" 입니다.
     *
     * 관측 결과 (kt --list)
     *   orders
     *   orders-retry-1000
     *   orders-retry-2000
     *   orders-retry-4000
     *   orders-dlt
     *
     * ① 접미사에 붙은 숫자는 인덱스(0,1,2)가 아니라 "그 단계의 백오프 지연값(ms)" 입니다.
     *    1000 / 2000 / 4000 은 delay=1000, multiplier=2.0 의 결과와 정확히 일치합니다.
     *
     * ② 따라서 topicSuffixingStrategy 의 기본값은 SUFFIX_WITH_DELAY_VALUE 입니다.
     *    많은 블로그가 "orders-retry-0 이 생긴다" 고 쓰는데, 그건 SUFFIX_WITH_INDEX_VALUE 를
     *    명시했을 때의 이야기입니다.
     *
     * ③ ★ 가장 중요한 함의 — 백오프를 바꾸면 "토픽 이름이 바뀝니다."
     *    delay 를 1000 → 1500 으로 조정하는 순간 orders-retry-1500 / 3000 / 6000 이
     *    새로 생기고, 기존 orders-retry-1000 에 남아 있던 미처리 메시지는
     *    "아무도 구독하지 않는 토픽" 에 갇힙니다. 컨슈머가 없으니 LAG 지표에도 안 잡힙니다.
     *    에러 없이 메시지가 증발하는 전형적인 경로입니다.
     *    → 그래서 이 코스는 SUFFIX_WITH_INDEX_VALUE 를 권합니다. 백오프를 바꿔도
     *      토픽 이름이 그대로라 기존 메시지가 그대로 소비됩니다.
     *
     * ④ DLT 접미사가 "-dlt" 인 것도 확인하세요. Step 07 의 "orders.DLT" 와 다른 토픽입니다.
     *    점 하나 차이로 기존 DLT 모니터링이 전부 무력화됩니다.
     */
    @Component
    @Profile("sol08-q1")
    public static class Q1DefaultSuffix {

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

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-sol-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                log.warn("★ 실패 {}", Fail.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event) {
            log.error("DLT 수신 {}", event.orderId());
        }
    }

    // ==================================================================
    // 정답 2. 코스 규약에 맞추기 + 백오프 실측
    // ==================================================================
    /*
     * 왜 이 답인가
     *
     * ① topicSuffixingStrategy = SUFFIX_WITH_INDEX_VALUE
     *    정답 1에서 본 이유 때문입니다. 백오프 변경에 토픽 이름이 흔들리지 않습니다.
     *
     * ② dltTopicSuffix = ".DLT"
     *    Step 07 에서 이미 orders.DLT 를 만들었고, DLT 모니터링·알림·재발행 도구가
     *    전부 그 이름을 보고 있습니다. 토픽을 새로 파는 대신 기존 것을 재사용합니다.
     *
     * ③ retryTopicSuffix = "-retry" 는 기본값과 "같습니다." 생략해도 결과가 같습니다.
     *    그래도 명시하는 편이 낫습니다. 리뷰어가 기본값을 외우고 있을 거라고 가정하면 안 됩니다.
     *    특히 dltTopicSuffix 는 바꾸는데 retryTopicSuffix 는 기본값에 맡기면
     *    "둘 중 하나만 커스터마이징됐다" 는 인상을 주어 오해를 부릅니다.
     *
     * ④ numPartitions = "3", replicationFactor = "1" 을 반드시 명시합니다.
     *    numPartitions 기본값은 -1 이고, 이는 "브로커의 num.partitions 를 따르라" 는 뜻입니다.
     *    실습 브로커는 3이라 우연히 맞지만, 운영에서 원본이 12 파티션이면
     *    retry 토픽만 3 파티션이 되어 재시도 처리량이 조용히 1/4 로 떨어집니다.
     *
     * ⑤ 실측값
     *      attempt=1 topic=orders          t=0ms
     *      attempt=2 topic=orders-retry-0  t=1006ms    (대기 1s)
     *      attempt=3 topic=orders-retry-1  t=3012ms    (대기 2s)
     *      attempt=4 topic=orders-retry-2  t=7021ms    (대기 4s)
     *    총 7.021초입니다. Step 07 의 블로킹 25초와 비교하면 짧지만,
     *    "본선 파티션이 멈춘 시간" 으로 따지면 25,046ms 대 0ms 입니다. 이게 요점입니다.
     */
    @Component
    @Profile("sol08-q2")
    public static class Q2CourseConvention {

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

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-sol-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                log.warn("★ 실패 {}", Fail.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event) {
            log.error("DLT 수신 {}", event.orderId());
        }
    }

    // ==================================================================
    // 정답 3. @DltHandler 로 실패 원인 복원
    // ==================================================================
    /*
     * 왜 이 답인가
     *
     * ① 상수 이름이 Step 07 과 다릅니다.
     *      Step 07 (DeadLetterPublishingRecoverer 직접 사용)
     *        KafkaHeaders.DLT_ORIGINAL_TOPIC  → "kafka_dlt-original-topic"
     *      Step 08 (@RetryableTopic 인프라)
     *        KafkaHeaders.ORIGINAL_TOPIC      → "kafka_original-topic"
     *    Step 07 의 코드를 복사해 오면 @Header 가 값을 못 찾아
     *    "Could not resolve method parameter" 로 DLT 리스너 자체가 실패합니다.
     *    그리고 DltStrategy 기본값이 FAIL_ON_ERROR 이므로 그 DLT 파티션이 막힙니다.
     *
     * ② ★ ORIGINAL_TOPIC 에는 "orders" 가 찍힙니다. "orders-retry-2" 가 아닙니다.
     *    헤더는 본선에서 첫 실패가 났을 때 한 번만 붙고, 이후 retry 단계에서는
     *    덮어쓰이지 않습니다(DeadLetterPublishingRecoverer 가 기존 original 헤더가
     *    있으면 그대로 둡니다). 그래서 DLT 에서 "이 메시지가 원래 어디서 왔는가" 를
     *    정확히 복원할 수 있습니다. 재발행할 때 이 값을 목적지로 쓰면 됩니다.
     *
     * ③ ORIGINAL_PARTITION 을 int, ORIGINAL_OFFSET 을 long 으로 선언하면
     *    Spring 이 4바이트/8바이트 빅엔디안 배열을 알아서 변환해 줍니다.
     *    ConsumerRecord.headers() 로 직접 순회할 때만 ByteBuffer 가 필요합니다(Step 07).
     *
     * ④ DLT 핸들러에는 재시도를 걸지 않습니다. 여기서 또 실패하면 갈 곳이 없습니다.
     *    알림을 보내다 실패해도 예외를 밖으로 던지지 말고 로그만 남기는 편이 안전합니다.
     */
    @Component
    @Profile("sol08-q3")
    public static class Q3DltHandler {

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

        @RetryableTopic(
                attempts = "3",
                backoff = @Backoff(delay = 500, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-sol-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                throw new RemoteApiException("재고 API 타임아웃: " + record.value().orderId());
            }
            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event,
                          @Header(KafkaHeaders.ORIGINAL_TOPIC)     String topic,
                          @Header(KafkaHeaders.ORIGINAL_PARTITION) int    partition,
                          @Header(KafkaHeaders.ORIGINAL_OFFSET)    long   offset,
                          @Header(KafkaHeaders.EXCEPTION_FQCN)     String exFqcn) {

            log.error("DLT 수신 {} origin={}-{}@{} ex={}",
                    event.orderId(), topic, partition, offset,
                    exFqcn.substring(exFqcn.lastIndexOf('.') + 1));
        }
    }

    // ==================================================================
    // 정답 4. ★핵심★ 순서 깨짐 재현
    // ==================================================================
    /*
     * 왜 이 답인가 — 그리고 이 문제의 결론
     *
     * 실측 타임라인
     *   | 시각      | 토픽           | 이벤트    | 결과                       | STATE     |
     *   |----------|---------------|----------|---------------------------|-----------|
     *   | t=2004ms | orders        | CREATED  | 실패 → retry-0 발행, 커밋   | (없음)     |
     *   | t=2011ms | orders        | UPDATED  | 성공                       | UPDATED   |
     *   | t=2014ms | orders        | CANCELLED| 성공                       | CANCELLED |
     *   | t=3018ms | orders-retry-0| CREATED  | ★성공★                    | CREATED   |
     *
     *   최종 상태 ORD-0002 = CREATED      ← 취소된 주문이 되살아났습니다
     *
     * 관찰 결과
     *   - ERROR 로그: 0줄. (WARN 한 줄만 있고 그건 우리가 직접 찍은 것입니다)
     *   - @DltHandler 호출: 안 됨. 재시도가 "성공" 했으니까요.
     *   - kcg --describe 의 LAG: 전부 0.
     *
     * ★ 이 문제의 역설
     *   이 데이터 오염을 만든 것은 "실패" 가 아니라 "성공" 입니다.
     *   재시도가 계속 실패해 DLT 로 갔다면 @DltHandler 가 알림을 보냈을 것이고,
     *   누군가 알아챘을 것입니다. 그런데 두 번째 시도에서 성공해 버리면
     *   시스템 어디에도 "무언가 잘못됐다" 는 신호가 남지 않습니다.
     *   모니터링이 잡을 수 없는 종류의 버그입니다.
     *
     * ★ 왜 블로킹에는 이 문제가 없는가
     *   DefaultErrorHandler 는 consumer.seek() 으로 "제자리에서" 반복합니다.
     *   offset 100 이 성공하거나 DLT 로 갈 때까지 101, 102 는 poll 조차 되지 않습니다.
     *   Step 07 의 "파티션이 멈춘다" 는 단점은, 정확히 이 순서 보장의 다른 이름입니다.
     *
     * ★ 해결 3가지 (본문 8-6)
     *   ① 순서가 중요하면 논블로킹을 쓰지 않는다. 블로킹 + 짧은 백오프 + 즉시 DLT.
     *
     *   ② 이벤트에 버전/시퀀스를 넣고 역행을 거부한다:
     *        Long current = VERSION.get(e.orderId());
     *        if (current != null && e.version() <= current) {
     *            log.warn("오래된 이벤트 무시 {} v{} <= v{}", e.orderId(), e.version(), current);
     *            return;
     *        }
     *        VERSION.put(e.orderId(), e.version());
     *      Step 13 의 멱등 컨슈머와 같은 구조입니다.
     *
     *   ③ 이벤트를 delta 가 아니라 snapshot 으로 설계한다.
     *      "생성/수정/취소" 대신 "현재 상태 전체 + updatedAt" 을 보내면
     *      늦게 온 오래된 스냅샷을 timestamp 비교만으로 버릴 수 있습니다.
     *
     * ⚠️ ②와 ③은 "오염을 막을" 뿐 "순서를 되돌리지는" 못합니다.
     *    재시도된 CREATED 는 그냥 버려집니다. 그것이 옳은 동작인지는 도메인이 정합니다.
     */
    @Component
    @Profile("sol08-q4")
    public static class Q4OrderState {

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

        static final Map<String, OrderState> STATE = new ConcurrentHashMap<>();
        static final AtomicBoolean FIRST_TRY = new AtomicBoolean(true);
        static final long T0 = System.currentTimeMillis();

        static long elapsed() {
            return System.currentTimeMillis() - T0;
        }

        @RetryableTopic(
                attempts = "4",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-sol-order-state")
        public void onMessage(ConsumerRecord<String, OrderEvent> record) {
            OrderEvent event = record.value();
            if (event.status() == OrderState.CREATED && FIRST_TRY.getAndSet(false)) {
                log.warn("★ 실패 {} {} (topic={} t={}ms)",
                        event.orderId(), event.status(), record.topic(), elapsed());
                throw new RemoteApiException("재고 API 타임아웃");
            }
            STATE.put(event.orderId(), event.status());
            log.info("상태 반영 {} → {}  (topic={} t={}ms)",
                    event.orderId(), event.status(), record.topic(), elapsed());
        }

        @DltHandler
        public void onDlt(OrderEvent event) {
            log.error("DLT 수신 {} {}", event.orderId(), event.status());
        }
    }

    @Component
    @Profile("sol08-q4")
    public static class Q4Publisher implements ApplicationRunner {

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

        private final KafkaTemplate<String, Object> template;

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

        @Override
        public void run(ApplicationArguments args) {
            Thread t = new Thread(() -> {
                sleep(2_000L);
                for (OrderState s : new OrderState[]{
                        OrderState.CREATED, OrderState.UPDATED, OrderState.CANCELLED}) {
                    OrderEvent e = OrderEvent.of("ORD-0002", s);
                    template.send("orders", e.orderId(), e);
                    log.info("발행 {} {} (t={}ms)", e.orderId(), s, Q4OrderState.elapsed());
                }
                sleep(3_000L);
                log.info("최종 상태 ORD-0002 = {}   ← ★ CANCELLED 가 아니면 순서가 깨진 것입니다",
                        Q4OrderState.STATE.get("ORD-0002"));
            }, "sol-q4-publisher");
            t.setDaemon(true);
            t.start();
        }

        private static void sleep(long ms) {
            try {
                Thread.sleep(ms);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt();
            }
        }
    }

    // ==================================================================
    // 정답 5. exclude + traversingCauses
    // ==================================================================
    /*
     * 왜 이 답인가
     *
     * ① exclude = { InvalidOrderException.class }
     *    include 를 비워 두면 "나열한 것 말고 전부 재시도" 가 됩니다(블랙리스트).
     *    반대로 include 만 주면 화이트리스트가 되어, 명시하지 않은 예외는
     *    전부 즉시 DLT 로 갑니다. 여기서는 RemoteApiException 을 재시도해야 하므로
     *    exclude 가 맞습니다. 새로운 예외가 추가돼도 기본이 "재시도" 라 안전합니다.
     *
     * ② ★ traversingCauses = "true" 가 이 문제의 핵심입니다.
     *    리스너가 던진 예외는 컨테이너에 도달할 때
     *      ListenerExecutionFailedException
     *        └─ caused by: InvalidOrderException
     *    으로 "감싸여" 옵니다. 기본값(false)에서는 최상위 예외 타입만 보고 분류하므로
     *    exclude 에 InvalidOrderException 을 적어도 매칭되지 않습니다.
     *    → 결과: 검증 실패인데 3번 재시도하고 retry 토픽 2개를 거칩니다.
     *      기능적으로는 결국 DLT 에 도착하므로 "동작은 한다" 고 착각하기 쉽습니다.
     *
     *    Step 07 의 addNotRetryableExceptions 도 정확히 같은 함정을 갖습니다.
     *    그쪽은 DefaultErrorHandler 가 내부적으로 cause 를 벗겨 주지만,
     *    RuntimeException 으로 한 번 더 감싸면 역시 분류가 깨집니다.
     *
     * ③ 확인 방법이 중요합니다. "DLT 에 도착했다" 만 보면 정답과 오답이 구분되지 않습니다.
     *    ★ orders-retry-0 이 비어 있는지를 봐야 합니다:
     *        kcc --topic orders-retry-0 --from-beginning --timeout-ms 5000
     *      → 아무것도 안 나오면 정답, ORD-0043 이 나오면 traversingCauses 를 빠뜨린 것입니다.
     *
     * ④ 로그에서도 구분됩니다.
     *      정답: "★ 검증 실패 ORD-0043 (topic=orders)" 가 딱 1번
     *      오답: 같은 줄이 topic=orders, orders-retry-0, orders-retry-1 로 3번
     */
    @Component
    @Profile("sol08-q5")
    public static class Q5Exclude {

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

        @RetryableTopic(
                attempts = "3",
                backoff = @Backoff(delay = 1000, multiplier = 2.0),
                topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE,
                retryTopicSuffix = "-retry",
                dltTopicSuffix = ".DLT",
                numPartitions = "3",
                replicationFactor = "1",
                exclude = { InvalidOrderException.class },   // ← 재시도 제외
                traversingCauses = "true",                   // ← ★ 원인 체인까지 따라간다
                kafkaTemplate = "kafkaTemplate")
        @KafkaListener(topics = "orders", groupId = "s08-sol-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            OrderCreated event = record.value();

            if (record.partition() == 1 && record.offset() == 42) {
                log.warn("★ 검증 실패 {} (topic={})", event.orderId(), record.topic());
                throw new InvalidOrderException("quantity 는 1 이상이어야 합니다: " + event.orderId());
            }

            if (record.partition() == 2 && record.offset() == 30) {
                log.warn("★ API 실패 {} (topic={})", event.orderId(), record.topic());
                throw new RemoteApiException("재고 API 타임아웃");
            }

            log.info("처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }

        @DltHandler
        public void onDlt(OrderCreated event,
                          @Header(KafkaHeaders.ORIGINAL_TOPIC) String topic,
                          @Header(KafkaHeaders.EXCEPTION_FQCN) String exFqcn) {
            log.error("DLT 수신 {} origin={} ex={}", event.orderId(), topic,
                    exFqcn.substring(exFqcn.lastIndexOf('.') + 1));
        }
    }

    // ==================================================================
    // 정답 6. RetryTopicConfiguration 빈으로 전역 설정
    // ==================================================================
    /*
     * 왜 이 답인가
     *
     * 애노테이션 속성과 빌더 메서드의 1:1 대조표
     *
     *   @RetryableTopic 속성                    | RetryTopicConfigurationBuilder
     *   ---------------------------------------|--------------------------------------------
     *   (리스너에 직접 부착)                     | .includeTopic("orders")
     *   attempts = "4"                          | .maxAttempts(4)
     *   @Backoff(delay=1000, multiplier=2.0,    | .exponentialBackoff(1000, 2.0, 10000)
     *            maxDelay=10000)                |
     *   @Backoff(delay=1000)                    | .fixedBackOff(1000)
     *   topicSuffixingStrategy                  | .suffixTopicsWithIndexValues()
     *     = SUFFIX_WITH_INDEX_VALUE             |
     *   retryTopicSuffix = "-retry"             | .retryTopicSuffix("-retry")
     *   dltTopicSuffix = ".DLT"                 | .dltSuffix(".DLT")
     *   concurrency = "1"                       | .concurrency(1)
     *   autoCreateTopics/numPartitions/         | .autoCreateTopics(true, 3, (short) 1)
     *     replicationFactor                     |
     *   exclude = { X.class }                   | .notRetryOn(X.class)
     *   include = { X.class }                   | .retryOn(X.class)
     *   traversingCauses = "true"               | .traversingCauses()
     *   dltStrategy = FAIL_ON_ERROR             | .doNotRetryOnDltFailure()
     *   dltStrategy = ALWAYS_RETRY_ON_ERROR     | .retryOnDltFailure()
     *   dltStrategy = NO_DLT                    | .doNotConfigureDlt()
     *
     * ① includeTopic("orders") 는 "이 토픽을 구독하는 모든 @KafkaListener" 에 적용됩니다.
     *    그래서 Q6Inventory 와 Q6Notification 둘 다에 재시도가 붙습니다.
     *    ⚠️ 이게 장점이자 함정입니다. 알림 리스너는 재시도가 필요 없는데도 붙어 버립니다.
     *      리스너마다 정책이 달라야 하면 애노테이션으로 돌아가거나
     *      excludeTopics 로 빼야 합니다.
     *
     * ② ★ 애노테이션과 빈이 둘 다 있으면 "애노테이션이 이깁니다."
     *    그래서 문제 지문이 "애노테이션을 전부 지운 상태에서 풀라" 고 한 것입니다.
     *    운영에서 이 규칙을 모르면 "전역 설정을 바꿨는데 왜 안 먹지?" 로 몇 시간을 씁니다.
     *    범인은 리스너 위에 남아 있는 옛 @RetryableTopic 한 줄입니다.
     *
     * ③ 컨슈머 그룹은 2개입니다(s08-sol-inventory, s08-sol-notification).
     *    retry 토픽과 DLT 는 각 리스너의 groupId 를 그대로 물려받으므로
     *    그룹이 늘어나지는 않습니다. kcg --list 로 확인하세요.
     *
     * ④ .create(template) 의 인자는 retry/DLT 발행에 쓸 KafkaTemplate 입니다.
     *    값 타입이 Object 인 템플릿을 넘겨야 여러 이벤트 타입을 한 설정으로 처리할 수 있습니다.
     */
    @Configuration
    @Profile("sol08-q6")
    public static class Q6BuilderConfig {

        @Bean
        public RetryTopicConfiguration ordersRetryTopicConfig(KafkaTemplate<String, Object> template) {
            return RetryTopicConfigurationBuilder
                    .newInstance()
                    .includeTopic("orders")
                    .maxAttempts(4)
                    .exponentialBackoff(1000, 2.0, 10_000)
                    .suffixTopicsWithIndexValues()
                    .retryTopicSuffix("-retry")
                    .dltSuffix(".DLT")
                    .concurrency(1)
                    .autoCreateTopics(true, 3, (short) 1)
                    .notRetryOn(InvalidOrderException.class)
                    .traversingCauses()
                    .doNotRetryOnDltFailure()
                    .create(template);
        }
    }

    @Component
    @Profile("sol08-q6")
    public static class Q6Inventory {

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

        @KafkaListener(topics = "orders", groupId = "s08-sol-inventory")
        public void onMessage(ConsumerRecord<String, OrderCreated> record) {
            if (Fail.shouldFail(record)) {
                log.warn("★ 실패 {}", Fail.stamp(record, record.value().orderId()));
                throw new RemoteApiException("재고 API 타임아웃");
            }
            log.info("재고 처리 완료 {}-{}@{}", record.topic(), record.partition(), record.offset());
        }
    }

    @Component
    @Profile("sol08-q6")
    public static class Q6Notification {

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

        @KafkaListener(topics = "orders", groupId = "s08-sol-notification")
        public void notifyCustomer(ConsumerRecord<String, OrderCreated> record) {
            log.info("알림 발송 완료 {} (topic={})", record.key(), record.topic());
        }
    }
}