학습 목표

시나리오

결제 서비스가 payments 토픽을 소비해 외부 정산 API로 전송합니다. 하루 90만 건, 파티션 3개, 컨슈머 인스턴스 3대입니다.

실패는 두 종류로 들어옵니다.

문제가 터진 지점은 둘을 구분하지 않았을 때였습니다. 파싱 불가 레코드 하나가 파티션 1의 맨 앞에 걸리자 그 파티션의 나머지 30만 건이 전부 밀렸습니다. lag 그래프는 파티션 1만 직선으로 치솟고 나머지 두 개는 정상이었습니다.

아키텍처

이 예제는 두 층의 재시도를 겹칩니다. 목적이 다릅니다.

블로킹 재시도와 논블로킹 재시도
구분 블로킹 재시도
(DefaultErrorHandler)
논블로킹 재시도
(@RetryableTopic)
동작 방식 같은 컨슈머 스레드에서 BackOff만큼 sleep 후 같은 레코드를 다시 처리 실패한 레코드를 retry 토픽으로 옮겨 쓰고 원본 오프셋은 즉시 전진
파티션 순서 유지 — 뒤 레코드는 대기 깨짐 — 실패 레코드가 뒤로 밀립니다
파티션 처리량 막힘 — 재시도 동안 뒤가 전부 정지 유지 — 뒤 레코드가 계속 처리됩니다
max.poll.interval.ms 총 백오프 시간이 이 값을 넘으면 그룹에서 축출됩니다 영향 없음
적합한 상황 수백 ms~수 초의 짧은 일시 장애. 순서가 중요한 토픽 수십 초~수 분이 필요한 외부 시스템 장애. 순서보다 처리량이 중요한 토픽

이 예제의 구성은 다음과 같습니다.

  1. 짧은 블로킹 재시도 3회 (1초 → 2초 → 4초). 순간적인 502를 흡수합니다.
  2. 그래도 실패하면 payments.DLT로 발행하고 원본 오프셋을 전진시킵니다.
  3. 파싱 실패 등 재시도 무의미한 예외는 재시도를 건너뛰고 즉시 DLQ로 갑니다.
  4. DLQ 컨슈머는 알림만 보내고 자동 재처리는 하지 않습니다(무한 루프 방지).

사전 요구사항

검증 환경 (2026-07 기준 안정 버전)
항목버전확인 방법
Apache Kafka (브로커) 4.3.1 예제 1의 클러스터
Spring Boot 4.1.0 Maven Central의 spring-boot-starter-parent 최신 정식 버전
Spring for Apache Kafka 4.1.0 Spring Boot 4.1.0이 관리하는 spring-kafka.version. 별도로 버전을 적지 않습니다
kafka-clients 4.3.1로 오버라이드 Spring Boot 4.1.0의 기본 관리 버전은 4.2.1입니다. 브로커와 맞추려고 <kafka.version>을 올립니다
Java 17 이상 Spring Framework 7 / Spring Boot 4는 Java 17을 기준선으로 합니다

전체 코드

디렉터리 구조
payment-dlq/
├── pom.xml
└── src/main/
    ├── java/com/example/payment/
    │   ├── PaymentDlqApplication.java   # 진입점
    │   ├── KafkaErrorHandlingConfig.java# 재시도 + DLQ 구성 (핵심)
    │   ├── PaymentListener.java         # 블로킹 재시도 + DLQ 리스너
    │   ├── SlowApiListener.java         # 논블로킹 재시도(@RetryableTopic) 리스너
    │   ├── DeadLetterListener.java      # DLQ 감시 리스너
    │   ├── SettlementApiClient.java     # 외부 API 클라이언트(장애 시뮬레이션)
    │   └── PaymentParseException.java   # 재시도 불가 예외
    └── resources/
        └── application.yml

pom.xml

payment-dlq/pom.xml
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
  <modelVersion>4.0.0</modelVersion>

  <parent>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-parent</artifactId>
    <!-- 2026-07 기준 최신 정식 버전. spring-kafka 4.1.0 을 관리합니다. -->
    <version>4.1.0</version>
    <relativePath/>
  </parent>

  <groupId>com.example</groupId>
  <artifactId>payment-dlq</artifactId>
  <version>1.0.0</version>

  <properties>
    <java.version>17</java.version>
    <!-- Spring Boot 4.1.0 의 기본 관리 버전은 4.2.1 입니다.
         브로커(4.3.1)와 클라이언트를 맞추려고 명시적으로 올립니다.
         kafka-clients 는 브로커보다 낮아도 동작하지만,
         버전을 맞춰 두면 새 기능·버그픽스 차이로 헤매는 일이 줄어듭니다. -->
    <kafka.version>4.3.1</kafka.version>
  </properties>

  <dependencies>
    <!-- Spring Boot 4 부터 제공되는 Kafka 스타터.
         3.x 의 org.springframework.kafka:spring-kafka 직접 의존을 대체합니다. -->
    <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-kafka</artifactId>
    </dependency>
    <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-json</artifactId>
    </dependency>
    <!-- /actuator/health, /actuator/metrics 로 컨슈머 상태를 노출합니다. -->
    <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-actuator</artifactId>
    </dependency>
    <dependency>
      <groupId>org.springframework.boot</groupId>
      <artifactId>spring-boot-starter-test</artifactId>
      <scope>test</scope>
    </dependency>
  </dependencies>

  <build>
    <plugins>
      <plugin>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-maven-plugin</artifactId>
      </plugin>
    </plugins>
  </build>
