Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions src/main/java/com/catchtable/global/config/KafkaConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -44,13 +48,19 @@ public ProducerFactory<String, Object> 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<String, Object> kafkaListenerContainerFactory() {
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory(KafkaTemplate<String, Object> kafkaTemplate) {
ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());

DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(
kafkaTemplate,
(record, ex) -> new TopicPartition(record.topic() + "-dlt", record.partition())
);
Comment on lines +58 to +61

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

원밸 픽이 여러 개의 파티션을 가지고 잋지만 DLT 픽(-dlt)은 다른 파티션 개수(예: 1개)를 가지고 잋을 경우, record.partition()을 그대로 사용하면 DLT 픽에 존재하지 않는 파티션 인덱스로 메시지를 전섑하려고 시도하여 InvalidPartitionException이 발생할 수 잋습니다.\n\n읻 문제를 방지하기 위해 TopicPartition 생성 시 파티션 번호로 -1을 전달하는 쑔이 안전합니다. DeadLetterPublishingRecoverer는 음수 파티션 값을 null로 처리하여, Kafka Producer의 기밸 파티셔너(Partitioner)가 DLT 픽의 파티션을 적절히 결정하도록 위임합니다.

        DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(\n                kafkaTemplate,\n                (record, ex) -> new TopicPartition(record.topic() + "-dlt", -1)\n        );

factory.setCommonErrorHandler(new DefaultErrorHandler(recoverer, new FixedBackOff(2000L, 3L)));

return factory;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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) {
Expand All @@ -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) {
Expand All @@ -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) {
Expand All @@ -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) {
Expand All @@ -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) {
Expand Down Expand Up @@ -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(() -> {
Expand Down
Loading