grep

Engineering

실패한 메시지는 어디로 가야 할까?- Kafka Retry/DLT 설계와 운영에서 밟은 3가지 함정

Jolly여기어때

2026년 8월 5일

원문에서 보기 ↗

글. 이지혜(Jolly) / 정산개발팀

안녕하세요, 여기어때 정산개발팀 백엔드 개발자 Jolly입니다.

정산 시스템의 Kafka 컨슈머들은 저마다 다른 방식으로 실패를 처리하고 있었습니다.

오늘은 non-blocking 재시도와 공통 DLT(Dead Letter Topic)를 어노테이션 하나로 표준화한 공통 라이브러리를 만들며 내린 설계 결정들과, 그 과정에서 만난 Spring Kafka의 함정들을 공유하려고 합니다.

Kafka 컨슈머의 실패 처리를 고민하고 계신 분들께 도움이 되었으면 좋겠습니다.

아래 구현은 Java 17, Spring Boot 3.5, Spring Kafka 3.3 기준입니다.

바쁘신 분은 실제로 구현하면서 밟았던 3가지 함정부터 보셔도 됩니다.

들어가며 — 실패한 메시지는 어디로 가는가

정산 도메인에서 Kafka 메시지 하나는 곧 돈입니다. 주문 완료 이벤트를 놓치면 매출이 잡히지 않고, 환불 이벤트를 놓치면 정산 금액이 틀어집니다.

“실패하면 로그 남기고 넘어간다”가 허용되지 않는 도메인입니다.

그런데 실제 운영에서 컨슈머 실패는 생각보다 자주, 그리고 대부분 일시적인 이유로 발생합니다.

대표적인 것이 이벤트 소비 순서 역전입니다.

[기대 순서]  매출마감 생성  →  order.completed 소비
[실제 순서]  order.completed 소비  →  (아직 매출마감이 없음 → 실패)  →  매출마감 생성

선행 데이터가 몇 분 뒤에 도착하는 상황이라, 몇 분 후에 다시 시도하면 대부분 성공합니다.

문제는 이 “몇 분 후 재시도”를 각 서비스가 저마다 구현하고 있었다는 점입니다.

어떤 서비스는 RetryTemplate으로 blocking 재시도를 하고, 어떤 서비스는 로그만 남기고 넘어가기도 했습니다. 최종 실패 메시지의 보관 위치와 포맷도 서비스마다 달랐습니다.

여기서 두 가지 요구사항이 나왔습니다.

  1. 재시도는 파티션을 막지 않아야 한다 (non-blocking retry)
  2. 최종 실패는 유실 없이 한 곳에 표준 포맷으로 모여야 한다 (공통 DLT)

그리고 하나 더. 컨슈머 서비스가 이걸 적용하는 비용이 어노테이션 한 줄이어야 했습니다. 적용이 번거로우면 결국 아무도 안 쓰게 될 테니까요.

전체 흐름을 그림으로 먼저 보면 이렇습니다.

첫 번째 시도 — 인메모리 큐와 30분 스케줄러

사실 retry topic 이전에 시도했던 방법이 있습니다.

처음에는 실패 메시지를 인메모리 큐에 넣고, 30분마다 다시 처리했습니다. 간단했고 동작도 했지만, 재배포나 비정상 종료가 발생하면 대기 중인 메시지가 사라질 수 있었습니다.

더 큰 문제는 메시지를 큐로 옮기는 순간 Kafka offset은 이미 커밋된다는 점이었습니다. 실패 상태가 Kafka 밖에만 존재하니, 유실되면 offset reset 같은 Kafka의 복구 수단도 사용할 수 없었습니다.

결국 실패 메시지를 잃지 않으려면 실패 상태도 애플리케이션 메모리가 아니라 Kafka 안에 있어야 한다고 판단했습니다.

이를 RFC로 정리해 팀에 공유하고(Special Thanks To 이스트), 함께 검토한 끝에 이 방향으로 확정했습니다.

Blocking이냐, Non-blocking이냐

가장 먼저 결정할 것은 재시도 방식이었습니다.

