From 7dae666ac04183126e281cbde6b4fa24a9c88cb4 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Thu, 14 May 2026 17:28:47 +0900 Subject: [PATCH 1/3] =?UTF-8?q?feat:=20=EB=A9=B1=EB=93=B1=EC=84=B1=20?= =?UTF-8?q?=EA=B2=80=EC=82=AC=EB=A1=9C=20=EC=A4=91=EB=B3=B5=20=EB=B0=A9?= =?UTF-8?q?=EC=96=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../application/StockCommandService.java | 18 +++++++-- .../domain/model/ProcessedEvent.java | 30 ++++++++++++++ .../repository/ProcessedEventRepository.java | 10 +++++ .../messaging/OrderEventConsumer.java | 11 ++++-- .../messaging/dto/StockRestoreMessage.java | 4 ++ .../JpaProcessedEventRepository.java | 8 ++++ .../ProcessedEventRepositoryImpl.java | 24 ++++++++++++ .../presentation/InternalStockController.java | 4 ++ .../presentation/dto/RestoreStockRequest.java | 2 +- .../application/StockCommandServiceTest.java | 39 +++++++++++++++++-- 10 files changed, 138 insertions(+), 12 deletions(-) create mode 100644 src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java create mode 100644 src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java create mode 100644 src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java create mode 100644 src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index 76ef241..26607ac 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -4,10 +4,12 @@ import com.michelet.inventory.application.dto.StockReservedEvent; import com.michelet.inventory.application.dto.StockRestoredEvent; import com.michelet.inventory.domain.exception.StockNotFoundException; +import com.michelet.inventory.domain.model.ProcessedEvent; import com.michelet.inventory.domain.model.Product; import com.michelet.inventory.domain.model.ProductOption; import com.michelet.inventory.domain.model.ProductStatus; import com.michelet.inventory.domain.model.Stock; +import com.michelet.inventory.domain.repository.ProcessedEventRepository; import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; @@ -26,6 +28,7 @@ public class StockCommandService { private final StockRepository stockRepository; private final ProductOptionRepository productOptionRepository; private final ProductRepository productRepository; + private final ProcessedEventRepository processedEventRepository; // OutboxHelper 주입 (KafkaTemplate 대체) private final InventoryOutboxHelper outboxHelper; @@ -71,16 +74,25 @@ public void reserveStock(ReserveStockRequest request) { @Transactional public void restoreStock(RestoreStockRequest request) { + // 1. 멱등성 검증 (Redisson Lock 내부이므로 동시성 중복 방어됨) + if (processedEventRepository.existsById(request.eventId())) { + log.warn("[Idempotency] 이미 처리된 카프카 메시지입니다. 중복 복구를 방지하고 반환합니다. eventId={}", request.eventId()); + return; // 중복 메시지는 로직을 건너뛰고 정상 완료 처리 + } + Stock stock = stockRepository.findById(request.optionId()) .orElseThrow(StockNotFoundException::new); - // 1. 도메인 로직: 재고 복구 + // 2. 도메인 로직: 재고 복구 stock.restore(request.quantity()); - // 2. DB 업데이트 (더티 체킹 후 flush) + // 3. DB 업데이트 stockRepository.save(stock); - // 3. 재고 복구 이벤트를 Outbox에 적재 + // 4. 처리 완료 기록 (멱등키 저장) + processedEventRepository.save(new ProcessedEvent(request.eventId())); + + // 5. 재고 복구 결과를 Outbox에 적재하여 카탈로그로 발송 StockRestoredEvent event = new StockRestoredEvent( stock.getOptionId(), stock.getTotalQuantity(), diff --git a/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java new file mode 100644 index 0000000..d7ae2bf --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java @@ -0,0 +1,30 @@ +package com.michelet.inventory.domain.model; + +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import java.time.LocalDateTime; +import java.util.UUID; +import lombok.AccessLevel; +import lombok.Getter; +import lombok.NoArgsConstructor; + +@Entity +@Table(name = "p_processed_events") +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +public class ProcessedEvent { + + @Id + @Column(name = "event_id", updatable = false, nullable = false) + private UUID eventId; + + @Column(nullable = false) + private LocalDateTime processedAt; + + public ProcessedEvent(UUID eventId) { + this.eventId = eventId; + this.processedAt = LocalDateTime.now(); + } +} diff --git a/src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java b/src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java new file mode 100644 index 0000000..0805679 --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/repository/ProcessedEventRepository.java @@ -0,0 +1,10 @@ +package com.michelet.inventory.domain.repository; + +import com.michelet.inventory.domain.model.ProcessedEvent; +import java.util.UUID; + +public interface ProcessedEventRepository { + boolean existsById(UUID eventId); + + ProcessedEvent save(ProcessedEvent event); +} 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 fe3bef8..9f9563f 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/OrderEventConsumer.java @@ -26,15 +26,18 @@ public void consumeStockRestoreRequest(StockRestoreMessage payload) { throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: payload는 null일 수 없습니다."); } - log.info("[Kafka Consumer] 재고 복구 이벤트 수신: optionId={}, quantity={}", payload.optionId(), payload.quantity()); + log.info("[Kafka Consumer] 재고 복구 이벤트 수신: eventId={}, optionId={}, quantity={}", payload.eventId(), + payload.optionId(), payload.quantity()); try { // 파사드를 통해 분산 락을 걸고 안전하게 재고 복구 로직 실행 - // TODO(`#26`): reservationId를 Kafka 페이로드에 포함시켜 멱등성 키로 활용 - RestoreStockRequest request = new RestoreStockRequest(payload.optionId(), payload.quantity(), null); + // 전달받은 eventId를 파사드 요청에 포함 + RestoreStockRequest request = new RestoreStockRequest(payload.optionId(), payload.quantity(), + payload.eventId()); stockLockFacade.restoreStockWithLock(request); - log.info("[Kafka Consumer] 재고 복구 완료! optionId: {}, quantity: {}", payload.optionId(), payload.quantity()); + log.info("[Kafka Consumer] 재고 복구 로직 처리 완료! eventId: {}, optionId: {}, quantity: {}", payload.eventId(), + payload.optionId(), payload.quantity()); } catch (IllegalArgumentException e) { // 데이터 검증 실패나 비즈니스 룰 위반 시 -> 즉각 DLT로 직행하도록 원본 예외 던짐 diff --git a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java index 43d3767..25f6b9b 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java +++ b/src/main/java/com/michelet/inventory/infrastructure/messaging/dto/StockRestoreMessage.java @@ -4,10 +4,14 @@ // 오더 서비스가 발행한 이벤트를 읽어들이기 위한 수신 전용 DTO public record StockRestoreMessage( + UUID eventId, // 멱등성 검증용 고유 이벤트 ID UUID optionId, Integer quantity ) { public StockRestoreMessage { + if (eventId == null) { + throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: eventId는 null일 수 없습니다."); + } if (optionId == null) { throw new IllegalArgumentException("재고 복구 이벤트 파싱 오류: optionId는 null일 수 없습니다."); } diff --git a/src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java b/src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java new file mode 100644 index 0000000..cf3a1c7 --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/repository/JpaProcessedEventRepository.java @@ -0,0 +1,8 @@ +package com.michelet.inventory.infrastructure.repository; + +import com.michelet.inventory.domain.model.ProcessedEvent; +import java.util.UUID; +import org.springframework.data.jpa.repository.JpaRepository; + +public interface JpaProcessedEventRepository extends JpaRepository { +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java b/src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java new file mode 100644 index 0000000..e56d61f --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/repository/ProcessedEventRepositoryImpl.java @@ -0,0 +1,24 @@ +package com.michelet.inventory.infrastructure.repository; + +import com.michelet.inventory.domain.model.ProcessedEvent; +import com.michelet.inventory.domain.repository.ProcessedEventRepository; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Repository; + +@Repository +@RequiredArgsConstructor +public class ProcessedEventRepositoryImpl implements ProcessedEventRepository { + + private final JpaProcessedEventRepository jpaRepository; + + @Override + public boolean existsById(UUID eventId) { + return jpaRepository.existsById(eventId); + } + + @Override + public ProcessedEvent save(ProcessedEvent event) { + return jpaRepository.save(event); + } +} diff --git a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java b/src/main/java/com/michelet/inventory/presentation/InternalStockController.java index bdc5c67..030b85d 100644 --- a/src/main/java/com/michelet/inventory/presentation/InternalStockController.java +++ b/src/main/java/com/michelet/inventory/presentation/InternalStockController.java @@ -27,6 +27,10 @@ public ResponseEntity> reserveStock( 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 diff --git a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java b/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java index e3e488e..acb0540 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java +++ b/src/main/java/com/michelet/inventory/presentation/dto/RestoreStockRequest.java @@ -7,6 +7,6 @@ public record RestoreStockRequest( @NotNull UUID optionId, @Positive int quantity, - UUID reservationId // 멱등키 — TODO 추후 중복 복구 방지 로직에서 활용 예정 (그땐 @NotNull) + @NotNull UUID eventId // 멱등키 필드 ) { } diff --git a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java index 67532c9..b9daac7 100644 --- a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java +++ b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java @@ -16,6 +16,7 @@ import com.michelet.inventory.domain.exception.SoldOutException; import com.michelet.inventory.domain.exception.StockNotFoundException; import com.michelet.inventory.domain.model.Stock; +import com.michelet.inventory.domain.repository.ProcessedEventRepository; import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; @@ -51,6 +52,10 @@ class StockCommandServiceTest { @Mock private InventoryOutboxHelper outboxHelper; + // 멱등성 검증을 위한 Mock 객체 추가 + @Mock + private ProcessedEventRepository processedEventRepository; + // 카프카로 전송된 이벤트를 낚아채서 내부 값을 검증하기 위한 Captor @Captor private ArgumentCaptor eventCaptor; @@ -139,18 +144,22 @@ void reserveStock_Fail_MaxLimitExceeded() { } @Test - @DisplayName("성공: 재고 복구 경로 정상 동작 및 Outbox 적재 테스트") + @DisplayName("성공: 재고 복구 경로 정상 동작 및 멱등키, Outbox 적재 테스트") void restoreStock_Success() { UUID optionId = UUID.randomUUID(); + UUID eventId = UUID.randomUUID(); // 멱등키(메시지 ID) + // 1. 초기 재고 100개, 일일 재고 50개 생성 Stock stock = Stock.create(optionId, 100, 50, 10); // 2. 소비된 상태를 시뮬레이션하기 위해 미리 2개를 차감 (total: 98, daily: 48) stock.reserve(2); - // 3. 다시 2개를 복구해 달라는 요청 - RestoreStockRequest request = new RestoreStockRequest(optionId, 2, null); + // 3. 다시 2개를 복구해 달라는 요청 - eventId를 포함하여 요청 객체 생성 + RestoreStockRequest request = new RestoreStockRequest(optionId, 2, eventId); + // 처음 들어온 메시지이므로 existsById는 false를 반환하도록 Mocking + given(processedEventRepository.existsById(eventId)).willReturn(false); given(stockRepository.findById(optionId)).willReturn(Optional.of(stock)); // when @@ -169,12 +178,34 @@ void restoreStock_Success() { assertThat(event.currentDailyStock()).isEqualTo(50); } + // 멱등성 보장 핵심 테스트 + @Test + @DisplayName("성공(멱등성): 이미 처리된 이벤트인 경우 이중 복구를 수행하지 않고 무시한다.") + void restoreStock_Idempotency() { + UUID optionId = UUID.randomUUID(); + UUID eventId = UUID.randomUUID(); + RestoreStockRequest request = new RestoreStockRequest(optionId, 2, eventId); + + // DB에 이미 eventId가 저장되어 있다고(처리되었다고) Mocking + given(processedEventRepository.existsById(eventId)).willReturn(true); + + // when + stockCommandService.restoreStock(request); + + // then: DB 조회나 아웃박스 저장이 단 한 번도 호출되지 않아야 함! (이중 복구 방지 증명) + verifyNoInteractions(stockRepository); + verifyNoInteractions(outboxHelper); + } + @Test @DisplayName("실패: 재고 복구 시 존재하지 않는 옵션 예외") void restoreStock_Fail_NotFound() { UUID optionId = UUID.randomUUID(); - RestoreStockRequest request = new RestoreStockRequest(optionId, 5, null); + UUID eventId = UUID.randomUUID(); + RestoreStockRequest request = new RestoreStockRequest(optionId, 5, eventId); + // 멱등성 검사 통과 설정 + given(processedEventRepository.existsById(eventId)).willReturn(false); given(stockRepository.findById(optionId)).willReturn(Optional.empty()); assertThatThrownBy(() -> stockCommandService.restoreStock(request)) From 81df64c4495ad4cf72fb55264298b10e27863840 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Thu, 14 May 2026 22:07:59 +0900 Subject: [PATCH 2/3] feat: LocalDateTime -> Instant --- .../michelet/inventory/domain/model/ProcessedEvent.java | 7 ++++--- .../inventory/presentation/dto/ReserveStockRequest.java | 2 +- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java index d7ae2bf..fdc117f 100644 --- a/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java +++ b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java @@ -4,7 +4,8 @@ import jakarta.persistence.Entity; import jakarta.persistence.Id; import jakarta.persistence.Table; -import java.time.LocalDateTime; +import java.time.Instant; +import java.time.temporal.ChronoUnit; import java.util.UUID; import lombok.AccessLevel; import lombok.Getter; @@ -21,10 +22,10 @@ public class ProcessedEvent { private UUID eventId; @Column(nullable = false) - private LocalDateTime processedAt; + private Instant processedAt; public ProcessedEvent(UUID eventId) { this.eventId = eventId; - this.processedAt = LocalDateTime.now(); + this.processedAt = Instant.now().truncatedTo(ChronoUnit.MILLIS); } } diff --git a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java b/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java index 1d01b4d..91cee31 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java +++ b/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java @@ -7,6 +7,6 @@ public record ReserveStockRequest( @NotNull UUID optionId, @NotNull @Min(1) Integer quantity, - UUID reservationId // 멱등키 — TODO 추후 중복 차감 방지 로직에서 활용 예정 (그땐 @NotNull) + UUID reservationId // 멱등키 — TODO 추후 동기 통신(Feign)에서의 재시도로 인한 중복 차감을 막기 위한 로직에서 활용 예정 (그땐 @NotNull) ) { } From 6b9f3b009517b8f67712103f0acadef70f06f2b3 Mon Sep 17 00:00:00 2001 From: ji-circle Date: Thu, 14 May 2026 23:17:12 +0900 Subject: [PATCH 3/3] =?UTF-8?q?fix:=20ProcessedEvent=20,=20ReserveStockReq?= =?UTF-8?q?uest=20=EC=A1=B0=EA=B1=B4=20=EC=88=98=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../inventory/domain/model/ProcessedEvent.java | 3 ++- .../presentation/dto/ReserveStockRequest.java | 10 +++++++--- 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java index fdc117f..48aea8d 100644 --- a/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java +++ b/src/main/java/com/michelet/inventory/domain/model/ProcessedEvent.java @@ -6,6 +6,7 @@ import jakarta.persistence.Table; import java.time.Instant; import java.time.temporal.ChronoUnit; +import java.util.Objects; import java.util.UUID; import lombok.AccessLevel; import lombok.Getter; @@ -25,7 +26,7 @@ public class ProcessedEvent { private Instant processedAt; public ProcessedEvent(UUID eventId) { - this.eventId = eventId; + this.eventId = Objects.requireNonNull(eventId, "eventId는 null일 수 없습니다"); this.processedAt = Instant.now().truncatedTo(ChronoUnit.MILLIS); } } diff --git a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java b/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java index 91cee31..8b337fc 100644 --- a/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java +++ b/src/main/java/com/michelet/inventory/presentation/dto/ReserveStockRequest.java @@ -5,8 +5,12 @@ import java.util.UUID; public record ReserveStockRequest( - @NotNull UUID optionId, - @NotNull @Min(1) Integer quantity, - UUID reservationId // 멱등키 — TODO 추후 동기 통신(Feign)에서의 재시도로 인한 중복 차감을 막기 위한 로직에서 활용 예정 (그땐 @NotNull) + @NotNull(message = "optionId는 필수입니다.") + UUID optionId, + @NotNull(message = "quantity는 필수입니다.") + @Min(value = 1, message = "quantity는 1 이상이어야 합니다.") + Integer quantity, + @NotNull(message = "멱등키(reservationId)는 필수입니다.") + UUID reservationId ) { }