</project>

application.yml

src/main/resources/application.yml
spring:
  application:
    name: payment-dlq

  kafka:
    # 예제 1 클러스터의 호스트용 리스너. 초기 연결 이중화를 위해 3개를 적습니다.
    bootstrap-servers: localhost:29092,localhost:39092,localhost:49092

    consumer:
      group-id: payment-settlement
      # 새 그룹이 처음 붙을 때 처음부터 읽습니다. 기본값은 latest 입니다.
      auto-offset-reset: earliest
      # 자동 커밋을 끕니다. spring-kafka 의 컨테이너가 리스너 성공 후에 커밋합니다.
      # 자동 커밋이 켜져 있으면 처리 실패 레코드의 오프셋이 먼저 전진해 유실됩니다.
      enable-auto-commit: false
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      # 한 번의 poll 로 가져오는 최대 건수. 기본값 500.
      # 블로킹 재시도가 있으면 배치가 클수록 max.poll.interval.ms 초과 위험이 커집니다.
      max-poll-records: 50
      # poll 사이 최대 간격. 기본값 300000(5분).
      # 블로킹 재시도 총 백오프(이 예제는 7초) + 배치 처리 시간이 이 값 안에 들어와야 합니다.
      max-poll-interval: 300s
      properties:
        # spring.kafka.consumer 아래에 전용 키가 없는 설정은 properties 로 넘깁니다.
        # 상류가 트랜잭션 프로듀서를 쓴다면 read_committed 가 필요합니다(예제 5).
        isolation.level: read_committed

    producer:
      # DLQ 발행도 유실되면 사고 원인이 사라집니다. 내구성 설정을 그대로 적용합니다.
      acks: all
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.apache.kafka.common.serialization.StringSerializer
      properties:
        # spring.kafka.producer 에는 enable-idempotence 전용 키가 없습니다.
        # 4.3 기본값도 true 이지만 의도를 남깁니다.
        enable.idempotence: true
        # 4.0 에서 기본값이 0 -> 5 로 바뀐 설정입니다.
        linger.ms: 5
        delivery.timeout.ms: 120000

    listener:
      # RECORD: 리스너가 레코드 하나를 성공 처리할 때마다 오프셋을 커밋합니다.
      #         가장 안전하지만 커밋 요청이 많습니다.
      # BATCH(기본값): poll 배치 처리가 끝나면 커밋합니다.
      # 결제는 재처리 비용이 크므로 RECORD 를 씁니다.
      ack-mode: RECORD
      # 파티션이 3개이므로 인스턴스당 스레드 3개까지 의미가 있습니다.
      # 파티션 수보다 크게 잡으면 남는 스레드는 유휴 상태가 됩니다.
      concurrency: 3
      # 구독 대상 토픽이 없으면 기동에 실패하게 합니다(오타 조기 발견).
      missing-topics-fatal: true
      # poll 타임아웃. 기본값 5s.
      poll-timeout: 3s

app:
  topics:
    payments: payments
    payments-dlt: payments.DLT
    slow-api: payments.slow-api

management:
  endpoints:
    web:
      exposure:
        include: health,metrics,info

logging:
  level:
    # 재시도/복구 로그를 보기 위해 올립니다. 프로덕션에서는 INFO 로 둡니다.
    org.springframework.kafka.listener: INFO
    com.example.payment: DEBUG

재시도 + DLQ 구성 — 이 예제의 핵심

src/main/java/com/example/payment/KafkaErrorHandlingConfig.java
package com.example.payment;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.errors.SerializationException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.CommonErrorHandler;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.kafka.listener.RetryListener;
import org.springframework.kafka.support.ExponentialBackOffWithMaxRetries;
import org.springframework.kafka.support.converter.ConversionException;

/**
 * 재시도와 DLQ 를 한곳에서 구성합니다.
 *
 * Spring Boot 의 Kafka 자동 설정은 컨텍스트에 CommonErrorHandler 빈이 있으면
 * ConcurrentKafkaListenerContainerFactory 에 자동으로 적용합니다.
 * 따라서 컨테이너 팩토리를 직접 만들 필요가 없습니다.
 */
@Configuration
public class KafkaErrorHandlingConfig {

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

    /**
     * 실패한 레코드를 DLQ 토픽으로 발행합니다.
     *
     * 목적지 결정 규칙: 원본 토픽명 + ".DLT" 로 보내고
     * 원본과 같은 파티션 번호를 유지합니다.
     * 파티션을 유지하면 "어느 파티션에서 막혔는지" 를 DLQ 에서도 볼 수 있습니다.
     */
    @Bean
    public DeadLetterPublishingRecoverer deadLetterPublishingRecoverer(
            KafkaTemplate<Object, Object> template) {

        DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
                template,
                (ConsumerRecord<?, ?> record, Exception ex) -> {
                    // DLQ 파티션 수가 원본보다 적을 수 있으므로 -1 을 넘겨
                    // 파티셔너에 맡기는 것이 더 안전한 경우도 있습니다.
                    // 여기서는 원본과 DLQ 의 파티션 수가 같으므로(3) 번호를 유지합니다.
                    TopicPartition dlt =
                            new TopicPartition(record.topic() + ".DLT", record.partition());
                    log.error("DLQ 로 이동: {}-{}@{} key={} cause={}",
                            record.topic(), record.partition(), record.offset(),
                            record.key(), ex.getClass().getName());
                    return dlt;
                });

