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
Expand Up @@ -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);
Expand All @@ -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) {
Expand Down Expand Up @@ -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는 필수입니다.");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand All @@ -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() {
Expand Down Expand Up @@ -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);
Expand All @@ -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";
};
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -13,8 +15,13 @@
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 com.michelet.inventory.infrastructure.messaging.dto.OrderApprovedEvent;
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;
Expand Down Expand Up @@ -100,4 +107,64 @@ 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;
}

Map<UUID, Integer> groupedItems = msg.items().stream()
.collect(Collectors.groupingBy(
OrderCreatedMessage.OrderItemDto::optionId,
Collectors.summingInt(OrderCreatedMessage.OrderItemDto::quantity)
));

List<Stock> modifiedStocks = new ArrayList<>();

// 합산된 수량으로 재고 차감 시도
for (Map.Entry<UUID, Integer> entry : groupedItems.entrySet()) {
UUID optionId = entry.getKey();
int quantity = entry.getValue();

Stock stock = stockRepository.findById(optionId)
.orElseThrow(StockNotFoundException::new);

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();
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());
}
}
Original file line number Diff line number Diff line change
@@ -1,7 +1,14 @@
package com.michelet.inventory.application;

import com.michelet.inventory.presentation.dto.ReserveStockRequest;
import com.michelet.inventory.presentation.dto.RestoreStockRequest;
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 java.util.ArrayList;
import java.util.List;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
Expand Down Expand Up @@ -67,4 +74,43 @@ public void restoreStockWithLock(RestoreStockRequest request) {
}
}
}

// 3. 다중 품목 비동기 재고 선점 (OrderCreatedMessage 수신용)
public void reserveOrderStocksWithLock(
OrderCreatedMessage msg) {
// 데드락 방지를 위해 OptionId를 오름차순으로 정렬하여 락을 순차적으로 획득
List<UUID> sortedOptionIds = msg.items().stream()
.map(OrderItemDto::optionId)
.sorted()
.toList();

List<RLock> 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();
}
}
}
}
}
Original file line number Diff line number Diff line change
@@ -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;
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,8 @@ 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"
},
Expand Down
Loading