Kafka @RetryableTopic vs DefaultErrorHandler

seongwop·2026년 6월 26일

Spring

목록 보기
19/21

결제 서비스를 이벤트 기반 구조로 전환하면서 Kafka 메시지 처리 실패를 어떻게 다룰지 고민하게 되었다.

결제 서비스는 order.created, refund.requested, stock.failed 같은 이벤트를 소비한다. 이 이벤트들은 단순 로그성 메시지가 아니라 실제 결제 승인, 결제 취소, 환불 보상 흐름으로 이어진다. 따라서 메시지 처리 중 예외가 발생했을 때 그냥 로그만 남기고 끝낼 수 없었다.

필요했던 것은 크게 세 가지였다.

  • 일시적인 오류는 재시도
  • 잘못된 메시지는 DLT로 격리
  • 성공한 메시지만 offset commit

Spring Kafka에서는 이런 요구를 처리할 수 있는 대표적인 방식으로 @RetryableTopic과 DefaultErrorHandler가 있다. 처음에는 둘 중 어떤 방식을 선택해야 할지 고민했다.

@RetryableTopic

@RetryableTopic은 Spring Kafka에서 제공하는 non-blocking retry 방식이다.

리스너에서 예외가 발생하면 같은 consumer thread 안에서 계속 붙잡고 재시도하는 것이 아니라, 실패한 메시지를 별도의 retry topic으로 보낸다. 그리고 일정 시간이 지난 뒤 retry topic을 다시 소비한다.

예를 들면 이런 식이다.

@RetryableTopic(
        attempts = "3",
        backoff = @Backoff(delay = 1000),
        dltTopicSuffix = ".DLT"
)
@KafkaListener(topics = "order.created")
public void consume(String message) {
    // message 처리
}

이 방식의 장점은 명확하다. 지연 재시도에 강하다. 1분 뒤, 5분 뒤, 30분 뒤처럼 retry topic을 기반으로 재처리 흐름을 분리할 수 있다. 메시지를 처리하는 listener thread를 오래 점유하지 않는다는 점도 장점이다.

하지만 그만큼 운영해야 하는 topic도 늘어난다. 원본 topic 외에 retry topic, DLT topic이 추가된다. 이벤트 종류가 많아질수록 topic 구조도 복잡해진다.

즉 @RetryableTopic은 “실패 메시지를 별도 topic으로 보내고, 나중에 다시 소비하는 구조”에 가깝다.

DefaultErrorHandler

DefaultErrorHandler는 Spring Kafka listener container 레벨에서 예외를 처리하는 방식이다.

리스너에서 예외가 발생하면 DefaultErrorHandler가 해당 메시지를 재시도할지, 재시도를 끝내고 복구 처리할지 결정한다. DLT로 보내고 싶다면 DeadLetterPublishingRecoverer를 함께 사용한다.

예시는 이런 형태다.

DeadLetterPublishingRecoverer recoverer =
        new DeadLetterPublishingRecoverer(kafkaTemplate);

DefaultErrorHandler errorHandler = new DefaultErrorHandler(
        recoverer,
        new FixedBackOff(1000L, 2L)
);

errorHandler.addNotRetryableExceptions(
        NonRetryablePaymentException.class
);

이 방식은 retry topic을 늘리지 않고, listener container 내부에서 재시도 정책을 공통으로 관리할 수 있다. 짧은 재시도 후 DLT로 보내는 구조에 잘 맞는다.

반대로 긴 지연 재시도에는 적합하지 않을 수 있다. backoff 동안 listener 처리가 지연될 수 있기 때문이다.

즉 DefaultErrorHandler는 “현재 consumer 처리 흐름 안에서 재시도하고, 끝까지 실패하면 DLT로 격리하는 구조”에 가깝다.

@RetryableTopic vs DefaultErrorHandler

처음에는 @RetryableTopic이 더 좋아 보였다. 어노테이션 하나로 retry topic과 DLT를 구성할 수 있고, 코드도 직관적이기 때문이다.

하지만 프로젝트 상황을 다시 보니 지금 필요한 것은 복잡한 지연 재시도가 아니었다.

결제 서비스에서 우선 필요했던 것은 다음과 같았다.

  • consumer 공통 실패 처리 정책
  • 짧은 재시도 후 DLT 이동
  • 역직렬화 실패 같은 포이즌 필 즉시 격리
  • 수동 ack와 함께 성공한 메시지만 offset commit
  • Inbox, Redis 멱등성 레이어와의 조합