        // DLQ 발행이 실패했는데 성공한 것처럼 넘어가면 레코드가 사라집니다.
        // 발행 결과를 확인하고 실패 시 예외를 던지게 합니다.
        recoverer.setFailIfSendResultIsError(true);
        // 발행 결과를 기다리는 시간. delivery.timeout.ms(120초)보다 짧으면
        // 아직 재시도 중인데 실패로 판단할 수 있습니다.
        recoverer.setWaitForSendResultTimeout(java.time.Duration.ofSeconds(30));
        // 목적지를 결정하지 못하면(위 BiFunction 이 null 반환) 예외를 던집니다.
        recoverer.setThrowIfNoDestinationReturned(true);
        return recoverer;
    }

    /**
     * 블로킹 재시도 + DLQ.
     *
     * 백오프: 1초 → 2초 → 4초, 최대 3회 재시도(총 4회 처리 시도).
     * 총 대기 7초로, max.poll.interval.ms(300000) 대비 충분히 안전합니다.
     */
    @Bean
    public CommonErrorHandler paymentErrorHandler(DeadLetterPublishingRecoverer recoverer) {

        // ExponentialBackOff 는 기본적으로 "경과 시간" 으로 종료를 판단하는데,
        // spring-kafka 의 ExponentialBackOffWithMaxRetries 는 "횟수" 로 판단합니다.
        // 재시도 횟수를 명확히 통제하려면 이 클래스를 씁니다.
        ExponentialBackOffWithMaxRetries backOff = new ExponentialBackOffWithMaxRetries(3);
        backOff.setInitialInterval(1_000L);   // 1초
        backOff.setMultiplier(2.0);           // 1s → 2s → 4s
        backOff.setMaxInterval(10_000L);      // 상한 10초 (여기서는 도달하지 않습니다)

        DefaultErrorHandler handler = new DefaultErrorHandler(recoverer, backOff);

        // --- 재시도해도 결과가 같은 예외는 즉시 recoverer(DLQ) 로 보냅니다 ---------
        // 이 분류가 없으면 파싱 불가 레코드 하나가 7초 동안 파티션을 막습니다.
        handler.addNotRetryableExceptions(
                // 우리가 정의한 파싱 실패
                PaymentParseException.class,
                // 역직렬화 실패 — 바이트가 잘못된 것이므로 재시도 무의미
                SerializationException.class,
                ConversionException.class,
                // 코드 버그성 예외 — 재시도해도 같은 결과입니다
                NullPointerException.class,
                IllegalArgumentException.class,
                ClassCastException.class);

        // 복구(DLQ 발행)에 성공하면 그 레코드의 오프셋을 커밋합니다.
        // false 로 두면 오프셋이 전진하지 않아 재시작 시 같은 레코드를 다시 읽습니다.
        handler.setCommitRecovered(true);

        // 재시도/복구 이벤트를 관측 가능하게 만듭니다.
        // 이 로그가 없으면 "왜 lag 이 안 줄어드는지" 알 수 없습니다.
        handler.setRetryListeners(new LoggingRetryListener());

        return handler;
    }

    /** 재시도 횟수와 복구 결과를 로그로 남깁니다. */
    static class LoggingRetryListener implements RetryListener {

        @Override
        public void failedDelivery(ConsumerRecord<?, ?> record, Exception ex, int deliveryAttempt) {
            log.warn("재시도 {}회차 실패: {}-{}@{} key={} cause={}",
                    deliveryAttempt, record.topic(), record.partition(), record.offset(),
                    record.key(), ex.getClass().getSimpleName());
        }

        @Override
        public void recovered(ConsumerRecord<?, ?> record, Exception ex) {
            log.error("재시도 소진 → DLQ 발행 완료: {}-{}@{} key={}",
                    record.topic(), record.partition(), record.offset(), record.key());
        }

        @Override
        public void recoveryFailed(ConsumerRecord<?, ?> record, Exception original, Exception failure) {
            // DLQ 발행까지 실패한 최악의 경우입니다. 반드시 즉시 알림을 보내야 합니다.
            log.error("DLQ 발행 실패 — 레코드 유실 위험! {}-{}@{} original={} failure={}",
                    record.topic(), record.partition(), record.offset(),
                    original.getClass().getName(), failure.getClass().getName(), failure);
        }
    }
}

리스너와 예외

src/main/java/com/example/payment/PaymentParseException.java
package com.example.payment;

/**
 * 레코드 내용 자체가 처리 불가능함을 나타냅니다.
 * 재시도해도 결과가 같으므로 KafkaErrorHandlingConfig 에서
 * addNotRetryableExceptions 로 등록해 즉시 DLQ 로 보냅니다.
 */
public class PaymentParseException extends RuntimeException {

    public PaymentParseException(String message) {
        super(message);
    }