리스너 안에서 기다렸다 다시 시도하는 blocking 방식과, 실패 메시지를 retry topic으로 보내두고 별도 컨슈머가 지연 후 재소비하는 non-blocking 방식이 있습니다.

우리의 실패 패턴은 “몇 분 뒤에 다시 시도”였습니다.

3분을 blocking으로 기다리면 그 파티션의 뒤 메시지가 전부 밀리고, 재시도가 겹치면 max.poll.interval.ms를 넘겨 컨슈머가 그룹에서 쫓겨나는 2차 장애로 번집니다. 그래서 non-blocking retry를 선택했습니다.

순서가 깨지는 대가는 도메인 특성으로 흡수했습니다.

매출 집계는 이벤트의 원본 좌표(topic/partition/offset)와 멱등 처리로 수렴하도록 설계되어 있어, 재시도 메시지가 뒤로 밀려도 최종 상태는 같습니다.

non-blocking retry를 도입하려면 컨슈머가 순서 역전과 중복에 안전한지부터 확인해야 합니다. 이게 안 되는 컨슈머라면 blocking이 맞는 선택일 수 있습니다.

Kafka Streams window는 왜 아니었나

“선행 데이터가 도착하면 처리한다”는 요구만 보면, Kafka Streams로 두 스트림을 window join하는 방법도 떠오릅니다.

하지만 우리 문제에는 전제가 하나 어긋납니다.

window는 “이 시간 구간 안에 상대가 온다”를 알아야 쓸 수 있는 도구인데, 선행 데이터가 언제 생길지는 예측할 수 없습니다.

window를 작게 잡으면 늦게 오는 데이터를 놓치고, 놓치지 않으려고 크게 잡으면 그만큼 state store가 커지고 처리도 늦어집니다.

어떤 크기를 잡아도 “언젠가는 온다”를 보장할 수 없었습니다.

반면 재시도 방식은 도착 시점을 몰라도 됩니다.

“일정 간격으로 다시 확인하고, 끝내 안 오면 DLT로 보내 사람이 확인한다”는 모델이 문제의 성격과 맞았습니다.

재시도 정책: 3회 × 3분, 그리고 단일 retry topic

retry topic 네이밍은 규칙으로 못 박았습니다.

retry.{originTopic}
예: retry.order.completed

재시도 대상은 opt-in — 안 되는 실패는 바로 DLT로

모든 예외를 재시도하는 건 낭비일뿐더러, 진짜 문제를 가리는 일이기도 합니다.

역직렬화 실패, 비즈니스 검증 실패는 3분 뒤에 다시 시도해도 똑같이 실패합니다. 그런 실패를 재시도하면 DLT에 도달하기까지 6분(3분 × 재시도 2회)을 통째로 낭비하는 셈입니다. 그래서 기본 정책을 뒤집었습니다.

기본은 재시도 없음. 라이브러리가 정의한 KafkaRetryException을 던진 경우에만 재시도하고, 그 외 모든 예외는 즉시 DLT로 보낸다.

구현은 간단합니다.

공통 어노테이션의 기본값에 include = KafkaRetryException.class 하나만 박아두면 됩니다. 컨슈머 코드에서 "이건 기다리면 되는 실패"라고 판단한 지점에서만 throw new KafkaRetryException(...)으로 변환해 던집니다.

재시도가 의미 있는 실패와 사람이 봐야 하는 실패를 코드에서 명시적으로 구분하게 만든 것이, 이 설계에서 가장 효과가 컸던 결정이라고 생각합니다.

DLT는 서비스별이 아니라 팀 공통 1개

Spring Kafka의 기본값은 리스너마다 {topic}-dlt를 만드는 것입니다.

우리는 반대로 모든 서비스의 최종 실패를 단일 공통 topic으로 모았습니다.

common.dead-letter

이유는 실패 처리의 뒷단이 전부 공통이기 때문입니다.

DB 적재, Slack 알림, 재처리 도구까지. 이걸 서비스 수만큼 만들 이유가 없다고 판단했습니다.

대신 payload를 표준화해야 합니다.

