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
Original file line number Diff line number Diff line change
@@ -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 {
}
Original file line number Diff line number Diff line change
@@ -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 전용 팩토리 사용
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
public void consumeInventoryDLT(ConsumerRecord<String, String> 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<String, String> record, String headerKey) {
Header header = record.headers().lastHeader(headerKey);
if (header != null && header.value() != null) {
return new String(header.value(), StandardCharsets.UTF_8);
}
return "알 수 없음";
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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로 직행함
Expand Down
6 changes: 4 additions & 2 deletions src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ spring:
profiles:
active: local
jpa:
open-in-view: false
show-sql: false
properties:
hibernate:
Expand Down Expand Up @@ -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"
Expand All @@ -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"
Expand Down