    public PaymentParseException(String message, Throwable cause) {
        super(message, cause);
    }
}
src/main/java/com/example/payment/SettlementApiClient.java
package com.example.payment;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;

import java.util.concurrent.atomic.AtomicInteger;

/**
 * 외부 정산 API 클라이언트 (실습용 시뮬레이션).
 *
 * 실제 구현은 RestClient / WebClient 를 쓰고,
 * HTTP 5xx·타임아웃은 재시도 가능 예외로,
 * 4xx(400/422)는 재시도 불가 예외로 매핑합니다.
 */
@Component
public class SettlementApiClient {

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

    // "amount" 가 999 로 끝나는 결제는 앞의 2회를 실패시키고 3회차에 성공시킵니다.
    // → 블로킹 재시도가 실제로 동작하는지 확인하는 장치입니다.
    private final AtomicInteger transientFailures = new AtomicInteger();

    public void settle(String paymentId, long amount) {
        if (amount % 1000 == 999) {
            int n = transientFailures.incrementAndGet();
            if (n % 3 != 0) {
                log.debug("정산 API 일시 실패 시뮬레이션 (n={}) paymentId={}", n, paymentId);
                // 재시도하면 성공할 수 있는 실패입니다.
                // RuntimeException 계열이면 DefaultErrorHandler 가 재시도합니다.
                throw new IllegalStateException(
                        "정산 API 502 Bad Gateway (일시적) paymentId=" + paymentId);
            }
        }
        log.debug("정산 완료 paymentId={} amount={}", paymentId, amount);
    }
}
src/main/java/com/example/payment/PaymentListener.java
package com.example.payment;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;

import java.util.regex.Matcher;
import java.util.regex.Pattern;

/**
 * 결제 이벤트를 소비해 정산 API 로 보냅니다.
 *
 * 실패 처리는 이 클래스에 없습니다 —
 * KafkaErrorHandlingConfig 의 DefaultErrorHandler 가 담당합니다.
 * 리스너는 "실패하면 예외를 던진다" 만 지키면 됩니다.
 * try-catch 로 예외를 삼키면 에러 핸들러가 동작하지 못하고 유실이 됩니다.
 */
@Component
public class PaymentListener {

    private static final Logger log = LoggerFactory.getLogger(PaymentListener.class);
    private static final Pattern AMOUNT = Pattern.compile("\"amount\"\\s*:\\s*(\\d+)");
    private static final Pattern PAYMENT_ID = Pattern.compile("\"paymentId\"\\s*:\\s*\"([^\"]+)\"");

    private final SettlementApiClient settlementApi;

    public PaymentListener(SettlementApiClient settlementApi) {
        this.settlementApi = settlementApi;
    }

    @KafkaListener(
            topics = "${app.topics.payments}",
            groupId = "payment-settlement",
            // 컨테이너 팩토리를 지정하지 않으면 Boot 가 자동 설정한 것을 씁니다.
            // 그 팩토리에는 위에서 정의한 CommonErrorHandler 가 적용되어 있습니다.
            concurrency = "3")
    public void onPayment(ConsumerRecord<String, String> record,
                          @Header(name = KafkaHeaders.DELIVERY_ATTEMPT, required = false)
                          Integer deliveryAttempt) {

        // DELIVERY_ATTEMPT 헤더는 컨테이너의 deliveryAttemptHeader 가 켜져 있을 때만 옵니다.
        // DefaultErrorHandler 는 이 헤더를 지원합니다(deliveryAttemptHeader() == true).
        log.debug("수신 {}-{}@{} key={} attempt={}",
                record.topic(), record.partition(), record.offset(), record.key(),
                deliveryAttempt);

        String paymentId = extract(PAYMENT_ID, record.value(), "paymentId");
        long amount = Long.parseLong(extract(AMOUNT, record.value(), "amount"));

        // 실패하면 예외가 밖으로 나가야 합니다. 절대 삼키지 않습니다.
        settlementApi.settle(paymentId, amount);
    }

    /**
     * 필드 추출. 실패는 재시도 불가 예외로 변환합니다.
     * 이 예외가 addNotRetryableExceptions 에 등록되어 있어 즉시 DLQ 로 갑니다.
     */
    private String extract(Pattern pattern, String json, String fieldName) {
        if (json == null) {
            throw new PaymentParseException("value 가 null 입니다");
        }
        Matcher m = pattern.matcher(json);
        if (!m.find()) {
            throw new PaymentParseException(
                    "필수 필드 없음: " + fieldName + " payload=" + truncate(json));
        }
        return m.group(1);
    }

    private static String truncate(String s) {
        return s.length() <= 200 ? s : s.substring(0, 200) + "...";
    }
}

논블로킹 재시도 — 오래 걸리는 외부 장애용

외부 시스템이 수 분 동안 죽어 있는 경우에는 블로킹 재시도를 쓸 수 없습니다. @RetryableTopic은 실패한 레코드를 {topic}-retry-0, -retry-1 토픽으로 옮겨 쓰고 원본 오프셋을 즉시 전진시킵니다. 파티션이 막히지 않습니다.

src/main/java/com/example/payment/SlowApiListener.java
package com.example.payment;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.DltHandler;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.annotation.RetryableTopic;
import org.springframework.kafka.retrytopic.TopicSuffixingStrategy;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.messaging.handler.annotation.Header;
import org.springframework.stereotype.Component;

