[async] consumer에서 오류가 난다면?

maxxyoung·2025년 8월 8일

Async API

목록 보기
5/5

오류의 상황

컨슈머에서 오류가 났다면 크게 나누어 봤을 때 재시도 할 수 없는 오류와 재시도 할 수 있는 오류로 나눌 수 있다.

재시도 불가능한 오류 (Non-retryable error)

특징

  • 메시지 데이터를 다시 읽어도 항상 실패할 가능성이 높은 경우
  • 데이터 자체가 잘못되었거나, 비즈니스 로직 상 유효하지 않은 경우
  • 재시도를 해봤자 같은 오류가 반복됨 → 즉, 장애 원인이 메시지에 내재됨

예시 상황

  • 필수 필드 누락 (null 값이 들어온 주문 ID)
  • 잘못된 포맷(JSON 파싱 실패, 스키마 불일치)
  • 도메인 규칙 위반 (예: 가격이 음수)
  • 메시지가 너무 오래되어 이미 처리할 가치가 없는 경우(타임아웃된 주문)

처리 방식

  • 해당 메시지를 Dead Letter Topic(DLT)로 이동
  • 에러 로그 및 모니터링 전송
  • 필요시 운영자가 DLT 데이터를 검토 후 재처리 결정

재시도 가능한 오류 (Retryable error)

특징

  • 메시지 데이터는 유효하지만, 외부 요인이나 일시적인 문제로 실패한 경우
  • 재시도 시 성공 가능성이 높음
  • 장애 원인이 일시적이거나 환경적

예시 상황

  • 외부 API 서버 장애나 타임아웃
  • DB 연결 일시적 끊김
  • 네트워크 단절
  • 일시적인 락 경합(동시성 문제)

처리 방식

  • 재시도 로직 적용
  • 백오프 전략 적용 (exponential backoff)
  • 최대 재시도 횟수 초과 시 DLT로 이동

오류의 처리

나의 경우는 외부 API 연동이 실패했을 때 이므로 백오프 시도 후 -> DLT에 적재로 설정했다.

@Configuration
class KafkaTopicsConfig {
    @Bean
    fun priceCompareRequests() = NewTopic("price-compare-requests", 3, 1)

    @Bean
    fun priceCompareRequestsDlt() = NewTopic("price-compare-requests.DLT", 3, 1)
}

로컬에서 실행할 때의 편의 성을 위해 토픽을 빈으로 생성했다. 이렇게 설정하면 스프링이 처음 시작하면서 토픽도 같이 생성이 된다. DLT 토픽을 생성한다.

@Bean
    fun errorHandler(kafkaTemplate: KafkaTemplate<Any, Any>): DefaultErrorHandler {
        // 재시도: 최대 3회, 1초 → 2배씩 → 최대 5초
        val backOff = ExponentialBackOffWithMaxRetries(3).apply {
            initialInterval = 1000L
            multiplier = 2.0
            maxInterval = 5000L
        }

        val recoverer = DeadLetterPublishingRecoverer(kafkaTemplate) { record, _ ->
            TopicPartition(record.topic() + ".DLT", record.partition())
        }

        return DefaultErrorHandler(recoverer, backOff).apply {
            addNotRetryableExceptions(
                org.apache.kafka.common.errors.SerializationException::class.java,
                org.springframework.kafka.support.serializer.DeserializationException::class.java,
                com.fasterxml.jackson.core.JsonProcessingException::class.java
            )
        }
    }
    
@Bean
    fun kafkaListenerContainerFactory(
        consumerFactory: ConsumerFactory<String, String>,
        errorhandler: DefaultErrorHandler
    ): ConcurrentKafkaListenerContainerFactory<String, String> {
        val factory = ConcurrentKafkaListenerContainerFactory<String, String>()
        factory.consumerFactory = consumerFactory
        factory.isBatchListener = false
        factory.containerProperties.ackMode = ContainerProperties.AckMode.RECORD
        factory.setCommonErrorHandler(errorhandler)
        return factory
    }    

에러 핸들러를 등록한다. json 파싱 예외, 직렬화, 역직렬화의 경우 마로 DLT에 들어가게 설정했다. 이외의 예외는 최대 3회, 1초 → 2배씩 → 최대 5초 시도하고 실패 시 DLT에 들어간다.
DLT에는 원본 메시지의 메타데이터(토픽, 파티션, 오프셋, 예외 정보)가 헤더로 기록되어 재처리 및 디버깅에 활용 가능하다.

profile
오직 나만을 위한 글. 틀린 부분 말씀해 주시면 감사드립니다.

0개의 댓글