서로 다른 스키마의 메시지가 한 topic에 섞이므로, DLT 발행 시점에 원본 payload를 CommonDltMessage라는 공통 포맷으로 한 번 감쌉니다. 여기서 중요한 건 두 가지입니다.

원본 좌표(topic/partition/offset)를 retry topic이 아니라 origin 기준으로 담는 것, 그리고 정상 payload는 파싱 가능한 JSON 문자열로, 역직렬화에 실패한 payload는 raw 문자열로 최대한 원문 그대로 보존하는 것.

이렇게까지 보존하는 건 재처리(re-drive) 때문입니다. DLT는 끝이 아니라, 사람이 확인한 후 다시 흘려보낼 대기열이라고 봤습니다.

한 가지 더, 라이브러리의 책임은 DLT 발행에서 끝입니다.

DLT 소비(DB 저장, Slack 알림, 재처리)는 발행 라이브러리 밖의 read-side 책임으로 분리했고, 어노테이션에도 autoStartDltHandler = "false"를 설정해 컨슈머 서비스가 실수로 DLT 리스너를 띄우지 않게 했습니다.

발행 측 라이브러리에 소비 로직까지 넣으면, 모든 컨슈머 서비스가 DB 의존성과 Slack 설정을 끌고 들어가게 되기 때문입니다.

구현 — Spring Kafka의 확장 포인트 두 곳을 갈아끼우기

적용하는 쪽 코드부터 보겠습니다. 목표했던 “어노테이션 한 줄”입니다.

@CommonKafkaRetry(backoff = @Backoff(delayExpression = "${kafka.retry.order-completed.delay-ms}"))
@KafkaListener(
    topics = "${kafka.consumer.order-completed.topic}",
    groupId = "${kafka.consumer.order-completed.group-id}",
    ...
)
public void consume(ConsumerRecord<String, OrderCompleted> record) {
    orderCompletedService.completedOrder(record.value());
}

@CommonKafkaRetry는 @RetryableTopic을 감싼 meta-annotation 입니다. 정책 기본값(3회, 3분, SINGLE_TOPIC, opt-in 예외)을 어노테이션에 고정하고, 서비스별로 조정할 값만 @AliasFor로 열어두었습니다.

기본값을 벗어나려면 어노테이션 속성을 직접 바꿔야 하니, 정책에서 벗어나는 코드는 리뷰에서 바로 눈에 띕니다.

문제는 topic 네이밍과 DLT payload입니다.

Spring Kafka의 기본 동작(-retry-N/-dlt suffix, 원본 payload 그대로 발행)을 우리 규칙으로 바꿔야 하는데, 이 두 가지는 각각 다른 확장 포인트에 숨어 있습니다.

RetryTopicConfigurationSupport 확장

Spring Kafka는 retry topic 인프라 전체를 하나의 팩토리로 조립하는데, 이 팩토리를 통째로 갈아끼울 수 있는 공식 확장점이 RetryTopicConfigurationSupport입니다.

우리는 여기서 네이밍과 DLT 발행기, 딱 두 곳만 교체했습니다.

@Configuration
public class CommonRetryTopicConfigurationSupport extends RetryTopicConfigurationSupport {
  @Override
    protected RetryTopicComponentFactory createComponentFactory() {
        return new RetryTopicComponentFactory() {
            @Override
            public RetryTopicNamesProviderFactory retryTopicNamesProviderFactory() {
                return namingStrategy;   // (1) topic/group/client-id 네이밍 교체
            }
            @Override
            public DeadLetterPublishingRecovererFactory deadLetterPublishingRecovererFactory(
                    DestinationTopicResolver resolver) {
                DeadLetterPublishingRecovererFactory factory =
                        super.deadLetterPublishingRecovererFactory(resolver);
                factory.setDeadLetterPublisherCreator((templates, destinations) ->
                        new CommonDltPublishingRecoverer(...));   // (2) DLT payload 교체
                return factory;
            }
        };
    }
}

(1)의 네이밍 전략에서 주의할 점이 있습니다.

topic 이름만 바꾸면 안 됩니다. retry 컨슈머의 consumer group, client-id, endpoint id 에도 일관된 suffix(-retry, -dlt)를 줘야 main 컨슈머와 그룹이 분리되고, 모니터링에서 어느 컨슈머가 retry인지 식별할 수 있습니다.