/**
 * 논블로킹 재시도 예시.
 *
 * 실패한 레코드는 지연 토픽으로 이동하고 원본 파티션은 계속 진행합니다.
 * 대가는 "순서가 깨진다" 는 것입니다 —
 * 순서가 필요한 토픽에는 쓰지 마세요.
 *
 * 자동으로 만들어지는 토픽:
 *   payments.slow-api-retry-1000
 *   payments.slow-api-retry-5000
 *   payments.slow-api-retry-25000
 *   payments.slow-api-dlt
 */
@Component
public class SlowApiListener {

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

    @RetryableTopic(
            // 원본 1회 + 재시도 3회 = 총 4회 처리 시도
            attempts = "4",
            // 지연 1초 → 5초 → 25초 (multiplier 5)
            backoff = @org.springframework.retry.annotation.Backoff(
                    delay = 1_000L, multiplier = 5.0, maxDelay = 60_000L),
            // 지연 값을 토픽 이름 접미사로 씁니다(운영 중 어느 단계인지 바로 보입니다).
            topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_DELAY_VALUE,
            // 재시도 토픽·DLT 토픽을 이 값으로 만듭니다.
            // 원본 토픽과 파티션 수를 맞추는 것이 안전합니다.
            numPartitions = "3",
            replicationFactor = "3",
            // 재시도해도 같은 예외는 논블로킹 재시도도 건너뛰고 바로 DLT 로 보냅니다.
            exclude = { PaymentParseException.class, IllegalArgumentException.class })
    @KafkaListener(topics = "${app.topics.slow-api}", groupId = "slow-api-consumer")
    public void onSlowApiEvent(ConsumerRecord<String, String> record,
                               @Header(KafkaHeaders.RECEIVED_TOPIC) String receivedTopic) {
        // receivedTopic 을 보면 원본인지 재시도 토픽인지 알 수 있습니다.
        log.info("처리 시도 topic={} key={} offset={}", receivedTopic, record.key(), record.offset());

        // 실제 로직. 실패하면 예외를 던집니다.
        if (record.value().contains("\"forceFail\":true")) {
            throw new IllegalStateException("외부 시스템 장애 (시뮬레이션) key=" + record.key());
        }
    }

    /**
     * 모든 재시도를 소진한 레코드가 여기로 옵니다.
     * @RetryableTopic 이 만든 DLT 토픽을 자동으로 구독합니다.
     */
    @DltHandler
    public void onDlt(ConsumerRecord<String, String> record,
                      @Header(KafkaHeaders.RECEIVED_TOPIC) String receivedTopic,
                      @Header(name = KafkaHeaders.EXCEPTION_MESSAGE, required = false) String errorMessage) {
        log.error("DLT 도착 topic={} key={} offset={} error={}",
                receivedTopic, record.key(), record.offset(), errorMessage);
        // 여기서 알림을 보내고 자동 재처리는 하지 않습니다.
    }
}

DLQ 감시 — 자동 재처리를 하지 않는 이유

src/main/java/com/example/payment/DeadLetterListener.java
package com.example.payment;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.common.header.Header;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.concurrent.atomic.AtomicLong;

/**
 * DLQ 감시 리스너.
 *
 * 여기서 자동 재처리를 하면 안 됩니다.
 * 원인이 그대로인 상태에서 원본 토픽으로 되돌리면
 * DLQ ↔ 원본 사이를 무한히 도는 루프가 만들어집니다.
 * DLQ 는 "사람이 원인을 고친 뒤 의도적으로 재처리하는 곳" 입니다.
 */
@Component
public class DeadLetterListener {

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

    private final AtomicLong dlqCount = new AtomicLong();

    @KafkaListener(
            topics = "${app.topics.payments-dlt}",
            groupId = "payment-dlq-monitor",
            // DLQ 는 건수가 적으므로 단일 스레드로 충분합니다.
            concurrency = "1")
    public void onDeadLetter(ConsumerRecord<String, String> record) {
        long total = dlqCount.incrementAndGet();

        // DeadLetterPublishingRecoverer 가 붙여 준 헤더로 원인을 추적합니다.
        // 헤더 이름은 KafkaHeaders 상수로 정의되어 있습니다.
        String originalTopic = header(record, KafkaHeaders.DLT_ORIGINAL_TOPIC);
        String originalPartition = intHeader(record, KafkaHeaders.DLT_ORIGINAL_PARTITION);
        String originalOffset = longHeader(record, KafkaHeaders.DLT_ORIGINAL_OFFSET);
        String exceptionFqcn = header(record, KafkaHeaders.DLT_EXCEPTION_FQCN);
        String exceptionMessage = header(record, KafkaHeaders.DLT_EXCEPTION_MESSAGE);

        log.error("""
                        DLQ 레코드 #{}
                          원본        : {}-{}@{}
                          키          : {}
                          예외        : {}
                          메시지      : {}
                          payload     : {}""",
                total, originalTopic, originalPartition, originalOffset,
                record.key(), exceptionFqcn, exceptionMessage, record.value());

        // 실제 서비스에서는 여기서 알림(Slack/PagerDuty)을 보내고,
        // 원인별 건수를 메트릭으로 올립니다.
        // 자동 재처리는 하지 않습니다.
    }

