diff --git a/src/main/java/com/michelet/inventory/infrastructure/config/JpaConfig.java b/src/main/java/com/michelet/inventory/infrastructure/config/JpaConfig.java new file mode 100644 index 0000000..a492f14 --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/config/JpaConfig.java @@ -0,0 +1,10 @@ +package com.michelet.inventory.infrastructure.config; + +import org.springframework.context.annotation.Configuration; +import org.springframework.data.jpa.repository.config.EnableJpaRepositories; + +@Configuration +// Redis와 스캔 충돌을 막기 위해 패키지 명시 +@EnableJpaRepositories(basePackages = "com.michelet.inventory.infrastructure.repository") +public class JpaConfig { +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java new file mode 100644 index 0000000..581c40f --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java @@ -0,0 +1,57 @@ +package com.michelet.inventory.infrastructure.messaging; + +import java.nio.charset.StandardCharsets; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.common.header.Header; +import org.springframework.kafka.annotation.KafkaListener; +import org.springframework.kafka.support.KafkaHeaders; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class DeadLetterConsumer { + + /** + * Inventory 서비스의 주요 도메인 로직(재고 복구 등) 실패 시 격리된 메시지를 수신 - 3회 재시도 후에도 실패한 '독성 메시지'를 분석하기 위한 전용 컨슈머 + */ + @KafkaListener( + topics = { + "${inventory.kafka.topic.restore-request:order.stock-restore.requested}.DLT", + "${inventory.kafka.topic.restored:stock.restored}.DLT", + "${inventory.kafka.topic.reserved:stock.reserved}.DLT" + }, + groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}-dlt", + containerFactory = "dltListenerContainerFactory" // String 전용 팩토리 사용 + ) + public void consumeInventoryDLT(ConsumerRecord record) { + String originalTopic = extractHeaderAsString(record, KafkaHeaders.DLT_ORIGINAL_TOPIC); + String exceptionMessage = extractHeaderAsString(record, KafkaHeaders.DLT_EXCEPTION_MESSAGE); + String stackTrace = extractHeaderAsString(record, KafkaHeaders.DLT_EXCEPTION_STACKTRACE); + + log.error(""" + ================================================================================ + [CRITICAL ALERT] 인벤토리 DLT 에러 메시지 격리 수신 (최종 실패) + 원본 토픽 : {} + 에러 원인 : {} + 원본 데이터 : {} + 상세 트레이스 : + {} + ================================================================================""", + originalTopic, + exceptionMessage, + record.value() != null ? record.value() : "데이터 없음", + stackTrace + ); + } + + private String extractHeaderAsString(ConsumerRecord record, String headerKey) { + Header header = record.headers().lastHeader(headerKey); + if (header != null && header.value() != null) { + return new String(header.value(), StandardCharsets.UTF_8); + } + return "알 수 없음"; + } +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java index 1760ea3..fe3bef8 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -16,10 +16,10 @@ public class OrderEventConsumer { private final StockLockFacade stockLockFacade; @KafkaListener( - topics = "${inventory.kafka.topic.restored:stock.restored}", + topics = "${inventory.kafka.topic.restore-request:order.stock-restore.requested}", groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" ) - public void consumeStockRestoredEvent(StockRestoreMessage payload) { + public void consumeStockRestoreRequest(StockRestoreMessage payload) { if (payload == null) { log.error("[Kafka Consumer] 잘못된 복구 이벤트 수신: payload가 null입니다 (Tombstone 메시지일 가능성)."); // IllegalArgumentException을 던지면 아래 catch 블록에서 잡지 않고 밖으로 던져져서 DLT로 직행함 diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index eb0fe4b..aebc366 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -4,6 +4,7 @@ spring: profiles: active: local jpa: + open-in-view: false show-sql: false properties: hibernate: @@ -31,7 +32,7 @@ spring: properties: # 실제 역직렬화를 수행할 델리게이트 클래스 지정 spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer - spring.json.trusted.packages: "com.michelet.inventory.infrastructure.messaging.dto" + spring.json.trusted.packages: "com.michelet.inventory.infrastructure.messaging.dto,com.michelet.order.application.dto" spring.json.use.type.headers: true # 서로 다른 패키지의 이벤트를 1:1로 매핑 spring.json.type.mapping: "com.michelet.order.application.dto.StockRestoreEventPayload:com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage" @@ -47,7 +48,8 @@ inventory: topic: product-created: "product.created" reserved: "stock.reserved" - restored: "stock.restored" + restore-request: "order.stock-restore.requested" # 오더의 '명령'을 수신할 토픽 + restored: "stock.restored" # 카탈로그로 '결과'를 발송할 토픽 status-changed: "product.status-changed" daily-reset: "stock.daily-reset" product-updated: "product.updated"