diff --git a/src/main/java/com/catchtable/global/config/KafkaConfig.java b/src/main/java/com/catchtable/global/config/KafkaConfig.java index f821ce1..9b67708 100644 --- a/src/main/java/com/catchtable/global/config/KafkaConfig.java +++ b/src/main/java/com/catchtable/global/config/KafkaConfig.java @@ -3,6 +3,7 @@ import io.micrometer.core.instrument.MeterRegistry; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.springframework.beans.factory.annotation.Value; @@ -11,6 +12,9 @@ import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.*; +import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; +import org.springframework.kafka.listener.DefaultErrorHandler; +import org.springframework.util.backoff.FixedBackOff; import java.util.HashMap; import java.util.Map; @@ -44,13 +48,19 @@ public ProducerFactory producerFactory() { return factory; } - // @RetryableTopic 이 retry/dlt 토픽 라우팅과 ErrorHandler 를 자체 관리한다. - // 여기서 commonErrorHandler 를 강제하면 RetryTopicConfigurer 의 ErrorHandler swap 이 - // 차단되어 main listener container 가 partition assignment 단계 진입에 실패한다. + // FixedBackOff(2000L, 3L): 2초 간격으로 3번 재시도. 메모리 retry 후 최종 실패 시 DLT publish. + // DLT 토픽 명은 "${original}-dlt" 규약을 따른다. broker 에 이미 그 이름으로 존재. @Bean - public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory() { + public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory(KafkaTemplate kafkaTemplate) { ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); + + DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer( + kafkaTemplate, + (record, ex) -> new TopicPartition(record.topic() + "-dlt", record.partition()) + ); + factory.setCommonErrorHandler(new DefaultErrorHandler(recoverer, new FixedBackOff(2000L, 3L))); + return factory; } diff --git a/src/main/java/com/catchtable/notification/service/NotificationKafkaConsumer.java b/src/main/java/com/catchtable/notification/service/NotificationKafkaConsumer.java index d9c17aa..93fd1c2 100644 --- a/src/main/java/com/catchtable/notification/service/NotificationKafkaConsumer.java +++ b/src/main/java/com/catchtable/notification/service/NotificationKafkaConsumer.java @@ -14,12 +14,7 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.data.redis.core.StringRedisTemplate; -import org.springframework.kafka.annotation.DltHandler; import org.springframework.kafka.annotation.KafkaListener; -import org.springframework.kafka.annotation.RetryableTopic; -import org.springframework.kafka.retrytopic.DltStrategy; -import org.springframework.kafka.support.KafkaHeaders; -import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Payload; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @@ -42,7 +37,6 @@ public class NotificationKafkaConsumer { private final StringRedisTemplate redisTemplate; private final VacancyService vacancyService; - @RetryableTopic(attempts = "3", dltStrategy = DltStrategy.ALWAYS_RETRY_ON_ERROR) @KafkaListener(topics = "notification.reservation.confirmed", groupId = "catchtable-notification-group") @Transactional public void handleReservationConfirmed(@Payload ReservationConfirmedEvent event) { @@ -64,7 +58,6 @@ public void handleReservationConfirmed(@Payload ReservationConfirmedEvent event) ); } - @RetryableTopic(attempts = "3", dltStrategy = DltStrategy.ALWAYS_RETRY_ON_ERROR) @KafkaListener(topics = "notification.reservation.canceled", groupId = "catchtable-notification-group") @Transactional public void handleReservationCanceled(@Payload ReservationCanceledEvent event) { @@ -86,7 +79,6 @@ public void handleReservationCanceled(@Payload ReservationCanceledEvent event) { ); } - @RetryableTopic(attempts = "3", dltStrategy = DltStrategy.ALWAYS_RETRY_ON_ERROR) @KafkaListener(topics = "notification.reservation.changed", groupId = "catchtable-notification-group") @Transactional public void handleReservationChanged(@Payload ReservationChangedEvent event) { @@ -110,7 +102,6 @@ public void handleReservationChanged(@Payload ReservationChangedEvent event) { ); } - @RetryableTopic(attempts = "3", dltStrategy = DltStrategy.ALWAYS_RETRY_ON_ERROR) @KafkaListener(topics = "notification.reservation.visited", groupId = "catchtable-notification-group") @Transactional public void handleReservationVisited(@Payload ReservationVisitedEvent event) { @@ -130,7 +121,6 @@ public void handleReservationVisited(@Payload ReservationVisitedEvent event) { ); } - @RetryableTopic(attempts = "3", dltStrategy = DltStrategy.ALWAYS_RETRY_ON_ERROR) @KafkaListener(topics = "notification.vacancy.opened", groupId = "catchtable-notification-group") @Transactional public void handleVacancyOpened(@Payload VacancyEvent event) { @@ -181,11 +171,6 @@ public void handleVacancyOpened(@Payload VacancyEvent event) { log.info("[Kafka Consumer] {}명에게 빈자리 알림을 생성했습니다. key={}", users.size(), redisKey); } - @DltHandler - public void handleDlt(Object message, @Header(KafkaHeaders.RECEIVED_TOPIC) String topic) { - log.error("[DLT] 메시지 처리 최종 실패. Topic: {}, Message: {}", topic, message.toString()); - } - private User findUserOrThrow(Long userId) { return userRepository.findById(userId) .orElseThrow(() -> {