CommonDltPublishingRecoverer — retry hop과 final DLT의 분기

DeadLetterPublishingRecoverer는 retry topic 발행과 DLT 발행에 같은 코드 경로를 씁니다. 그래서 createProducerRecord를 override해 목적지에 따라 분기했습니다.

공통 DLT일 때만 payload를 CommonDltMessage로 교체하고, retry hop은 Spring Kafka의 기본 발행 흐름을 유지합니다.

최종 DLT에서는 역직렬화 실패 원문까지 보존할 수 있도록 raw key/value도 함께 전달하고, DLT의 실제 partition은 producer가 선택하도록 두었습니다.

@Override
protected ProducerRecord<Object, Object> createProducerRecord(
        ConsumerRecord<?, ?> record,
        TopicPartition tp,
        Headers headers,
        byte[] key,
        byte[] value) {

    if (commonDltTopic.equals(tp.topic())) {
        CommonDltMessage payload =
                dltMessageFactory.from(record, currentException.get(), key, value);

        Object producerKey = key != null ? key : record.key();

        // 공통 DLT의 실제 파티션 중 하나를 producer가 선택
        return new ProducerRecord<>(
                tp.topic(), null, producerKey, payload, headers);
    }

    return super.createProducerRecord(record, tp, headers, key, value);
}

이 짧은 분기에서도 생각보다 여러 함정을 밟았는데, 바로 아래에서 하나씩 소개하겠습니다.

AutoConfiguration — 호스트 서비스와의 격리

라이브러리가 컨슈머 서비스의 기존 빈과 충돌하면 도입 장벽이 됩니다.

그래서 발행 경로 전체(ObjectMapper, ProducerFactory, KafkaTemplate)를 전용 빈으로 격리해 호스트 서비스의 직렬화 설정에 영향받지 않게 했습니다.

특히 날짜 포맷은 라이브러리가 팀 표준(yyyy-MM-dd HH:mm:ss)으로 고정합니다. 발행 측 서비스마다 포맷이 다르면 DLT를 소비하는 쪽이 깨지기 때문입니다.

정상 메시지만 테스트하면 놓치는 3가지 함정

사실 이 글에서 제일 하고 싶었던 이야기입니다.

공식 문서만으로는 놓치기 쉬워, 실제 구현과 디버깅 과정에서 확인한 것들입니다.

함정 1. 역직렬화된 값만 사용하면 실패 원문이 사라진다

처음에는 createProducerRecord()로 전달되는 raw byte[]를 그대로 발행했습니다.

그런데 retry 템플릿이 JsonSerializer와 StringSerializer만 사용하고 있어서 문제가 생겼습니다. value는 byte[]가 다시 JSON으로 직렬화되면서 Base64 형태가 됐고, key는 StringSerializer에 byte[]가 들어가면서 ClassCastException이 발생했습니다.

그래서 record.key()와 record.value(), 즉 역직렬화된 값으로 교체해서 발행했습니다. 정상 메시지만 놓고 보면 이걸로 해결된 것처럼 보였습니다.

하지만 역직렬화 자체가 실패한 메시지는 달랐습니다.

이 경우 record.value()는 null이고, 실제 원문은 createProducerRecord()의 raw key/value에 남아 있습니다. 이를 무조건 record.key()와 record.value()로 교체하면 정작 분석과 재처리에 필요한 실패 원문이 사라집니다.

그래서 전용 ProducerFactory에 DelegatingByTypeSerializer를 적용했습니다.

byte[]는 ByteArraySerializer, 일반 객체는 JsonSerializer, 문자열 key는 StringSerializer로 보내도록 타입별로 serializer를 나눴습니다.

정상 메시지는 기존 객체 형태로 발행하고, 역직렬화에 실패한 메시지는 retry hop에서 raw bytes를 유지합니다. 최종 DLT에서는 raw key/value가 있으면 이를 UTF-8 문자열로 변환해 CommonDltMessage에 우선 담습니다.

