실무 예제 · 6
DLQ + 재시도 패턴
컨슈머가 레코드 하나를 처리하다 실패하면 세 가지 선택지밖에 없습니다. 재시도하거나, 건너뛰거나, 어딘가에 치워 두거나. 아무 조치도 하지 않으면 spring-kafka는 기본적으로 재시도를 반복하고, 그 사이 그 파티션의 뒤에 있는 모든 레코드가 멈춥니다. 이 예제는 재시도 가능/불가능을 분류하고, 파티션을 막지 않는 재시도를 구성하고, 마지막에 실패 원인을 함께 담은 DLQ로 보내는 전체 구조를 만듭니다.
학습 목표
- 블로킹 재시도(
DefaultErrorHandler)와 논블로킹 재시도(@RetryableTopic)의 차이와 선택 기준을 설명할 수 있습니다. - 재시도해도 결과가 같은 예외를 즉시 DLQ로 보내는 분류를 구성할 수 있습니다.
DeadLetterPublishingRecoverer가 DLQ 레코드에 붙이는 헤더로 원인 추적을 할 수 있습니다.- DLQ에 쌓인 레코드를 안전하게 재처리하는 절차를 설계할 수 있습니다.
시나리오
결제 서비스가 payments 토픽을 소비해 외부 정산 API로 전송합니다.
하루 90만 건, 파티션 3개, 컨슈머 인스턴스 3대입니다.
실패는 두 종류로 들어옵니다.
- 일시적 실패 — 정산 API가 5초간 502를 반환합니다. 잠시 뒤 재시도하면 성공합니다.
- 영구적 실패 — 상류가 잘못된 스키마로 발행한 레코드가 섞여 있습니다. 몇 번 재시도해도 항상 파싱에 실패합니다.
문제가 터진 지점은 둘을 구분하지 않았을 때였습니다. 파싱 불가 레코드 하나가 파티션 1의 맨 앞에 걸리자 그 파티션의 나머지 30만 건이 전부 밀렸습니다. lag 그래프는 파티션 1만 직선으로 치솟고 나머지 두 개는 정상이었습니다.
아키텍처
이 예제는 두 층의 재시도를 겹칩니다. 목적이 다릅니다.
| 구분 | 블로킹 재시도 ( DefaultErrorHandler) |
논블로킹 재시도 ( @RetryableTopic) |
|---|---|---|
| 동작 방식 | 같은 컨슈머 스레드에서 BackOff만큼 sleep 후 같은 레코드를 다시 처리 |
실패한 레코드를 retry 토픽으로 옮겨 쓰고 원본 오프셋은 즉시 전진 |
| 파티션 순서 | 유지 — 뒤 레코드는 대기 | 깨짐 — 실패 레코드가 뒤로 밀립니다 |
| 파티션 처리량 | 막힘 — 재시도 동안 뒤가 전부 정지 | 유지 — 뒤 레코드가 계속 처리됩니다 |
max.poll.interval.ms |
총 백오프 시간이 이 값을 넘으면 그룹에서 축출됩니다 | 영향 없음 |
| 적합한 상황 | 수백 ms~수 초의 짧은 일시 장애. 순서가 중요한 토픽 | 수십 초~수 분이 필요한 외부 시스템 장애. 순서보다 처리량이 중요한 토픽 |
이 예제의 구성은 다음과 같습니다.
- 짧은 블로킹 재시도 3회 (1초 → 2초 → 4초). 순간적인 502를 흡수합니다.
- 그래도 실패하면
payments.DLT로 발행하고 원본 오프셋을 전진시킵니다. - 파싱 실패 등 재시도 무의미한 예외는 재시도를 건너뛰고 즉시 DLQ로 갑니다.
- DLQ 컨슈머는 알림만 보내고 자동 재처리는 하지 않습니다(무한 루프 방지).
사전 요구사항
| 항목 | 버전 | 확인 방법 |
|---|---|---|
| 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
<?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
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 구성 — 이 예제의 핵심
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);
}
}
}
리스너와 예외
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);
}
}
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);
}
}
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 … 토픽으로 옮겨 쓰고
원본 오프셋을 즉시 전진시킵니다. 파티션이 막히지 않습니다.
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 감시 — 자동 재처리를 하지 않는 이유
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());
}
}
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가 붙이는 헤더입니다.
이 헤더가 있어야 "무엇이 왜 실패했는지"를 나중에 알 수 있습니다.
| 상수 | 헤더 이름 | 내용 |
|---|---|---|
DLT_ORIGINAL_TOPIC | kafka_dlt-original-topic | 원본 토픽명 (문자열) |
DLT_ORIGINAL_PARTITION | kafka_dlt-original-partition | 원본 파티션 (4바이트 int) |
DLT_ORIGINAL_OFFSET | kafka_dlt-original-offset | 원본 오프셋 (8바이트 long) |
DLT_ORIGINAL_TIMESTAMP | kafka_dlt-original-timestamp | 원본 타임스탬프 (8바이트 long) |
DLT_ORIGINAL_TIMESTAMP_TYPE | kafka_dlt-original-timestamp-type | CreateTime 또는 LogAppendTime |
DLT_ORIGINAL_CONSUMER_GROUP | kafka_dlt-original-consumer-group | 실패한 컨슈머 그룹 ID |
DLT_EXCEPTION_FQCN | kafka_dlt-exception-fqcn | 예외 클래스 전체 이름 |
DLT_EXCEPTION_CAUSE_FQCN | kafka_dlt-exception-cause-fqcn | 근본 원인 예외 클래스 |
DLT_EXCEPTION_MESSAGE | kafka_dlt-exception-message | 예외 메시지 |
DLT_EXCEPTION_STACKTRACE | kafka_dlt-exception-stacktrace | 스택트레이스 (크기가 큽니다) |
DLT_KEY_EXCEPTION_FQCN | kafka_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. 일시 실패가 재시도로 회복되는가
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로 가는가
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의 뒷 레코드는 모두 정지했을 것입니다.
./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 이어야 합니다.
# 레코드가 retry 토픽으로 "옮겨졌기" 때문입니다.
./kcli kafka-consumer-groups.sh --describe --group slow-api-consumer
4. 블로킹 재시도가 파티션을 막는 것을 직접 관찰
이것이 이 예제에서 가장 배울 것이 많은 실험입니다.
addNotRetryableExceptions를 일부러 주석 처리하고
백오프를 크게 늘린 뒤, 파싱 실패 레코드 하나를 넣고 그 뒤에 정상 레코드 10건을 넣으세요.
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에 쌓인 레코드를 되돌리는 것은 코드가 아니라 절차의 문제입니다. 원인을 고치기 전에 되돌리면 무한 루프가 됩니다.
- 원인 분류 —
kafka_dlt-exception-fqcn헤더로 그룹화합니다. 원인이 여러 가지면 각각 다르게 처리해야 합니다. - 수정 배포 — 코드 버그면 컨슈머를 고쳐 배포합니다. 데이터 문제면 상류 프로듀서를 고칩니다.
- 소량 재처리로 검증 — DLQ에서 몇 건만 원본 토픽으로 되돌려 성공하는지 확인합니다.
- 전량 재처리 — 확인 후 나머지를 되돌립니다.
되돌리는 방법은 두 가지입니다.
# dlq-replay.properties
# 같은 클러스터 안에서 토픽 이름만 바꿔 복사합니다.
# IdentityReplicationPolicy 를 쓰면 원본 토픽명에 접두어가 붙지 않습니다.
clusters = local
local.bootstrap.servers = kafka-1:19092,kafka-2:19092,kafka-3:19092
# 같은 클러스터로 되돌리므로 별칭을 하나만 쓰는 대신
# 실제 운영에서는 replay 전용 컨슈머 그룹으로 처리하는 편이 단순합니다.
# 자세한 MirrorMaker 2 구성은 예제 11을 보세요.
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 아티팩트에서 직접 확인했습니다.
- Consumer Configs —
max.poll.interval.ms=300000,max.poll.records=500,enable.auto.commit=true,auto.commit.interval.ms=5000,isolation.level=read_uncommitted,auto.offset.reset=latest - Producer Configs —
acks=all,enable.idempotence=true,linger.ms=5,delivery.timeout.ms=120000 - Broker Configs —
message.max.bytes=1048588,auto.create.topics.enable=true,num.partitions=1,default.replication.factor=1 - Topic Configs —
retention.ms=604800000(7일) - Connect Configs —
errors.tolerance=none,errors.log.enable=false,errors.retry.timeout=0,errors.retry.delay.max.ms=60000,errors.deadletterqueue.topic.replication.factor=3,errors.deadletterqueue.context.headers.enable=false - Spring for Apache Kafka Reference — Handling Exceptions —
DefaultErrorHandler,DeadLetterPublishingRecoverer,ExponentialBackOffWithMaxRetries,addNotRetryableExceptions,RetryListener - Spring for Apache Kafka Reference — Non-Blocking Retries —
@RetryableTopic,@DltHandler,TopicSuffixingStrategy - Spring Boot Reference — Apache Kafka Support —
spring.kafka.*프로퍼티,CommonErrorHandler빈 자동 적용