컨슈머에서 오류가 났다면 크게 나누어 봤을 때 재시도 할 수 없는 오류와 재시도 할 수 있는 오류로 나눌 수 있다.
특징
예시 상황
처리 방식
특징
예시 상황
처리 방식
나의 경우는 외부 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에는 원본 메시지의 메타데이터(토픽, 파티션, 오프셋, 예외 정보)가 헤더로 기록되어 재처리 및 디버깅에 활용 가능하다.