    private static String header(ConsumerRecord<?, ?> record, String name) {
        Header h = record.headers().lastHeader(name);
        return h == null ? "(없음)" : new String(h.value(), StandardCharsets.UTF_8);
    }

    /** DLT_ORIGINAL_PARTITION 은 4바이트 big-endian int 로 저장됩니다. */
    private static String intHeader(ConsumerRecord<?, ?> record, String name) {
        Header h = record.headers().lastHeader(name);
        if (h == null || h.value().length < Integer.BYTES) {
            return "(없음)";
        }
        return String.valueOf(java.nio.ByteBuffer.wrap(h.value()).getInt());
    }

    /** DLT_ORIGINAL_OFFSET 은 8바이트 big-endian long 으로 저장됩니다. */
    private static String longHeader(ConsumerRecord<?, ?> record, String name) {
        Header h = record.headers().lastHeader(name);
        if (h == null || h.value().length < Long.BYTES) {
            return "(없음)";
        }
        return String.valueOf(java.nio.ByteBuffer.wrap(h.value()).getLong());
    }
}
src/main/java/com/example/payment/PaymentDlqApplication.java
package com.example.payment;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class PaymentDlqApplication {

    public static void main(String[] args) {
        SpringApplication.run(PaymentDlqApplication.class, args);
    }
}

DLQ 헤더 목록

DeadLetterPublishingRecoverer가 붙이는 헤더입니다. 이 헤더가 있어야 "무엇이 왜 실패했는지"를 나중에 알 수 있습니다.

spring-kafka 4.1.0 KafkaHeaders 상수와 실제 헤더 이름
상수 헤더 이름 내용
DLT_ORIGINAL_TOPICkafka_dlt-original-topic원본 토픽명 (문자열)
DLT_ORIGINAL_PARTITIONkafka_dlt-original-partition원본 파티션 (4바이트 int)
DLT_ORIGINAL_OFFSETkafka_dlt-original-offset원본 오프셋 (8바이트 long)
DLT_ORIGINAL_TIMESTAMPkafka_dlt-original-timestamp원본 타임스탬프 (8바이트 long)
DLT_ORIGINAL_TIMESTAMP_TYPEkafka_dlt-original-timestamp-typeCreateTime 또는 LogAppendTime
DLT_ORIGINAL_CONSUMER_GROUPkafka_dlt-original-consumer-group실패한 컨슈머 그룹 ID
DLT_EXCEPTION_FQCNkafka_dlt-exception-fqcn예외 클래스 전체 이름
DLT_EXCEPTION_CAUSE_FQCNkafka_dlt-exception-cause-fqcn근본 원인 예외 클래스
DLT_EXCEPTION_MESSAGEkafka_dlt-exception-message예외 메시지
DLT_EXCEPTION_STACKTRACEkafka_dlt-exception-stacktrace스택트레이스 (크기가 큽니다)
DLT_KEY_EXCEPTION_FQCNkafka_dlt-key-exception-fqcn키 역직렬화 실패 시의 예외

실행 방법

순서대로 실행
# 0. 예제 1의 클러스터. payments / payments.DLT 는 create-topics.sh 가 만듭니다.
cd kafka-lab && docker compose ps

# 1. @RetryableTopic 예제용 원본 토픽을 추가로 만듭니다.
#    (retry / dlt 토픽은 spring-kafka 가 자동 생성합니다)
./kcli kafka-topics.sh --create --if-not-exists \
  --topic payments.slow-api --partitions 3 --replication-factor 3 \
  --config min.insync.replicas=2

# 2. 애플리케이션 기동
cd ../payment-dlq
mvn -q spring-boot:run

# 3. 다른 터미널에서 정상 결제 3건
cd kafka-lab
for i in 1 2 3; do
  echo "PAY-$i:{\"paymentId\":\"PAY-$i\",\"amount\":$((10000 + i))}"
done | ./kcli kafka-console-producer.sh --topic payments \
        --property parse.key=true --property key.separator=:

# 4. 일시 실패 → 재시도 후 성공하는 건 (amount 가 999 로 끝남)
echo 'PAY-T1:{"paymentId":"PAY-T1","amount":10999}' \
  | ./kcli kafka-console-producer.sh --topic payments \
      --property parse.key=true --property key.separator=:

# 5. 파싱 실패 → 재시도 없이 즉시 DLQ 로 가는 건 (amount 필드 없음)
echo 'PAY-BAD:{"paymentId":"PAY-BAD"}' \
  | ./kcli kafka-console-producer.sh --topic payments \
      --property parse.key=true --property key.separator=:

# 6. 논블로킹 재시도 확인 (retry 토픽으로 이동)
echo 'SLOW-1:{"id":"SLOW-1","forceFail":true}' \
  | ./kcli kafka-console-producer.sh --topic payments.slow-api \
      --property parse.key=true --property key.separator=:

검증 방법

1. 일시 실패가 재시도로 회복되는가

4단계 실행 후 애플리케이션 로그
DEBUG c.e.payment.PaymentListener  : 수신 payments-1@4 key=PAY-T1 attempt=1
DEBUG c.e.p.SettlementApiClient    : 정산 API 일시 실패 시뮬레이션 (n=1) paymentId=PAY-T1
 WARN c.e.p.KafkaErrorHandlingConfig: 재시도 1회차 실패: payments-1@4 key=PAY-T1 cause=IllegalStateException