함정 2. 확장 포인트를 직접 교체하면 숨은 wiring도 책임져야 한다

createProducerRecord()에는 exception이 전달되지 않아, 진입점인 accept()에서 ThreadLocal에 보관한 뒤 꺼내 사용했습니다.

이 방식은 두 메서드가 같은 리스너 스레드에서 실행된다는 전제를 갖습니다.

또 RetryTopicConfigurationSupport를 직접 제공하면 Boot가 등록하던 backoff용 scheduler가 사라집니다. 그래서 라이브러리에서 RetryTopicSchedulerWrapper를 직접 등록했습니다.

Spring Kafka의 확장 포인트를 교체할 때는 메서드뿐 아니라 기존 auto-configuration이 제공하던 실행 전제와 보조 빈까지 확인해야 했습니다.

함정 3. 원본 좌표와 실제 DLT partition은 다르다

retry topic을 거친 레코드의 record.topic()/partition()/offset()은 원본이 아니라 retry topic의 좌표입니다.

원본 좌표는 Spring Kafka가 남긴 ORIGINAL_TOPIC/PARTITION/OFFSET 헤더의 첫 번째 값에서 복원해야 합니다.

이걸 놓치면 DLT에 쌓인 메시지가 전부 “retry topic에서 실패했다”고 기록되어, 재처리 시 원본 추적이 불가능해집니다.

반대로 공통 DLT에 발행할 때는 원본 partition을 그대로 지정하면 안 됩니다. 원본 토픽과 DLT의 파티션 수가 다를 수 있기 때문입니다.

따라서 실제 DLT partition은 producer가 선택하도록 null로 두고, 원본 topic/partition/offset은 CommonDltMessage에 별도로 보존했습니다.

테스트 — 빈이 아니라 wire format을 검증한다

전용 ObjectMapper, 전용 Serializer, 전용 Template처럼 wiring을 갈아끼우는 작업은 잘못돼도 어디서도 예외가 나지 않습니다.

깨진 포맷의 메시지가 조용히 발행되고, 한참 뒤에 다른 서비스의 역직렬화 실패로 나타납니다.

그래서 빈 wiring 테스트가 아니라 직렬화된 바이트(wire format) 자체를 단언하는 회귀 테스트를 두었습니다. 날짜 포맷 하나만 틀어져도 이 테스트가 먼저 깨집니다.

end-to-end 검증은 @EmbeddedKafka 통합 테스트로 재시도 발행부터 DLT payload의 origin 필드까지 확인합니다. backoff 3분은 테스트에서만 100ms로 override합니다.

소비하는 쪽 — DLT 컨슈머의 방어선

공통 DLT 컨슈머는 저장 전에 원본 좌표(topic + partition + offset) 기준 중복 체크를 합니다.

역직렬화가 깨져 좌표를 확인할 수 없는 레코드는 유실보다 중복이 낫다고 보고 우선 저장합니다.

중복만큼 중요하게 본 것이 저장 자체의 실패였습니다.

DB 저장에 실패하면 offset을 진행하지 않습니다.

ExponentialBackOff와 ContainerPausingBackOffHandler를 사용해 1초부터 최대 1분까지 간격을 늘리며 DLT 컨테이너를 pause/resume하고, 저장에 성공한 레코드만 커밋합니다.

DB의 TEXT byte 한도를 넘는 payload가 영구 실패 레코드가 되지 않도록 recordValue는 UTF-8 byte 기준으로 축약합니다.

저장 후에는 도메인별 Slack 채널로 알림을 보내며, Slack 일시 실패는 최대 3회 재시도합니다.

알림이 끝내 실패하더라도 DB에 저장된 실패 기록은 유지됩니다.

DLT는 단순히 쌓아두는 토픽이 아니라 사람이 확인하고 복구하는 운영 경로라,저장과 알림까지 있어야 완성이라고 생각했습니다.

마치며

남은 과제도 있습니다.

이번 작업에서 배운 것 두 가지만 남기면 이렇습니다.

이 글이 Kafka 컨슈머의 실패 처리와 DLT 운영을 고민하시는 분들께 조금이나마 도움이 되었으면 합니다.

감사합니다.