이 요구에는 DefaultErrorHandler가 더 잘 맞았다.

@RetryableTopic을 사용하면 retry topic이 추가되고, topic 운영 복잡도가 올라간다. 지금 단계에서는 “몇 분 뒤 재처리”보다 “실패 메시지를 명확히 재시도하고, 안 되면 DLT로 보낸다”가 더 중요했다.

그래서 최종적으로 DefaultErrorHandler를 선택했다.

적용 방식

결제 서비스에서는 Kafka consumer 설정에 DefaultErrorHandler를 등록했다.

@Bean
public CommonErrorHandler kafkaCommonErrorHandler(
        KafkaTemplate<String, String> kafkaTemplate
) {
    DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
            kafkaTemplate,
            (record, ex) -> new TopicPartition(resolveDltTopic(record), record.partition())
    );

    DefaultErrorHandler errorHandler = new DefaultErrorHandler(
            recoverer,
            new FixedBackOff(1000L, 2L)
    );

    errorHandler.addNotRetryableExceptions(NonRetryablePaymentException.class);

    errorHandler.setCommitRecovered(true);

    return errorHandler;
}

핵심은 네 가지다.

첫 번째는 DeadLetterPublishingRecoverer다.
재시도가 모두 실패한 메시지를 DLT로 보내는 역할을 한다. 프로젝트에서는 원본 topic에 따라 DLT topic을 명확히 매핑했다.

private String resolveDltTopic(ConsumerRecord<?, ?> record) {
    return switch (record.topic()) {
        case KafkaTopics.ORDER_CREATED -> KafkaTopics.ORDER_CREATED_DLT;
        case KafkaTopics.REFUND_REQUESTED -> KafkaTopics.REFUND_REQUESTED_DLT;
        case KafkaTopics.STOCK_FAILED -> KafkaTopics.STOCK_FAILED_DLT;
        default -> record.topic() + ".DLT";
    };
}

두 번째는 FixedBackOff다.
일시적인 오류를 몇 번 재시도할지 정한다. 예를 들어 new FixedBackOff(1000L, 2L)는 1초 간격으로 2번 재시도한다. 최초 시도까지 포함하면 총 3번 처리 기회가 생긴다.

세 번째는 not-retryable exception 설정이다.
역직렬화 실패처럼 메시지 자체가 잘못된 경우는 재시도해도 해결되지 않는다. 이런 메시지를 포이즌 필이라고 볼 수 있다. 그래서 바로 DLT로 보내는 것이 낫다.

이런 예외들은 NonRetryablePaymentException를 뱉게 만들고 해당 예외를 not-retryable로 등록하면, 잘못된 payload는 불필요한 재시도 없이 DLT로 이동한다.

네 번째는 MANUAL_IMMEDIATE ack mode다.

factory.getContainerProperties()
        .setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);

그리고 listener에서는 모든 처리가 끝난 뒤 직접 ack를 호출한다.

@KafkaListener(topics = KafkaTopics.ORDER_CREATED)
public void consumeOrderCreated(String message, Acknowledgment acknowledgment) {
    OrderCreatedEvent event = readValue(message, OrderCreatedEvent.class, KafkaTopics.ORDER_CREATED);
    paymentEventService.handleOrderCreated(event);
    acknowledgment.acknowledge();
}

이렇게 하면 역직렬화나 비즈니스 로직에서 예외가 발생했을 때 acknowledge()가 호출되지 않는다. 실패 메시지는 DefaultErrorHandler의 재시도와 DLT 흐름을 탄다.


성능만 놓고 보면 재처리에 consumer 스레드가 묶이지 않는 retry topic 기반의 @RetryableTopic이 더 유리할 수도 있다고 생각한다.

실패 메시지를 원본 consumer 흐름에서 분리하기 때문에 긴 backoff가 필요하거나 실패율이 높은 환경에서는 정상 메시지 처리를 덜 방해한다.

하지만 이번 프로젝트에서는 결제라는 도메인의 특성으로 인해 짧은 재시도 후 DLT 격리가 목적이었고, retry topic 증가로 인한 운영 복잡도를 아직 감당할 필요가 없다고 판단했다. 그것이 DefaultErrorHandler를 선택하게 된 이유이다.

0개의 댓글