DEBUG c.e.payment.PaymentListener  : 수신 payments-1@4 key=PAY-T1 attempt=2
 WARN c.e.p.KafkaErrorHandlingConfig: 재시도 2회차 실패: payments-1@4 key=PAY-T1 cause=IllegalStateException
DEBUG c.e.payment.PaymentListener  : 수신 payments-1@4 key=PAY-T1 attempt=3
DEBUG c.e.p.SettlementApiClient    : 정산 완료 paymentId=PAY-T1 amount=10999

확인 포인트는 두 가지입니다. attempt가 1→2→3으로 증가하고, 로그 사이 간격이 1초 → 2초로 벌어지는 것입니다. 마지막에 성공하면 DLQ로 가지 않습니다.

2. 재시도 불가 예외가 즉시 DLQ로 가는가

5단계 실행 후 애플리케이션 로그
 WARN c.e.p.KafkaErrorHandlingConfig: 재시도 1회차 실패: payments-2@7 key=PAY-BAD cause=PaymentParseException
ERROR c.e.p.KafkaErrorHandlingConfig: DLQ 로 이동: payments-2@7 key=PAY-BAD cause=com.example.payment.PaymentParseException
ERROR c.e.p.KafkaErrorHandlingConfig: 재시도 소진 → DLQ 발행 완료: payments-2@7 key=PAY-BAD
ERROR c.e.payment.DeadLetterListener : DLQ 레코드 #1
  원본        : payments-2@7
  키          : PAY-BAD
  예외        : com.example.payment.PaymentParseException
  메시지      : 필수 필드 없음: amount payload={"paymentId":"PAY-BAD"}
  payload     : {"paymentId":"PAY-BAD"}

핵심은 "재시도 1회차 실패" 뒤 곧바로 DLQ로 갔다는 것입니다. addNotRetryableExceptions가 없었다면 여기서 1초·2초·4초를 기다린 뒤 DLQ로 갔을 것이고, 그 7초 동안 파티션 2의 뒷 레코드는 모두 정지했을 것입니다.

DLQ 토픽에서 헤더까지 직접 확인
./kcli kafka-console-consumer.sh --topic payments.DLT --from-beginning \
  --timeout-ms 8000 \
  --property print.key=true --property print.headers=true 2>/dev/null
기대 출력 (헤더가 한 줄에 이어져 나옵니다)
kafka_dlt-original-topic:payments,kafka_dlt-original-partition:...,kafka_dlt-original-offset:...,kafka_dlt-original-consumer-group:payment-settlement,kafka_dlt-exception-fqcn:com.example.payment.PaymentParseException,kafka_dlt-exception-message:필수 필드 없음: amount payload={"paymentId":"PAY-BAD"}	PAY-BAD	{"paymentId":"PAY-BAD"}

3. 논블로킹 재시도가 파티션을 막지 않는가

자동 생성된 재시도 토픽 확인
./kcli kafka-topics.sh --list | grep slow-api
기대 출력 — 지연 값이 토픽 이름에 붙습니다
payments.slow-api
payments.slow-api-dlt
payments.slow-api-retry-1000
payments.slow-api-retry-25000
payments.slow-api-retry-5000
원본 토픽의 lag이 0으로 유지되는지 확인
# 실패 레코드가 있어도 원본 토픽의 lag 은 0 이어야 합니다.
# 레코드가 retry 토픽으로 "옮겨졌기" 때문입니다.
./kcli kafka-consumer-groups.sh --describe --group slow-api-consumer

4. 블로킹 재시도가 파티션을 막는 것을 직접 관찰

이것이 이 예제에서 가장 배울 것이 많은 실험입니다. addNotRetryableExceptions일부러 주석 처리하고 백오프를 크게 늘린 뒤, 파싱 실패 레코드 하나를 넣고 그 뒤에 정상 레코드 10건을 넣으세요.

파티션별 lag 관찰
watch -n 2 "docker exec kafka-1 /opt/kafka/bin/kafka-consumer-groups.sh \
  --bootstrap-server kafka-1:19092 --describe --group payment-settlement"

실패 레코드가 있는 파티션만 lag이 쌓이고 나머지는 0인 모양이 보입니다. 시나리오에서 설명한 "파티션 1만 직선으로 치솟는 lag 그래프"가 정확히 이 모양입니다. 이 패턴을 알아보는 것만으로 장애 조사 시간이 크게 줄어듭니다 (케이스 2).

DLQ 재처리 절차

DLQ에 쌓인 레코드를 되돌리는 것은 코드가 아니라 절차의 문제입니다. 원인을 고치기 전에 되돌리면 무한 루프가 됩니다.

  1. 원인 분류kafka_dlt-exception-fqcn 헤더로 그룹화합니다. 원인이 여러 가지면 각각 다르게 처리해야 합니다.
  2. 수정 배포 — 코드 버그면 컨슈머를 고쳐 배포합니다. 데이터 문제면 상류 프로듀서를 고칩니다.
  3. 소량 재처리로 검증 — DLQ에서 몇 건만 원본 토픽으로 되돌려 성공하는지 확인합니다.
  4. 전량 재처리 — 확인 후 나머지를 되돌립니다.

