From 4d367a096a5983aeec66475c8434ba4dedb41f03 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Thu, 14 May 2026 15:10:00 +0900 Subject: [PATCH 1/2] =?UTF-8?q?fix:=20DLT=20=EC=BB=A8=EC=8A=88=EB=A8=B8=20?= =?UTF-8?q?=EA=B5=AC=EC=B6=95,=20=ED=86=A0=ED=94=BD=20=EC=9D=B4=EB=A6=84?= =?UTF-8?q?=20=EC=B6=A9=EB=8F=8C=20=EA=B2=B0=ED=95=A8=20=EC=88=98=EC=A0=95?= =?UTF-8?q?,=20Kafka=20=EC=97=AD=EC=A7=81=EB=A0=AC=ED=99=94=20=EC=8B=A0?= =?UTF-8?q?=EB=A2=B0=20=EC=A0=95=EC=B1=85=20=EC=99=84=ED=99=94,=20OSIV=20?= =?UTF-8?q?=EB=B9=84=ED=99=9C=EC=84=B1=ED=99=94=20=EB=B0=8F=20JpaConfig=20?= =?UTF-8?q?=EB=B6=84=EB=A6=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../infrastructure/config/JpaConfig.java | 10 ++++ .../messaging/DeadLetterConsumer.java | 57 +++++++++++++++++++ .../messaging/OrderEventConsumer.java | 2 +- src/main/resources/application.yml | 6 +- 4 files changed, 72 insertions(+), 3 deletions(-) create mode 100644 src/main/java/com/michelet/inventory/infrastructure/config/JpaConfig.java create mode 100644 src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java 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..eb67317 --- /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.stock-restored:stock.restored}.DLT", + "${inventory.kafka.topic.stock-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..a51a148 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -16,7 +16,7 @@ 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) { diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index eb0fe4b..1a85d67 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.*" 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" From 6d5ed3e3e40a7b4e04b7a35b91fb2eb6893fa7e2 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Thu, 14 May 2026 16:21:50 +0900 Subject: [PATCH 2/2] =?UTF-8?q?fix:=20=ED=86=A0=ED=94=BD=20=EC=9D=B4?= =?UTF-8?q?=EB=A6=84=20=EB=AF=B8=EC=8A=A4=EB=A7=A4=EC=B9=98=20=ED=95=B4?= =?UTF-8?q?=EA=B2=B0,=20=EB=A9=94=EC=84=9C=EB=93=9C=EB=AA=85=20=EC=88=98?= =?UTF-8?q?=EC=A0=95,=20=EC=8B=A0=EB=A2=B0=20=ED=8C=A8=ED=82=A4=EC=A7=80?= =?UTF-8?q?=20=EB=B2=94=EC=9C=84=20=EC=B6=95=EC=86=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../infrastructure/messaging/DeadLetterConsumer.java | 4 ++-- .../infrastructure/messaging/OrderEventConsumer.java | 2 +- src/main/resources/application.yml | 2 +- 3 files changed, 4 insertions(+), 4 deletions(-) diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java index eb67317..581c40f 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java @@ -20,8 +20,8 @@ public class DeadLetterConsumer { @KafkaListener( topics = { "${inventory.kafka.topic.restore-request:order.stock-restore.requested}.DLT", - "${inventory.kafka.topic.stock-restored:stock.restored}.DLT", - "${inventory.kafka.topic.stock-reserved:stock.reserved}.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 전용 팩토리 사용 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 a51a148..fe3bef8 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -19,7 +19,7 @@ public class OrderEventConsumer { 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 1a85d67..aebc366 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -32,7 +32,7 @@ spring: properties: # 실제 역직렬화를 수행할 델리게이트 클래스 지정 spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer - spring.json.trusted.packages: "com.michelet.*" + 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"