From 82c5aaece25219e92c587d89c6a212ad5a387ef9 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Sat, 16 May 2026 02:59:01 +0900 Subject: [PATCH 1/4] =?UTF-8?q?fix:=20=EC=A3=BC=EB=AC=B8=20=EC=83=9D?= =?UTF-8?q?=EC=84=B1=20=EB=B6=80=EB=B6=84=EB=8F=84=20=EB=B9=84=EB=8F=99?= =?UTF-8?q?=EA=B8=B0=EB=A1=9C=20=EB=B3=80=EA=B2=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../application/InventoryOutboxScheduler.java | 12 +++++ .../application/StockCommandService.java | 45 ++++++++++++++++++ .../application/StockLockFacade.java | 46 +++++++++++++++++++ .../messaging/DeadLetterConsumer.java | 3 +- .../messaging/OrderEventConsumer.java | 28 +++++++++++ .../messaging/dto/OrderApprovedEvent.java | 11 +++++ .../messaging/dto/OrderCreatedMessage.java | 33 +++++++++++++ .../messaging/dto/OrderRejectedEvent.java | 14 ++++++ src/main/resources/application.yml | 6 ++- 9 files changed, 195 insertions(+), 3 deletions(-) create mode 100644 src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java create mode 100644 src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java create mode 100644 src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java diff --git a/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java b/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java index 4ec3f37..3be604c 100644 --- a/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java +++ b/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java @@ -10,6 +10,8 @@ import com.michelet.inventory.domain.model.InventoryOutbox; import com.michelet.inventory.domain.model.OutboxStatus; import com.michelet.inventory.domain.repository.InventoryOutboxRepository; +import com.michelet.inventory.infrastructure.messaging.dto.OrderApprovedEvent; +import com.michelet.inventory.infrastructure.messaging.dto.OrderRejectedEvent; import java.util.List; import java.util.UUID; import lombok.RequiredArgsConstructor; @@ -39,6 +41,8 @@ public class InventoryOutboxScheduler { private static final String EVENT_STOCK_RESERVED = "STOCK_RESERVED"; private static final String EVENT_STOCK_RESTORED = "STOCK_RESTORED"; private static final String EVENT_DAILY_RESET = "DAILY_STOCK_RESET"; + private static final String EVENT_ORDER_APPROVED = "ORDER_APPROVED"; + private static final String EVENT_ORDER_REJECTED = "ORDER_REJECTED"; @Value("${inventory.kafka.topic.product-created:product.created}") private String topicProductCreated; @@ -52,6 +56,10 @@ public class InventoryOutboxScheduler { private String topicStockRestored; @Value("${inventory.kafka.topic.daily-reset:stock.daily-reset}") private String topicDailyReset; + @Value("${inventory.kafka.topic.order-approved:order.approved}") + private String topicOrderApproved; + @Value("${inventory.kafka.topic.order-rejected:order.rejected}") + private String topicOrderRejected; @Scheduled(fixedDelay = 5000) public void processOutboxEvents() { @@ -119,6 +127,8 @@ private Object deserializePayload(String eventType, String jsonPayload) throws E case EVENT_STOCK_RESERVED -> objectMapper.readValue(jsonPayload, StockReservedEvent.class); case EVENT_STOCK_RESTORED -> objectMapper.readValue(jsonPayload, StockRestoredEvent.class); case EVENT_DAILY_RESET -> objectMapper.readValue(jsonPayload, DailyStockResetEvent.class); + case EVENT_ORDER_APPROVED -> objectMapper.readValue(jsonPayload, OrderApprovedEvent.class); + case EVENT_ORDER_REJECTED -> objectMapper.readValue(jsonPayload, OrderRejectedEvent.class); // 매핑 안 된 이벤트를 String으로 보내면 직렬화 에러 발생! 예외를 던져서 스케줄러 재시도 루프로 넘김 default -> { log.warn("등록되지 않은 알 수 없는 이벤트 타입입니다: {}", eventType); @@ -135,6 +145,8 @@ private String resolveTopic(String eventType) { case EVENT_STOCK_RESERVED -> topicStockReserved; case EVENT_STOCK_RESTORED -> topicStockRestored; case EVENT_DAILY_RESET -> topicDailyReset; + case EVENT_ORDER_APPROVED -> topicOrderApproved; + case EVENT_ORDER_REJECTED -> topicOrderRejected; default -> "inventory.unknown.event"; }; } diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index 26607ac..7bf3da0 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -13,8 +13,12 @@ import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; +import com.michelet.inventory.infrastructure.messaging.dto.OrderApprovedEvent; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; import com.michelet.inventory.presentation.dto.ReserveStockRequest; import com.michelet.inventory.presentation.dto.RestoreStockRequest; +import java.util.ArrayList; +import java.util.List; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -100,4 +104,45 @@ public void restoreStock(RestoreStockRequest request) { ); outboxHelper.append("STOCK", request.optionId().toString(), "STOCK_RESTORED", event); } + + // 다중 재고 비동기 처리 트랜잭션 + @Transactional + public void processOrderCreation(OrderCreatedMessage msg) { + if (processedEventRepository.existsById(msg.eventId())) { + log.warn("[Idempotency] 이미 처리된 주문 메시지입니다. reservationId={}", msg.reservationId()); + return; + } + + List modifiedStocks = new ArrayList<>(); + + // 모든 아이템 재고 차감 시도 (실패 시 BusinessException 발생하여 Facade -> Consumer 로 롤백됨) + for (var item : msg.items()) { + Stock stock = stockRepository.findById(item.optionId()) + .orElseThrow(StockNotFoundException::new); + + stock.reserve(item.quantity()); + modifiedStocks.add(stock); + + // 품절 처리 이벤트 발송 준비 + if (stock.getTotalQuantity() == 0) { + ProductOption option = productOptionRepository.findById(stock.getOptionId()).orElseThrow(); + Product product = option.getProduct(); + if (product.getStatus() != ProductStatus.SOLDOUT && product.getStatus() != ProductStatus.DELETED + && product.getStatus() != ProductStatus.EXPIRED) { + product.changeStatus(ProductStatus.SOLDOUT); + productRepository.save(product); + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_STATUS_CHANGED", + new ProductStatusChangedEvent(product.getId(), product.getStatus().name())); + } + } + } + + stockRepository.saveAll(modifiedStocks); + processedEventRepository.save(new ProcessedEvent(msg.eventId())); + + // 모두 성공 시 승인 이벤트 적재 + outboxHelper.append("ORDER", msg.reservationId().toString(), "ORDER_APPROVED", + new OrderApprovedEvent(msg.reservationId())); + log.info("[Inventory Saga] 다중 주문 재고 선점 성공 -> 승인 이벤트 적재 완료: {}", msg.reservationId()); + } } diff --git a/src/main/java/com/michelet/inventory/application/StockLockFacade.java b/src/main/java/com/michelet/inventory/application/StockLockFacade.java index b7f6630..323b1c0 100644 --- a/src/main/java/com/michelet/inventory/application/StockLockFacade.java +++ b/src/main/java/com/michelet/inventory/application/StockLockFacade.java @@ -1,7 +1,14 @@ package com.michelet.inventory.application; +import com.michelet.common.exception.BusinessException; +import com.michelet.inventory.domain.exception.InventoryErrorCode; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage.OrderItemDto; import com.michelet.inventory.presentation.dto.ReserveStockRequest; import com.michelet.inventory.presentation.dto.RestoreStockRequest; +import java.util.ArrayList; +import java.util.List; +import java.util.UUID; import java.util.concurrent.TimeUnit; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -67,4 +74,43 @@ public void restoreStockWithLock(RestoreStockRequest request) { } } } + + // 3. 다중 품목 비동기 재고 선점 (OrderCreatedMessage 수신용) + public void reserveOrderStocksWithLock( + OrderCreatedMessage msg) { + // 데드락 방지를 위해 OptionId를 오름차순으로 정렬하여 락을 순차적으로 획득 + List sortedOptionIds = msg.items().stream() + .map(OrderItemDto::optionId) + .sorted() + .toList(); + + List locks = new ArrayList<>(); + try { + for (java.util.UUID id : sortedOptionIds) { + RLock lock = redissonClient.getLock("stock:" + id); + if (lock.tryLock(5, TimeUnit.SECONDS)) { + locks.add(lock); + } else { + log.error("[StockLockFacade] 다중 재고 선점 락 획득 실패 - OptionId: {}", id); + throw new BusinessException( + InventoryErrorCode.CONCURRENCY_ERROR); + } + } + // 모든 락 획득 성공 시 비즈니스 로직 실행 + stockCommandService.processOrderCreation(msg); + + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new BusinessException( + InventoryErrorCode.CONCURRENCY_ERROR); + } finally { + // 데드락 및 리소스 낭비 방지를 위해 역순으로 락 해제 + for (int i = locks.size() - 1; i >= 0; i--) { + RLock lock = locks.get(i); + if (lock != null && lock.isHeldByCurrentThread()) { + lock.unlock(); + } + } + } + } } 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 581c40f..5181505 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java @@ -21,7 +21,8 @@ public class DeadLetterConsumer { 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" + "${inventory.kafka.topic.reserved:stock.reserved}.DLT", + "${inventory.kafka.topic.order-created:order.created}.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 9f9563f..ffc59e4 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -1,6 +1,10 @@ package com.michelet.inventory.infrastructure.messaging; +import com.michelet.common.exception.BusinessException; +import com.michelet.inventory.application.InventoryOutboxHelper; import com.michelet.inventory.application.StockLockFacade; +import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; +import com.michelet.inventory.infrastructure.messaging.dto.OrderRejectedEvent; import com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage; import com.michelet.inventory.presentation.dto.RestoreStockRequest; import lombok.RequiredArgsConstructor; @@ -14,6 +18,30 @@ public class OrderEventConsumer { private final StockLockFacade stockLockFacade; + private final InventoryOutboxHelper outboxHelper; + + // 오더 생성 이벤트를 비동기로 받아서 다중 락 처리 + @KafkaListener( + topics = "${inventory.kafka.topic.order-created:order.created}", + groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" + ) + public void consumeOrderCreated(OrderCreatedMessage payload) { + log.info("[Kafka Consumer] 신규 주문 생성 메시지 수신 -> 재고 다중 차감 시도: reservationId={}", payload.reservationId()); + try { + stockLockFacade.reserveOrderStocksWithLock(payload); + } catch (IllegalArgumentException e) { + log.error("[Kafka Consumer] 비즈니스/검증 룰 위반 에러 (DLT 직행 대상). reservationId: {}", payload.reservationId(), e); + throw e; + } catch (BusinessException e) { + // 품절, 한도 초과 등 비즈니스 예외 시 오더 서비스에 거절 이벤트 전송 + log.warn("[Kafka Consumer] 비즈니스 로직에 의한 재고 차감 실패 (거절 이벤트 정상 발행). 사유: {}", e.getMessage()); + outboxHelper.append("ORDER", payload.reservationId().toString(), "ORDER_REJECTED", + new OrderRejectedEvent(payload.reservationId(), e.getMessage())); + } catch (Exception e) { + log.error("[Kafka Consumer] 재고 차감 중 일시적 에러 발생 (재시도 대상). reservationId: {}", payload.reservationId(), e); + throw new RuntimeException("재고 다중 차감 실패", e); + } + } @KafkaListener( topics = "${inventory.kafka.topic.restore-request:order.stock-restore.requested}", diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java new file mode 100644 index 0000000..f2a524c --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderApprovedEvent.java @@ -0,0 +1,11 @@ +package com.michelet.inventory.infrastructure.messaging.dto; + +import java.util.UUID; + +public record OrderApprovedEvent(UUID reservationId) { + public OrderApprovedEvent { + if (reservationId == null) { + throw new IllegalArgumentException("reservationId는 null일 수 없습니다."); + } + } +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java new file mode 100644 index 0000000..715a9ac --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java @@ -0,0 +1,33 @@ +package com.michelet.inventory.infrastructure.messaging.dto; + +import java.util.List; +import java.util.UUID; + +public record OrderCreatedMessage( + UUID eventId, + UUID reservationId, + List items +) { + public OrderCreatedMessage { + if (eventId == null) { + throw new IllegalArgumentException("메시지 파싱 오류: eventId는 null일 수 없습니다."); + } + if (reservationId == null) { + throw new IllegalArgumentException("메시지 파싱 오류: reservationId는 null일 수 없습니다."); + } + if (items == null || items.isEmpty()) { + throw new IllegalArgumentException("메시지 파싱 오류: 주문 항목은 1개 이상이어야 합니다."); + } + } + + public record OrderItemDto(UUID optionId, Integer quantity) { + public OrderItemDto { + if (optionId == null) { + throw new IllegalArgumentException("메시지 파싱 오류: optionId는 null일 수 없습니다."); + } + if (quantity == null || quantity <= 0) { + throw new IllegalArgumentException("메시지 파싱 오류: 수량은 1 이상이어야 합니다."); + } + } + } +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java new file mode 100644 index 0000000..115e0ea --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderRejectedEvent.java @@ -0,0 +1,14 @@ +package com.michelet.inventory.infrastructure.messaging.dto; + +import java.util.UUID; + +public record OrderRejectedEvent(UUID reservationId, String reason) { + public OrderRejectedEvent { + if (reservationId == null) { + throw new IllegalArgumentException("reservationId는 null일 수 없습니다."); + } + if (reason == null || reason.isBlank()) { + throw new IllegalArgumentException("거절 사유(reason)는 필수입니다."); + } + } +} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index aebc366..ca721d4 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -35,8 +35,7 @@ spring: 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" - + spring.json.type.mapping: "com.michelet.order.application.dto.StockRestoreEventPayload:com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage,com.michelet.order.application.dto.OrderCreatedEventPayload:com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage" server: port: 19900 @@ -53,6 +52,9 @@ inventory: status-changed: "product.status-changed" daily-reset: "stock.daily-reset" product-updated: "product.updated" + order-created: "order.created" + order-approved: "order.approved" + order-rejected: "order.rejected" consumer: retry: interval-ms: 1000 From 4b909e899688d0b3ddb21002de5c7aee08865b9c Mon Sep 17 00:00:00 2001 From: ji-circle Date: Sat, 16 May 2026 19:27:42 +0900 Subject: [PATCH 2/4] =?UTF-8?q?fix:=20=EC=A3=BC=EB=AC=B8=20=EC=83=9D?= =?UTF-8?q?=EC=84=B1=20=EA=B3=BC=EC=A0=95=EC=9D=98=20=EB=B9=84=EB=8F=99?= =?UTF-8?q?=EA=B8=B0=20=EC=A0=84=ED=99=98=EC=97=90=20=EB=94=B0=EB=A5=B8=20?= =?UTF-8?q?=EC=BD=94=EB=93=9C=20=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../application/InventoryOutboxHelper.java | 52 ++++++++++++++----- .../application/StockCommandService.java | 4 +- .../application/StockLockFacade.java | 4 +- .../dto/ReserveStockRequest.java | 2 +- .../dto/RestoreStockRequest.java | 2 +- .../messaging/DeadLetterConsumer.java | 6 +-- .../messaging/OrderEventConsumer.java | 52 ++++++++++++++----- .../messaging/dto/OrderCreatedMessage.java | 7 +++ .../presentation/InternalStockController.java | 41 --------------- src/main/resources/application.yml | 6 +-- .../application/StockCommandServiceTest.java | 4 +- .../StockConcurrencyIntegrationTest.java | 2 +- src/test/resources/application-test.yml | 2 +- 13 files changed, 100 insertions(+), 84 deletions(-) rename src/main/java/com/michelet/inventory/{presentation => application}/dto/ReserveStockRequest.java (90%) rename src/main/java/com/michelet/inventory/{presentation => application}/dto/RestoreStockRequest.java (84%) delete mode 100644 src/main/java/com/michelet/inventory/presentation/InternalStockController.java diff --git a/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java b/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java index 1f0e3bd..2dfae00 100644 --- a/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java +++ b/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java @@ -33,21 +33,10 @@ public void validateMaxRetries() { } } - // MANDATORY로 변경하여 부모 트랜잭션이 없으면 즉각 실패하도록 원자성 강제 + // MANDATORY로 변경하여 부모 트랜잭션이 없으면 즉각 실패하도록 원자성 강제 (정상 비즈니스 로직용) @Transactional(propagation = Propagation.MANDATORY) public void append(String aggregateType, String aggregateId, String eventType, Object payloadObj) { - if (aggregateType == null || aggregateType.isBlank()) { - throw new IllegalArgumentException("aggregateType은 필수입니다."); - } - if (aggregateId == null || aggregateId.isBlank()) { - throw new IllegalArgumentException("aggregateId는 필수입니다."); - } - if (eventType == null || eventType.isBlank()) { - throw new IllegalArgumentException("eventType은 필수입니다."); - } - if (payloadObj == null) { - throw new IllegalArgumentException("payloadObj는 필수입니다."); - } + validateInputs(aggregateType, aggregateId, eventType, payloadObj); try { String payloadJson = objectMapper.writeValueAsString(payloadObj); @@ -65,6 +54,27 @@ public void append(String aggregateType, String aggregateId, String eventType, O } } + // 부모 트랜잭션이 롤백된 후, Catch 블록 등에서 독립적으로 보상/거절 이벤트를 기록할 때 사용 (Saga 롤백용) + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void appendIndependent(String aggregateType, String aggregateId, String eventType, Object payloadObj) { + validateInputs(aggregateType, aggregateId, eventType, payloadObj); + + try { + String payloadJson = objectMapper.writeValueAsString(payloadObj); + InventoryOutbox outbox = InventoryOutbox.builder() + .aggregateType(aggregateType) + .aggregateId(aggregateId) + .eventType(eventType) + .payload(payloadJson) + .build(); + outboxRepository.save(outbox); + log.info("[Inventory Outbox] 독립 트랜잭션 이벤트 적재 요청 (롤백/거절용): type={}, id={}", eventType, aggregateId); + } catch (JsonProcessingException e) { + log.error("Outbox 페이로드 직렬화 실패. aggregateId={}, eventType={}", aggregateId, eventType, e); + throw new RuntimeException("Outbox 이벤트 생성 중 오류가 발생했습니다.", e); + } + } + // 2단계 스케줄러에서 상태 업데이트 시 사용할 독립 트랜잭션 메서드 @Transactional(propagation = Propagation.REQUIRES_NEW) public void markAsPublished(UUID outboxId) { @@ -95,4 +105,20 @@ public void handleFailure(UUID outboxId) { outboxRepository.save(outbox); }); } + + // 공통 파라미터 검증 메서드 + private void validateInputs(String aggregateType, String aggregateId, String eventType, Object payloadObj) { + if (aggregateType == null || aggregateType.isBlank()) { + throw new IllegalArgumentException("aggregateType은 필수입니다."); + } + if (aggregateId == null || aggregateId.isBlank()) { + throw new IllegalArgumentException("aggregateId는 필수입니다."); + } + if (eventType == null || eventType.isBlank()) { + throw new IllegalArgumentException("eventType은 필수입니다."); + } + if (payloadObj == null) { + throw new IllegalArgumentException("payloadObj는 필수입니다."); + } + } } diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index 7bf3da0..22bb134 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -1,6 +1,8 @@ package com.michelet.inventory.application; import com.michelet.inventory.application.dto.ProductStatusChangedEvent; +import com.michelet.inventory.application.dto.ReserveStockRequest; +import com.michelet.inventory.application.dto.RestoreStockRequest; import com.michelet.inventory.application.dto.StockReservedEvent; import com.michelet.inventory.application.dto.StockRestoredEvent; import com.michelet.inventory.domain.exception.StockNotFoundException; @@ -15,8 +17,6 @@ import com.michelet.inventory.domain.repository.StockRepository; import com.michelet.inventory.infrastructure.messaging.dto.OrderApprovedEvent; import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; import java.util.ArrayList; import java.util.List; import lombok.RequiredArgsConstructor; diff --git a/src/main/java/com/michelet/inventory/application/StockLockFacade.java b/src/main/java/com/michelet/inventory/application/StockLockFacade.java index 323b1c0..1b5dc2a 100644 --- a/src/main/java/com/michelet/inventory/application/StockLockFacade.java +++ b/src/main/java/com/michelet/inventory/application/StockLockFacade.java @@ -1,11 +1,11 @@ package com.michelet.inventory.application; import com.michelet.common.exception.BusinessException; +import com.michelet.inventory.application.dto.ReserveStockRequest; +import com.michelet.inventory.application.dto.RestoreStockRequest; import com.michelet.inventory.domain.exception.InventoryErrorCode; import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage.OrderItemDto; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; import java.util.ArrayList; import java.util.List; import java.util.UUID; diff --git a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java b/src/main/java/com/michelet/inventory/application/dto/ReserveStockRequest.java similarity index 90% rename from src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java rename to src/main/java/com/michelet/inventory/application/dto/ReserveStockRequest.java index 8b337fc..969b629 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java +++ b/src/main/java/com/michelet/inventory/application/dto/ReserveStockRequest.java @@ -1,4 +1,4 @@ -package com.michelet.inventory.presentation.dto; +package com.michelet.inventory.application.dto; import jakarta.validation.constraints.Min; import jakarta.validation.constraints.NotNull; diff --git a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java b/src/main/java/com/michelet/inventory/application/dto/RestoreStockRequest.java similarity index 84% rename from src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java rename to src/main/java/com/michelet/inventory/application/dto/RestoreStockRequest.java index acb0540..1fbd997 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java +++ b/src/main/java/com/michelet/inventory/application/dto/RestoreStockRequest.java @@ -1,4 +1,4 @@ -package com.michelet.inventory.presentation.dto; +package com.michelet.inventory.application.dto; import jakarta.validation.constraints.NotNull; import jakarta.validation.constraints.Positive; 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 5181505..fd54db2 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/DeadLetterConsumer.java @@ -19,10 +19,10 @@ public class DeadLetterConsumer { */ @KafkaListener( topics = { - "${inventory.kafka.topic.restore-request:order.stock-restore.requested}.DLT", + // 개별로 나뉘어 있던 토픽을 지우고, 통합 커맨드 토픽의 DLT 구독 + "${inventory.kafka.topic.inventory-command:inventory.command}.DLT", "${inventory.kafka.topic.restored:stock.restored}.DLT", - "${inventory.kafka.topic.reserved:stock.reserved}.DLT", - "${inventory.kafka.topic.order-created:order.created}.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 ffc59e4..8b71a18 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -3,50 +3,68 @@ import com.michelet.common.exception.BusinessException; import com.michelet.inventory.application.InventoryOutboxHelper; import com.michelet.inventory.application.StockLockFacade; +import com.michelet.inventory.application.dto.RestoreStockRequest; +import com.michelet.inventory.domain.exception.InventoryErrorCode; import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; import com.michelet.inventory.infrastructure.messaging.dto.OrderRejectedEvent; import com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +import org.springframework.kafka.annotation.KafkaHandler; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.stereotype.Component; @Slf4j @Component @RequiredArgsConstructor +// 클래스 레벨에서 단일 토픽(inventory.command)을 구독하도록 통합 +@KafkaListener( + topics = "${inventory.kafka.topic.inventory-command:inventory.command}", + groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" +) public class OrderEventConsumer { private final StockLockFacade stockLockFacade; private final InventoryOutboxHelper outboxHelper; - // 오더 생성 이벤트를 비동기로 받아서 다중 락 처리 - @KafkaListener( - topics = "${inventory.kafka.topic.order-created:order.created}", - groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" - ) + // TypeId 헤더를 보고 스프링 카프카가 이 메서드로 라우팅해 줌 + @KafkaHandler public void consumeOrderCreated(OrderCreatedMessage payload) { - log.info("[Kafka Consumer] 신규 주문 생성 메시지 수신 -> 재고 다중 차감 시도: reservationId={}", payload.reservationId()); + // payload null 체크 (Tombstone 메시지 방어 - NPE 방지) + if (payload == null) { + log.error("[Kafka Consumer] 잘못된 주문 생성 이벤트 수신: payload가 null입니다 (Tombstone 메시지일 가능성)."); + throw new IllegalArgumentException("주문 생성 이벤트 파싱 오류: payload는 null일 수 없습니다."); + } + + // null 체크 이후에 로깅 (In-Order 처리) + log.info("[Kafka Consumer] 신규 주문 생성 메시지 수신 -> 재고 다중 차감 시도 (In-Order 처리): reservationId={}", + payload.reservationId()); try { stockLockFacade.reserveOrderStocksWithLock(payload); } catch (IllegalArgumentException e) { log.error("[Kafka Consumer] 비즈니스/검증 룰 위반 에러 (DLT 직행 대상). reservationId: {}", payload.reservationId(), e); throw e; } catch (BusinessException e) { - // 품절, 한도 초과 등 비즈니스 예외 시 오더 서비스에 거절 이벤트 전송 + // 락 획득 실패(일시적 경합)인 경우, 주문 거절을 하지 않고 RuntimeException을 던져 카프카 재시도를 유도 + if (InventoryErrorCode.CONCURRENCY_ERROR.name().equals(e.getErrorCode())) { + log.error("[Kafka Consumer] 락 경합으로 인한 일시적 실패 (재시도 대상). reservationId={}", payload.reservationId(), e); + throw new RuntimeException("재고 락 경합으로 인한 주문 처리 지연", e); + } + + // 품절, 한도 초과 등 명확한 비즈니스 예외 시에만 오더 서비스에 거절 이벤트 전송 log.warn("[Kafka Consumer] 비즈니스 로직에 의한 재고 차감 실패 (거절 이벤트 정상 발행). 사유: {}", e.getMessage()); - outboxHelper.append("ORDER", payload.reservationId().toString(), "ORDER_REJECTED", + // 롤백된 트랜잭션 밖에서 아웃박스를 안전하게 저장하기 위해 appendIndependent 사용! + outboxHelper.appendIndependent("ORDER", payload.reservationId().toString(), "ORDER_REJECTED", new OrderRejectedEvent(payload.reservationId(), e.getMessage())); + } catch (Exception e) { log.error("[Kafka Consumer] 재고 차감 중 일시적 에러 발생 (재시도 대상). reservationId: {}", payload.reservationId(), e); throw new RuntimeException("재고 다중 차감 실패", e); } } - @KafkaListener( - topics = "${inventory.kafka.topic.restore-request:order.stock-restore.requested}", - groupId = "${spring.kafka.consumer.group-id:inventory-service-consumer}" - ) + // 복구 메시지가 들어오면 이 메서드로 라우팅해 줌 + @KafkaHandler public void consumeStockRestoreRequest(StockRestoreMessage payload) { if (payload == null) { log.error("[Kafka Consumer] 잘못된 복구 이벤트 수신: payload가 null입니다 (Tombstone 메시지일 가능성)."); @@ -54,7 +72,7 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: payload는 null일 수 없습니다."); } - log.info("[Kafka Consumer] 재고 복구 이벤트 수신: eventId={}, optionId={}, quantity={}", payload.eventId(), + log.info("[Kafka Consumer] 재고 복구 이벤트 수신 (In-Order 처리): eventId={}, optionId={}, quantity={}", payload.eventId(), payload.optionId(), payload.quantity()); try { @@ -77,4 +95,10 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { throw new RuntimeException("재고 복구 컨슈머 처리 실패", e); } } + + // 만약 예상치 못한 다른 객체가 inventory.command로 들어올 경우, 서버가 터지지 않고 로그만 남기고 무시하도록 방어 + @KafkaHandler(isDefault = true) + public void unknown(Object object) { + log.warn("[Kafka Consumer] 알 수 없는 타입의 커맨드가 inventory.command 토픽으로 들어왔습니다: {}", object); + } } diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java index 715a9ac..f7ea7b2 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/OrderCreatedMessage.java @@ -18,6 +18,13 @@ public record OrderCreatedMessage( if (items == null || items.isEmpty()) { throw new IllegalArgumentException("메시지 파싱 오류: 주문 항목은 1개 이상이어야 합니다."); } + // 리스트 내부 null 원소 검증 + if (items.stream().anyMatch(java.util.Objects::isNull)) { + throw new IllegalArgumentException("메시지 파싱 오류: 주문 항목에 null 값이 포함될 수 없습니다."); + } + + // 외부 조작 방지를 위한 불변 리스트 복사 + items = List.copyOf(items); } public record OrderItemDto(UUID optionId, Integer quantity) { diff --git a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java b/src/main/java/com/michelet/inventory/presentation/InternalStockController.java deleted file mode 100644 index 030b85d..0000000 --- a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java +++ /dev/null @@ -1,41 +0,0 @@ -package com.michelet.inventory.presentation; - -import com.michelet.common.response.ApiResponse; -import com.michelet.inventory.application.StockLockFacade; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; -import jakarta.validation.Valid; -import lombok.RequiredArgsConstructor; -import org.springframework.http.ResponseEntity; -import org.springframework.web.bind.annotation.PostMapping; -import org.springframework.web.bind.annotation.RequestBody; -import org.springframework.web.bind.annotation.RequestMapping; -import org.springframework.web.bind.annotation.RestController; - -@RestController -@RequestMapping("/internal/stocks") -@RequiredArgsConstructor -public class InternalStockController { - - private final StockLockFacade stockLockFacade; - - @PostMapping("/reserve") - public ResponseEntity> reserveStock( - @RequestBody @Valid ReserveStockRequest request - ) { - stockLockFacade.reserveStockWithLock(request); - return ResponseEntity.ok(ApiResponse.ok(null)); - } - - /** - * Kafka 비동기 통신(order.stock-restore.requested)으로 대체됨. 향후 삭제 예정 - */ - @Deprecated(since = "1.0", forRemoval = true) - @PostMapping("/restore") - public ResponseEntity> restoreStock( - @RequestBody @Valid RestoreStockRequest request - ) { - stockLockFacade.restoreStockWithLock(request); - return ResponseEntity.ok(ApiResponse.ok(null)); - } -} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index ca721d4..ef071bb 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.inventory.infrastructure.messaging.dto,com.michelet.order.application.dto" + spring.json.trusted.packages: "com.michelet.inventory.infrastructure.messaging.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,com.michelet.order.application.dto.OrderCreatedEventPayload:com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage" @@ -47,12 +47,12 @@ inventory: topic: product-created: "product.created" reserved: "stock.reserved" - restore-request: "order.stock-restore.requested" # 오더의 '명령'을 수신할 토픽 restored: "stock.restored" # 카탈로그로 '결과'를 발송할 토픽 status-changed: "product.status-changed" daily-reset: "stock.daily-reset" product-updated: "product.updated" - order-created: "order.created" + # 분리되어 있던 토픽을 지우고, 단일 통합 커맨드 토픽 추가 + inventory-command: "inventory.command" order-approved: "order.approved" order-rejected: "order.rejected" consumer: diff --git a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java index b9daac7..1e31839 100644 --- a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java +++ b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java @@ -9,6 +9,8 @@ import static org.mockito.Mockito.verify; import static org.mockito.Mockito.verifyNoInteractions; +import com.michelet.inventory.application.dto.ReserveStockRequest; +import com.michelet.inventory.application.dto.RestoreStockRequest; import com.michelet.inventory.application.dto.StockReservedEvent; import com.michelet.inventory.application.dto.StockRestoredEvent; import com.michelet.inventory.domain.exception.MaxLimitExceededException; @@ -20,8 +22,6 @@ import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; -import com.michelet.inventory.presentation.dto.RestoreStockRequest; import java.util.Optional; import java.util.UUID; import org.junit.jupiter.api.DisplayName; diff --git a/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java b/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java index d6b4ee4..ba7815d 100644 --- a/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java +++ b/src/test/java/com/michelet/inventory/application/StockConcurrencyIntegrationTest.java @@ -5,10 +5,10 @@ import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.BDDMockito.given; +import com.michelet.inventory.application.dto.ReserveStockRequest; import com.michelet.inventory.domain.model.Stock; import com.michelet.inventory.domain.repository.StockRepository; import com.michelet.inventory.infrastructure.repository.JpaStockRepository; -import com.michelet.inventory.presentation.dto.ReserveStockRequest; import java.util.UUID; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; diff --git a/src/test/resources/application-test.yml b/src/test/resources/application-test.yml index c334cb0..1376627 100644 --- a/src/test/resources/application-test.yml +++ b/src/test/resources/application-test.yml @@ -51,4 +51,4 @@ inventory: restored: "stock.restored" status-changed: "product.status-changed" daily-reset: "stock.daily-reset" - + inventory-command: "inventory.command" From 635776c6903c16e26e0240c5ba5cb033feb4ad48 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Sun, 17 May 2026 14:32:04 +0900 Subject: [PATCH 3/4] =?UTF-8?q?fix:=20=EB=8B=A4=EC=A4=91=EC=9E=AC=EA=B3=A0?= =?UTF-8?q?=20=EC=B9=B4=ED=94=84=EC=B9=B4=20=EB=88=84=EB=9D=BD=EC=BD=94?= =?UTF-8?q?=EB=93=9C=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../application/StockCommandService.java | 30 ++++++++++++++++--- .../messaging/OrderEventConsumer.java | 8 +++-- 2 files changed, 32 insertions(+), 6 deletions(-) diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index 22bb134..4d32f07 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -19,6 +19,9 @@ import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; import java.util.ArrayList; import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Collectors; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.stereotype.Service; @@ -113,16 +116,35 @@ public void processOrderCreation(OrderCreatedMessage msg) { return; } + Map groupedItems = msg.items().stream() + .collect(Collectors.groupingBy( + OrderCreatedMessage.OrderItemDto::optionId, + Collectors.summingInt(OrderCreatedMessage.OrderItemDto::quantity) + )); + List modifiedStocks = new ArrayList<>(); - // 모든 아이템 재고 차감 시도 (실패 시 BusinessException 발생하여 Facade -> Consumer 로 롤백됨) - for (var item : msg.items()) { - Stock stock = stockRepository.findById(item.optionId()) + // 합산된 수량으로 재고 차감 시도 + for (Map.Entry entry : groupedItems.entrySet()) { + UUID optionId = entry.getKey(); + int quantity = entry.getValue(); + + Stock stock = stockRepository.findById(optionId) .orElseThrow(StockNotFoundException::new); - stock.reserve(item.quantity()); + stock.reserve(quantity); modifiedStocks.add(stock); + // 카탈로그 동기화용 이벤트 발송 + // 파티션 키는 optionId - 동일 옵션에 대한 차감 순서 FIFO + StockReservedEvent reservedEvent = new StockReservedEvent( + stock.getOptionId(), + stock.getTotalQuantity(), + stock.getCurrentDailyStock() + ); + outboxHelper.append("STOCK", stock.getOptionId().toString(), "STOCK_RESERVED", reservedEvent); + log.info("[Inventory Saga] 다중 주문 옵션 차감 이벤트 적재 완료: optionId={}, 수량={}", stock.getOptionId(), quantity); + // 품절 처리 이벤트 발송 준비 if (stock.getTotalQuantity() == 0) { ProductOption option = productOptionRepository.findById(stock.getOptionId()).orElseThrow(); 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 8b71a18..c4899f0 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -96,9 +96,13 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { } } - // 만약 예상치 못한 다른 객체가 inventory.command로 들어올 경우, 서버가 터지지 않고 로그만 남기고 무시하도록 방어 + // 알 수 없는 타입의 객체가 inventory.command로 들어올 경우, 조용히 넘기지 않고 예외를 발생시켜 DLT로 격리 @KafkaHandler(isDefault = true) public void unknown(Object object) { - log.warn("[Kafka Consumer] 알 수 없는 타입의 커맨드가 inventory.command 토픽으로 들어왔습니다: {}", object); + log.error("[Kafka Consumer] 지원하지 않는 커맨드 타입 수신. object={}", object); + throw new IllegalArgumentException( + "지원하지 않는 inventory.command payload 타입: " + + (object == null ? "null" : object.getClass().getName()) + ); } } From 0b508067fabdca2759c8436dc59033f3ce00cd5d Mon Sep 17 00:00:00 2001 From: ji-circle Date: Sun, 17 May 2026 16:20:13 +0900 Subject: [PATCH 4/4] =?UTF-8?q?fix:=20=EB=8F=99=EC=8B=9C=EC=84=B1=20?= =?UTF-8?q?=EC=98=A4=EB=A5=98=EC=97=90=20=EB=8C=80=ED=95=9C=20=EC=9E=AC?= =?UTF-8?q?=EC=8B=9C=EB=8F=84=20=EB=A1=9C=EC=A7=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../messaging/OrderEventConsumer.java | 13 +++++++------ 1 file changed, 7 insertions(+), 6 deletions(-) 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 c4899f0..eff66b6 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -4,7 +4,7 @@ import com.michelet.inventory.application.InventoryOutboxHelper; import com.michelet.inventory.application.StockLockFacade; import com.michelet.inventory.application.dto.RestoreStockRequest; -import com.michelet.inventory.domain.exception.InventoryErrorCode; +import com.michelet.inventory.domain.exception.ConcurrencyFailureException; import com.michelet.inventory.infrastructure.messaging.dto.OrderCreatedMessage; import com.michelet.inventory.infrastructure.messaging.dto.OrderRejectedEvent; import com.michelet.inventory.infrastructure.messaging.dto.StockRestoreMessage; @@ -44,13 +44,14 @@ public void consumeOrderCreated(OrderCreatedMessage payload) { } catch (IllegalArgumentException e) { log.error("[Kafka Consumer] 비즈니스/검증 룰 위반 에러 (DLT 직행 대상). reservationId: {}", payload.reservationId(), e); throw e; - } catch (BusinessException e) { + + } catch (ConcurrencyFailureException e) { + // 여기서 동시성 에러를 가로채서 재시도 처리함! // 락 획득 실패(일시적 경합)인 경우, 주문 거절을 하지 않고 RuntimeException을 던져 카프카 재시도를 유도 - if (InventoryErrorCode.CONCURRENCY_ERROR.name().equals(e.getErrorCode())) { - log.error("[Kafka Consumer] 락 경합으로 인한 일시적 실패 (재시도 대상). reservationId={}", payload.reservationId(), e); - throw new RuntimeException("재고 락 경합으로 인한 주문 처리 지연", e); - } + log.error("[Kafka Consumer] 락 경합으로 인한 일시적 실패 (재시도 대상). reservationId={}", payload.reservationId(), e); + throw new RuntimeException("재고 락 경합으로 인한 주문 처리 지연", e); + } catch (BusinessException e) { // 품절, 한도 초과 등 명확한 비즈니스 예외 시에만 오더 서비스에 거절 이벤트 전송 log.warn("[Kafka Consumer] 비즈니스 로직에 의한 재고 차감 실패 (거절 이벤트 정상 발행). 사유: {}", e.getMessage()); // 롤백된 트랜잭션 밖에서 아웃박스를 안전하게 저장하기 위해 appendIndependent 사용!