되돌리는 방법은 두 가지입니다.

(A) MirrorMaker 2로 DLQ → 원본 복사 (권장)
# dlq-replay.properties
# 같은 클러스터 안에서 토픽 이름만 바꿔 복사합니다.
# IdentityReplicationPolicy 를 쓰면 원본 토픽명에 접두어가 붙지 않습니다.
clusters = local
local.bootstrap.servers = kafka-1:19092,kafka-2:19092,kafka-3:19092

# 같은 클러스터로 되돌리므로 별칭을 하나만 쓰는 대신
# 실제 운영에서는 replay 전용 컨슈머 그룹으로 처리하는 편이 단순합니다.
# 자세한 MirrorMaker 2 구성은 예제 11을 보세요.
(B) 전용 재처리 컨슈머 (실무에서 가장 많이 쓰는 방식)
package com.example.payment.replay;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.KafkaHeaders;
import org.springframework.stereotype.Component;

import java.nio.charset.StandardCharsets;
import java.util.Set;

/**
 * DLQ 재처리 전용 리스너.
 *
 * 평상시에는 꺼져 있고(app.replay.enabled=false),
 * 원인을 고친 뒤 의도적으로 켜서 실행합니다.
 * 이렇게 스위치를 두면 "실수로 무한 루프" 를 구조적으로 막을 수 있습니다.
 */
@Component
@ConditionalOnProperty(name = "app.replay.enabled", havingValue = "true")
public class DlqReplayListener {

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

    /** 이 예외로 실패한 것만 되돌립니다. 원인별 선별 재처리가 안전합니다. */
    private static final Set<String> REPLAYABLE = Set.of(
            "java.lang.IllegalStateException",
            "org.springframework.web.client.HttpServerErrorException");

    private final KafkaTemplate<String, String> template;

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

    @KafkaListener(
            topics = "${app.topics.payments-dlt}",
            // 재처리는 감시 그룹과 다른 그룹을 씁니다.
            // 같은 그룹을 쓰면 감시 리스너의 오프셋을 건드립니다.
            groupId = "payment-dlq-replay",
            concurrency = "1")
    public void replay(ConsumerRecord<String, String> record) {
        String fqcn = headerAsString(record, KafkaHeaders.DLT_EXCEPTION_FQCN);
        String originalTopic = headerAsString(record, KafkaHeaders.DLT_ORIGINAL_TOPIC);

        if (!REPLAYABLE.contains(fqcn)) {
            log.info("재처리 대상 아님 — 건너뜁니다. key={} cause={}", record.key(), fqcn);
            return;
        }

        // 원본 토픽으로 되돌립니다. 키를 유지해 원래 파티션으로 갑니다.
        // sendResult 를 확인해 실패 시 예외가 나게 합니다(오프셋이 전진하지 않습니다).
        template.send(originalTopic, record.key(), record.value())
                .join();
        log.info("재처리 발행 완료 key={} → {}", record.key(), originalTopic);
    }

    private static String headerAsString(ConsumerRecord<?, ?> record, String name) {
        var h = record.headers().lastHeader(name);
        return h == null ? "" : new String(h.value(), StandardCharsets.UTF_8);
    }
}

프로덕션 고려사항

로컬 예제와 프로덕션의 차이
항목이 예제프로덕션
DLQ 알림 로그만 남깁니다 DLQ 도착률을 메트릭으로 올리고 1건이라도 오면 알림을 보냅니다. 로그만 남기면 아무도 안 봅니다
DLQ 보관 기간 30일 조사·재처리에 필요한 기간 + 여유. 기본값 7일은 연휴를 넘기지 못합니다. 용량 산정도 함께 해야 합니다
예외 분류 클래스 6개 HTTP 상태 코드까지 매핑합니다. 5xx·타임아웃은 재시도, 400·422는 즉시 DLQ. 분류하지 않은 예외는 기본적으로 재시도된다는 점을 기억하세요
재시도 전략 블로킹 3회 (총 7초) 짧은 장애는 블로킹, 긴 장애는 논블로킹. 블로킹 총 백오프 < max.poll.interval.ms를 반드시 지켜야 합니다
순서 보장 블로킹만 고려 논블로킹 재시도는 순서를 깨뜨립니다. 상태 전이 이벤트(주문 생성 → 결제 → 배송)에는 쓰면 안 됩니다
멱등성 없음 재시도는 같은 처리를 여러 번 실행합니다. 외부 API에 멱등 키를 보내거나 처리 기록 테이블을 두어야 합니다(예제 5의 EOS 경계 논의)
역직렬화 실패 StringDeserializer라 거의 없음 Avro/JSON 역직렬화 실패는 리스너에 도달하기도 전에 터집니다. ErrorHandlingDeserializer로 감싸야 DLQ로 보낼 수 있습니다(예제 7)
DLQ 발행 실패 로그 recoveryFailed레코드 유실 직전 상태입니다. 최고 심각도 알림 + 로컬 디스크 백업을 함께 두세요

자주 하는 실수

이어서 볼 곳

공식 문서 출처

Kafka 설정 기본값은 Apache Kafka 4.3.1 문서에서, spring-kafka API는 Maven Central의 spring-kafka-4.1.0 아티팩트에서 직접 확인했습니다.