diff --git a/api-test/inventory.http b/api-test/inventory.http index 9b17997..0da7781 100644 --- a/api-test/inventory.http +++ b/api-test/inventory.http @@ -441,3 +441,27 @@ X-Internal-Token: {{inv_internal_token}} client.log("모든 재고 소진 완료! (Product 상태 SOLDOUT 자동 전이 및 이벤트 발행됨)"); }); %} + +### [Internal Token] 내부 통신용 토큰 갱신 +GET http://{{host}}:{{orderPort}}/test/internal-token?audience=inventory-service + +> {% + client.global.set("inv_internal_token", response.body); + client.log("전역변수 저장: inv_internal_token"); +%} + +### 14. [Internal] Inventory Outbox 스케줄러 강제 트리거 +# DB에 적재된 INIT 상태의 Outbox 이벤트들을 읽어 카프카로 일괄 발행 +POST http://{{host}}:{{inventoryPort}}/internal/scheduler/trigger-outbox +Content-Type: application/json +X-Internal-Token: {{inv_internal_token}} + +> {% + client.test("Inventory Outbox 스케줄러 실행 확인", function () { + client.assert(response.status === 200, "스케줄러 호출 실패"); + client.log("======================================="); + client.log("성공! Inventory 서버의 Outbox 스케줄러가 수동 가동되었습니다."); + client.log("서버 콘솔에 [Inventory Outbox Scheduler] 이벤트 발행 성공! 로그가 찍히는지 확인하기"); + client.log("======================================="); + }); +%} diff --git a/src/main/java/com/michelet/inventory/application/ExhibitionSchedulerService.java b/src/main/java/com/michelet/inventory/application/ExhibitionSchedulerService.java index eb20d1d..aedc3cb 100644 --- a/src/main/java/com/michelet/inventory/application/ExhibitionSchedulerService.java +++ b/src/main/java/com/michelet/inventory/application/ExhibitionSchedulerService.java @@ -9,17 +9,13 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Lazy; import org.springframework.data.domain.PageRequest; import org.springframework.data.domain.Slice; -import org.springframework.kafka.core.KafkaTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Propagation; import org.springframework.transaction.annotation.Transactional; -import org.springframework.transaction.support.TransactionSynchronization; -import org.springframework.transaction.support.TransactionSynchronizationManager; @Slf4j @Service @@ -27,7 +23,7 @@ public class ExhibitionSchedulerService { private final ProductRepository productRepository; - private final KafkaTemplate kafkaTemplate; + private final InventoryOutboxHelper outboxHelper; // Spring AOP의 프록시를 타기 위한 자기 자신 주입 @Lazy @@ -37,11 +33,7 @@ public class ExhibitionSchedulerService { private static final int CHUNK_SIZE = 100; private static final ZoneId SEOUL_ZONE = ZoneId.of("Asia/Seoul"); - @Value("${inventory.kafka.topic.status-changed:product.status-changed}") - private String topicStatusChanged; - - // 외부 진입점의 @Transactional을 제거하여 영속성 컨텍스트 비대화를 막음 - @Scheduled(cron = "0 0 * * * *", zone = "Asia/Seoul") // 매 정시(0분 0초)마다 실행 + @Scheduled(cron = "0 0 * * * *", zone = "Asia/Seoul") // 매 정시(0분 0초)마다 실행 public void updateExhibitionStatus() { log.info("정시 전시 상태 변경 스케줄러 시작..."); LocalDateTime now = LocalDateTime.now(SEOUL_ZONE); @@ -78,7 +70,7 @@ public int processOpeningChunk(LocalDateTime now) { for (Product product : slice.getContent()) { product.changeStatus(ProductStatus.ACTIVE); - publishStatusChangeEventAfterCommit(product); + publishStatusChangeEvent(product); } return slice.getNumberOfElements(); @@ -90,34 +82,14 @@ public int processClosingChunk(LocalDateTime now) { for (Product product : slice.getContent()) { product.changeStatus(ProductStatus.EXPIRED); - publishStatusChangeEventAfterCommit(product); + publishStatusChangeEvent(product); } return slice.getNumberOfElements(); } - // DB 커밋 성공 시에만 카프카로 이벤트 발행 보장 - private void publishStatusChangeEventAfterCommit(Product product) { + private void publishStatusChangeEvent(Product product) { ProductStatusChangedEvent event = new ProductStatusChangedEvent(product.getId(), product.getStatus().name()); - - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - sendKafkaEvent(product, event); - } - }); - } else { - sendKafkaEvent(product, event); - } - } - - private void sendKafkaEvent(Product product, ProductStatusChangedEvent event) { - kafkaTemplate.send(topicStatusChanged, product.getId().toString(), event) - .whenComplete((result, ex) -> { - if (ex != null) { - log.error("상품 상태 변경 카프카 이벤트 발행 실패: productId={}", product.getId(), ex); - } - }); + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_STATUS_CHANGED", event); } } diff --git a/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java b/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java new file mode 100644 index 0000000..1f0e3bd --- /dev/null +++ b/src/main/java/com/michelet/inventory/application/InventoryOutboxHelper.java @@ -0,0 +1,98 @@ +package com.michelet.inventory.application; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.michelet.inventory.domain.model.InventoryOutbox; +import com.michelet.inventory.domain.model.OutboxStatus; +import com.michelet.inventory.domain.repository.InventoryOutboxRepository; +import jakarta.annotation.PostConstruct; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +@Slf4j +@Component +@RequiredArgsConstructor +public class InventoryOutboxHelper { + + private final InventoryOutboxRepository outboxRepository; + private final ObjectMapper objectMapper; + + @Value("${inventory.outbox.max-retries:3}") + private int maxRetries; + + // 시작 시점에 maxRetries 값 검증 + @PostConstruct + public void validateMaxRetries() { + if (maxRetries < 1) { + throw new IllegalStateException("inventory.outbox.max-retries 설정 오류: 반드시 1 이상이어야 합니다."); + } + } + + // 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는 필수입니다."); + } + + 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) { + outboxRepository.findById(outboxId).ifPresentOrElse( + outbox -> { + if (outbox.getStatus() != OutboxStatus.INIT) { + return; + } + outbox.markAsPublished(); + outboxRepository.save(outbox); + }, + () -> log.warn("[Inventory Outbox] 상태 변경 대상이 없습니다. id={}", outboxId) + ); + } + + // 비동기 실패 시 재시도 횟수 및 상태 관리 로직 + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void handleFailure(UUID outboxId) { + outboxRepository.findById(outboxId).ifPresent(outbox -> { + if (outbox.getStatus() != OutboxStatus.INIT) { + return; + } + outbox.incrementRetryCount(); + if (outbox.getRetryCount() >= maxRetries) { // 3번 이상 실패 시 영구 실패 처리 + outbox.markAsFailed(); + log.error("[CRITICAL] Outbox 발행 영구 실패. 수동 확인 요망! id={}", outboxId); + } + outboxRepository.save(outbox); + }); + } +} diff --git a/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java b/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java new file mode 100644 index 0000000..4ec3f37 --- /dev/null +++ b/src/main/java/com/michelet/inventory/application/InventoryOutboxScheduler.java @@ -0,0 +1,141 @@ +package com.michelet.inventory.application; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.michelet.inventory.application.dto.DailyStockResetEvent; +import com.michelet.inventory.application.dto.ProductCreatedEvent; +import com.michelet.inventory.application.dto.ProductStatusChangedEvent; +import com.michelet.inventory.application.dto.ProductUpdatedEvent; +import com.michelet.inventory.application.dto.StockReservedEvent; +import com.michelet.inventory.application.dto.StockRestoredEvent; +import com.michelet.inventory.domain.model.InventoryOutbox; +import com.michelet.inventory.domain.model.OutboxStatus; +import com.michelet.inventory.domain.repository.InventoryOutboxRepository; +import java.util.List; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.kafka.core.KafkaTemplate; +import org.springframework.orm.ObjectOptimisticLockingFailureException; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +@Slf4j +@Component +@RequiredArgsConstructor +public class InventoryOutboxScheduler { + + private final InventoryOutboxRepository outboxRepository; + private final InventoryOutboxHelper outboxHelper; + private final KafkaTemplate kafkaTemplate; + + // JSON 문자열을 객체로 복원하기 위한 매퍼 주입 + private final ObjectMapper objectMapper; + + // 문자열 상수 추출 + private static final String EVENT_PRODUCT_CREATED = "PRODUCT_CREATED"; + private static final String EVENT_PRODUCT_UPDATED = "PRODUCT_UPDATED"; + private static final String EVENT_STATUS_CHANGED = "PRODUCT_STATUS_CHANGED"; + 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"; + + @Value("${inventory.kafka.topic.product-created:product.created}") + private String topicProductCreated; + @Value("${inventory.kafka.topic.product-updated:product.updated}") + private String topicProductUpdated; + @Value("${inventory.kafka.topic.status-changed:product.status-changed}") + private String topicStatusChanged; + @Value("${inventory.kafka.topic.reserved:stock.reserved}") + private String topicStockReserved; + @Value("${inventory.kafka.topic.restored:stock.restored}") + private String topicStockRestored; + @Value("${inventory.kafka.topic.daily-reset:stock.daily-reset}") + private String topicDailyReset; + + @Scheduled(fixedDelay = 5000) + public void processOutboxEvents() { + // 1. OOM 방지 및 순서 보장을 위해 Top N 배치 조회 + List pendingEvents = outboxRepository.findTop50ByStatusOrderByCreatedAtAsc(OutboxStatus.INIT); + if (pendingEvents.isEmpty()) { + return; + } + + log.info("[Inventory Outbox Scheduler] {}개의 미발행 이벤트를 찾아 Kafka 전송을 시도합니다.", pendingEvents.size()); + + for (InventoryOutbox event : pendingEvents) { + try { + String topic = resolveTopic(event.getEventType()); + + // String(JSON)을 다시 원본 Event 객체로 복원 + Object originalEventObject = deserializePayload(event.getEventType(), event.getPayload()); + + // 블로킹(.get) 제거 -> 비동기 발송 콜백(.whenComplete) 적용 + kafkaTemplate.send(topic, event.getAggregateId(), originalEventObject) + .whenComplete((result, ex) -> { + if (ex == null) { + try { + outboxHelper.markAsPublished(event.getId()); + log.info("[Inventory Outbox Scheduler] 이벤트 발행 성공! Outbox ID: {}", event.getId()); + } catch (ObjectOptimisticLockingFailureException oole) { + log.info("[Inventory Outbox Scheduler] 낙관적 락 방어 (동시성 경합). Outbox ID: {}", + event.getId()); + } catch (Exception updateEx) { + log.error("[Inventory Outbox Scheduler] DB 상태 업데이트 실패. Outbox ID: {}", event.getId(), + updateEx); + } + } else { + log.error("[Inventory Outbox Scheduler] 카프카 이벤트 발행 실패. Outbox ID: {}", event.getId(), ex); + safeHandleFailure(event.getId()); + } + }); + + } catch (Exception e) { + // 역직렬화 실패, 토픽 변환 실패 등 무한 에러 유발 시 + log.error("[Inventory Outbox Scheduler] 이벤트 전송 준비 중 예외 발생. Outbox ID: {}", event.getId(), e); + safeHandleFailure(event.getId()); + } + } + } + + // 재시도 횟수 처리 및 상태 변경을 돕는 실패 처리 메서드 + private void safeHandleFailure(UUID eventId) { + try { + outboxHelper.handleFailure(eventId); + } catch (ObjectOptimisticLockingFailureException oole) { + log.info("[Inventory Outbox Scheduler] 실패 마킹 중 낙관적 락 방어. Outbox ID: {}", eventId); + } catch (Exception e) { + log.error("[Inventory Outbox Scheduler] 실패 상태 업데이트 중 예외 발생. Outbox ID: {}", eventId, e); + } + } + + // 이벤트 타입에 따른 발행 토픽 라우팅 + // JSON 문자열을 원래 DTO 클래스로 변환 + private Object deserializePayload(String eventType, String jsonPayload) throws Exception { + return switch (eventType) { + case EVENT_PRODUCT_CREATED -> objectMapper.readValue(jsonPayload, ProductCreatedEvent.class); + case EVENT_PRODUCT_UPDATED -> objectMapper.readValue(jsonPayload, ProductUpdatedEvent.class); + case EVENT_STATUS_CHANGED -> objectMapper.readValue(jsonPayload, ProductStatusChangedEvent.class); + 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); + // 매핑 안 된 이벤트를 String으로 보내면 직렬화 에러 발생! 예외를 던져서 스케줄러 재시도 루프로 넘김 + default -> { + log.warn("등록되지 않은 알 수 없는 이벤트 타입입니다: {}", eventType); + throw new IllegalArgumentException("Unknown event type: " + eventType); + } + }; + } + + private String resolveTopic(String eventType) { + return switch (eventType) { + case EVENT_PRODUCT_CREATED -> topicProductCreated; + case EVENT_PRODUCT_UPDATED -> topicProductUpdated; + case EVENT_STATUS_CHANGED -> topicStatusChanged; + case EVENT_STOCK_RESERVED -> topicStockReserved; + case EVENT_STOCK_RESTORED -> topicStockRestored; + case EVENT_DAILY_RESET -> topicDailyReset; + default -> "inventory.unknown.event"; + }; + } +} diff --git a/src/main/java/com/michelet/inventory/application/ProductCommandService.java b/src/main/java/com/michelet/inventory/application/ProductCommandService.java index b92b99a..5de9dad 100644 --- a/src/main/java/com/michelet/inventory/application/ProductCommandService.java +++ b/src/main/java/com/michelet/inventory/application/ProductCommandService.java @@ -15,18 +15,13 @@ import com.michelet.inventory.domain.repository.ProductOptionRepository; import com.michelet.inventory.domain.repository.ProductRepository; import com.michelet.inventory.domain.repository.StockRepository; -import java.math.BigDecimal; import java.util.ArrayList; import java.util.List; import java.util.UUID; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; -import org.springframework.transaction.support.TransactionSynchronization; -import org.springframework.transaction.support.TransactionSynchronizationManager; @Slf4j @Service @@ -39,41 +34,40 @@ public class ProductCommandService { private final ProductExhibitionRepository productExhibitionRepository; private final StockRepository stockRepository; - private final KafkaTemplate kafkaTemplate; - - @Value("${inventory.kafka.topic.product-created:product.created}") - private String topicProductCreated; - - @Value("${inventory.kafka.topic.product-updated:product.updated}") - private String topicProductUpdated; - - @Value("${inventory.kafka.topic.status-changed:product.status-changed}") - private String topicStatusChanged; + // OutboxHelper로 교체 + private final InventoryOutboxHelper outboxHelper; @Transactional(readOnly = true) public String checkHealth() { return "Inventory Command Service is Healthy"; } - @Transactional public ProductResult createProduct(CreateProductCommand command) { // 1. p_products 저장 Product product = Product.create( - command.restaurantId(), command.name(), command.category(), - command.basePrice(), command.attributes() + command.restaurantId(), + command.name(), + command.category(), + command.basePrice(), + command.attributes() ); productRepository.save(product); // 2. p_product_exhibitions 저장 - productExhibitionRepository.save(ProductExhibition.create( - product, command.exhibition().startAt(), command.exhibition().endAt() - )); + productExhibitionRepository.save( + ProductExhibition.create( + product, + command.exhibition().startAt(), + command.exhibition().endAt() + )); // 3. 옵션(Options) 리스트 생성 및 일괄 저장 List options = command.options().stream() - .map(opt -> ProductOption.create(product, opt.name(), opt.addPrice())) - .toList(); - + .map(opt -> ProductOption.create( + product, + opt.name(), + opt.addPrice() + )).toList(); // saveAll을 호출하면 ID가 채워진 저장된 리스트가 반환됨 List savedOptions = productOptionRepository.saveAll(options); @@ -86,38 +80,33 @@ public ProductResult createProduct(CreateProductCommand command) { // 리스트 크기 불일치 시 명확한 에러 발생 if (savedOptions.size() != requestOptions.size()) { throw new IllegalStateException( - String.format("저장된 옵션 개수(%d)와 요청된 옵션 개수(%d)가 일치하지 않습니다.", - savedOptions.size(), requestOptions.size()) - ); + String.format("저장된 옵션 개수(%d)와 요청된 옵션 개수(%d)가 일치하지 않습니다.", savedOptions.size(), requestOptions.size())); } for (int i = 0; i < savedOptions.size(); i++) { CreateProductCommand.OptionCommand reqOption = requestOptions.get(i); ProductOption dbOption = savedOptions.get(i); - // Positional Arguments 실수 방지를 위해 명시적 지역변수 선언 - BigDecimal addPrice = dbOption.getAddPrice(); - Integer totalQuantity = reqOption.totalQuantity(); - Integer dailyLimit = reqOption.dailyLimit(); - Integer currentDailyStock = Math.min(dailyLimit, totalQuantity); - - // DB 저장을 위한 Stock 객체 생성 - stocks.add(Stock.create( + // DB 저장을 위한 Stock 객체를 먼저 생성 + Stock stock = Stock.create( dbOption.getId(), - totalQuantity, - dailyLimit, + reqOption.totalQuantity(), + reqOption.dailyLimit(), reqOption.maxLimit() - )); + ); + stocks.add(stock); // 카프카 전송을 위한 Event DTO 생성 - optionEventDtos.add(new ProductCreatedEvent.OptionEventDto( - dbOption.getId(), - dbOption.getName(), - addPrice, - totalQuantity, - currentDailyStock, - dailyLimit - )); + // 서비스에서 Math.min을 직접 계산하지 않고, 도메인(Stock)이 계산한 최종 값을 DTO에 세팅 + optionEventDtos.add( + new ProductCreatedEvent.OptionEventDto( + dbOption.getId(), + dbOption.getName(), + dbOption.getAddPrice(), + stock.getTotalQuantity(), + stock.getCurrentDailyStock(), // Stock 객체에서 꺼내씀! + stock.getDailyLimit() + )); } stockRepository.saveAll(stocks); @@ -134,29 +123,16 @@ public ProductResult createProduct(CreateProductCommand command) { optionEventDtos ); - // 6. DB 커밋 완료 후에만 카프카 메시지 전송 (정합성 보장) 및 Callback 확인 - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - sendKafkaMessageWithCallback(topicProductCreated, event.productId(), event); - } - }); - } else { - // 단위 테스트 등 트랜잭션 동기화가 활성화되지 않은 환경을 위한 폴백 - sendKafkaMessageWithCallback(topicProductCreated, event.productId(), event); - } + // 6. 트랜잭션 동기화 및 카프카 직접 전송 로직 제거 후 Outbox 저장 + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_CREATED", event); // 7. 결과 반환 (옵션 리스트 포함) List optionResults = savedOptions.stream() - .map(opt -> new ProductResult.OptionResult(opt.getId(), opt.getName())) - .toList(); - + .map(opt -> new ProductResult.OptionResult(opt.getId(), opt.getName())).toList(); return new ProductResult(product.getId(), optionResults); } // 상품 수정 로직 - @Transactional public void updateProduct(UUID productId, UpdateProductCommand command) { Product product = productRepository.findById(productId) .orElseThrow(() -> new IllegalArgumentException("상품을 찾을 수 없습니다.")); @@ -172,21 +148,10 @@ public void updateProduct(UUID productId, UpdateProductCommand command) { product.getBasePrice(), product.getAttributes() ); - - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - sendKafkaMessageWithCallback(topicProductUpdated, product.getId(), event); - } - }); - } else { - sendKafkaMessageWithCallback(topicProductUpdated, product.getId(), event); - } + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_UPDATED", event); } // 상품 삭제 로직 - @Transactional public void deleteProduct(UUID productId) { Product product = productRepository.findById(productId) .orElseThrow(() -> new IllegalArgumentException("상품을 찾을 수 없습니다.")); @@ -195,31 +160,6 @@ public void deleteProduct(UUID productId) { product.changeStatus(ProductStatus.DELETED); ProductStatusChangedEvent event = new ProductStatusChangedEvent(product.getId(), product.getStatus().name()); - - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - sendKafkaMessageWithCallback(topicStatusChanged, product.getId(), event); - } - }); - } else { - sendKafkaMessageWithCallback(topicStatusChanged, product.getId(), event); - } - } - - // 범용적으로 쓸 수 있게 파라미터 변경 - private void sendKafkaMessageWithCallback(String topic, UUID key, Object event) { - log.info("DB 커밋 완료. 카프카 이벤트 발행 요청: topic={}, key={}", topic, key); - kafkaTemplate.send(topic, key.toString(), event) - .whenComplete((result, ex) -> { - if (ex == null) { - log.info("이벤트 발행 실제 성공: topic={}, key={}, offset={}", - topic, key, result.getRecordMetadata().offset()); - } else { - log.error("이벤트 발행 실패 (Dead Letter Queue 처리 필요): topic={}, key={}", - topic, key, ex); - } - }); + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_STATUS_CHANGED", event); } } diff --git a/src/main/java/com/michelet/inventory/application/StockCommandService.java b/src/main/java/com/michelet/inventory/application/StockCommandService.java index f08a076..76ef241 100644 --- a/src/main/java/com/michelet/inventory/application/StockCommandService.java +++ b/src/main/java/com/michelet/inventory/application/StockCommandService.java @@ -15,12 +15,8 @@ import com.michelet.inventory.presentation.dto.RestoreStockRequest; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; -import org.springframework.transaction.support.TransactionSynchronization; -import org.springframework.transaction.support.TransactionSynchronizationManager; @Slf4j @Service @@ -30,24 +26,16 @@ public class StockCommandService { private final StockRepository stockRepository; private final ProductOptionRepository productOptionRepository; private final ProductRepository productRepository; - private final KafkaTemplate kafkaTemplate; - @Value("${inventory.kafka.topic.reserved:stock.reserved}") - private String topicStockReserved; + // OutboxHelper 주입 (KafkaTemplate 대체) + private final InventoryOutboxHelper outboxHelper; - @Value("${inventory.kafka.topic.restored:stock.restored}") - private String topicStockRestored; - - @Value("${inventory.kafka.topic.status-changed:product.status-changed}") - private String topicStatusChanged; - - // while문, Thread.sleep, TransactionTemplate 모두 제거 및 순수 비즈니스 로직만 남김 @Transactional public void reserveStock(ReserveStockRequest request) { Stock stock = stockRepository.findById(request.optionId()) .orElseThrow(StockNotFoundException::new); - // 1. 도메인 로직을 통한 3중 검증 및 차감 + // 1. 도메인 로직 검증 및 차감 stock.reserve(request.quantity()); stockRepository.save(stock); @@ -67,22 +55,18 @@ public void reserveStock(ReserveStockRequest request) { ProductStatusChangedEvent statusEvent = new ProductStatusChangedEvent(product.getId(), product.getStatus().name()); - publishKafkaEvent( - topicStatusChanged, - product.getId().toString(), - statusEvent, - "SOLDOUT 상태 변경" - ); + // 카프카 직접 발송 대신 Outbox에 적재 + outboxHelper.append("PRODUCT", product.getId().toString(), "PRODUCT_STATUS_CHANGED", statusEvent); } } - // 3. Kafka 이벤트 발행 (DB 커밋 이후에 실행되도록 공통 메서드 분리) + // 3. 재고 차감 이벤트 Outbox 적재 StockReservedEvent event = new StockReservedEvent( stock.getOptionId(), stock.getTotalQuantity(), stock.getCurrentDailyStock() ); - publishKafkaEvent(topicStockReserved, request.optionId().toString(), event, "재고차감"); + outboxHelper.append("STOCK", request.optionId().toString(), "STOCK_RESERVED", event); } @Transactional @@ -96,40 +80,12 @@ public void restoreStock(RestoreStockRequest request) { // 2. DB 업데이트 (더티 체킹 후 flush) stockRepository.save(stock); - // 2. Kafka 이벤트 발행 + // 3. 재고 복구 이벤트를 Outbox에 적재 StockRestoredEvent event = new StockRestoredEvent( stock.getOptionId(), stock.getTotalQuantity(), stock.getCurrentDailyStock() ); - publishKafkaEvent(topicStockRestored, request.optionId().toString(), event, "재고복구"); - } - - // DB 커밋 완료 후에만 Kafka가 발행되도록 보장하는 공통 메서드 - private void publishKafkaEvent(String topic, String key, Object event, String logPrefix) { - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - sendToKafka(topic, key, event, logPrefix); - } - }); - } else { - sendToKafka(topic, key, event, logPrefix); - } - } - - private void sendToKafka(String topic, String key, Object event, String logPrefix) { - kafkaTemplate.send(topic, key, event) - .whenComplete((result, ex) -> { - if (ex != null) { - log.error("{} Kafka 메시지 발행 실패 (데이터 불일치 위험)! key: {}", logPrefix, key, ex); - } else { - long offset = - (result != null && result.getRecordMetadata() != null) ? result.getRecordMetadata().offset() - : -1; - log.info("{} Kafka 메시지 발행 성공! key: {}, offset: {}", logPrefix, key, offset); - } - }); + outboxHelper.append("STOCK", request.optionId().toString(), "STOCK_RESTORED", event); } } diff --git a/src/main/java/com/michelet/inventory/application/StockSchedulerService.java b/src/main/java/com/michelet/inventory/application/StockSchedulerService.java index 1102258..19bf1a7 100644 --- a/src/main/java/com/michelet/inventory/application/StockSchedulerService.java +++ b/src/main/java/com/michelet/inventory/application/StockSchedulerService.java @@ -6,13 +6,9 @@ import java.time.ZoneId; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.beans.factory.annotation.Value; -import org.springframework.kafka.core.KafkaTemplate; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; -import org.springframework.transaction.support.TransactionSynchronization; -import org.springframework.transaction.support.TransactionSynchronizationManager; @Slf4j @Service @@ -20,17 +16,10 @@ public class StockSchedulerService { private final StockRepository stockRepository; - private final KafkaTemplate kafkaTemplate; - - @Value("${inventory.kafka.topic.daily-reset:stock.daily-reset}") - private String topicDailyReset; - - // "초 분 시 일 월 요일" 순서임 - // 매일 자정(0시 0분 0초)에 실행하려면: @Scheduled(cron = "0 0 0 * * *") - // (테스트용) 1분마다 실행하려면: @Scheduled(cron = "0 * * * * *") + private final InventoryOutboxHelper outboxHelper; @Transactional - @Scheduled(cron = "0 0 0 * * *", zone = "Asia/Seoul") // 한국 시각 기준 매일 자정 + @Scheduled(cron = "0 0 0 * * *", zone = "Asia/Seoul") public void resetDailyStocks() { log.info("일일 재고 초기화 스케줄러 시작..."); @@ -38,31 +27,12 @@ public void resetDailyStocks() { log.info("일일 재고 초기화 완료! 업데이트된 상품 옵션 수: {}", updatedCount); - // 카탈로그 서비스(MongoDB)에도 이 초기화 사실을 알려야해서 - 여기서 전체 재고 초기화 이벤트를 Kafka로 던져줌 if (updatedCount > 0) { LocalDate todayInSeoul = LocalDate.now(ZoneId.of("Asia/Seoul")); DailyStockResetEvent event = new DailyStockResetEvent(todayInSeoul, updatedCount); - // DB 트랜잭션이 정상적으로 커밋된 이후에만 발행되도록! - if (TransactionSynchronizationManager.isSynchronizationActive()) { - TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { - @Override - public void afterCommit() { - sendKafkaEvent(event); - } - }); - } else { - sendKafkaEvent(event); - } + // 현재 트랜잭션의 일부로 Outbox 테이블에 적재됨 (커밋 시 함께 저장) + outboxHelper.append("STOCK", "ALL_STOCKS", "DAILY_STOCK_RESET", event); } } - - private void sendKafkaEvent(DailyStockResetEvent event) { - kafkaTemplate.send(topicDailyReset, "ALL_STOCKS", event) - .whenComplete((result, ex) -> { - if (ex != null) { - log.error("일일 재고 초기화 카프카 이벤트 발행 실패", ex); - } - }); - } } diff --git a/src/main/java/com/michelet/inventory/domain/model/InventoryOutbox.java b/src/main/java/com/michelet/inventory/domain/model/InventoryOutbox.java new file mode 100644 index 0000000..8e897b0 --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/model/InventoryOutbox.java @@ -0,0 +1,77 @@ +package com.michelet.inventory.domain.model; + +import com.michelet.common.entity.BaseEntity; +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.EnumType; +import jakarta.persistence.Enumerated; +import jakarta.persistence.GeneratedValue; +import jakarta.persistence.GenerationType; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import jakarta.persistence.Version; +import java.util.UUID; +import lombok.AccessLevel; +import lombok.Builder; +import lombok.Getter; +import lombok.NoArgsConstructor; +import org.hibernate.annotations.JdbcTypeCode; +import org.hibernate.type.SqlTypes; + +@Entity +@Table(name = "p_inventory_outbox") +@Getter +@NoArgsConstructor(access = AccessLevel.PROTECTED) +public class InventoryOutbox extends BaseEntity { + + @Id + @GeneratedValue(strategy = GenerationType.UUID) + private UUID id; + + @Column(nullable = false, length = 50) + private String aggregateType; + + @Column(nullable = false, length = 50) + private String aggregateId; + + @Column(nullable = false, length = 50) + private String eventType; + + @JdbcTypeCode(SqlTypes.JSON) + @Column(columnDefinition = "jsonb", nullable = false) + private String payload; + + @Enumerated(EnumType.STRING) + @Column(nullable = false, length = 20) + private OutboxStatus status; + + @Column(nullable = false) + private int retryCount = 0; // 재시도 횟수 관리 (기본값 0) + + @Version // 스케줄러 동시성 제어를 위한 낙관적 락 + private Long version; + + @Builder + private InventoryOutbox(String aggregateType, String aggregateId, String eventType, String payload) { + this.aggregateType = aggregateType; + this.aggregateId = aggregateId; + this.eventType = eventType; + this.payload = payload; + this.status = OutboxStatus.INIT; + this.retryCount = 0; + } + + public void markAsPublished() { + this.status = OutboxStatus.PUBLISHED; + } + + // 영구 실패 상태 전이 + public void markAsFailed() { + this.status = OutboxStatus.FAILED; + } + + // 재시도 횟수 증가 로직 + public void incrementRetryCount() { + this.retryCount++; + } +} diff --git a/src/main/java/com/michelet/inventory/domain/model/OutboxStatus.java b/src/main/java/com/michelet/inventory/domain/model/OutboxStatus.java new file mode 100644 index 0000000..9791961 --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/model/OutboxStatus.java @@ -0,0 +1,7 @@ +package com.michelet.inventory.domain.model; + +public enum OutboxStatus { + INIT, // 발행 대기 + PUBLISHED, // 발행 완료 + FAILED // 최대 재시도 초과로 영구 실패 (격리) +} diff --git a/src/main/java/com/michelet/inventory/domain/repository/InventoryOutboxRepository.java b/src/main/java/com/michelet/inventory/domain/repository/InventoryOutboxRepository.java new file mode 100644 index 0000000..3ce1c50 --- /dev/null +++ b/src/main/java/com/michelet/inventory/domain/repository/InventoryOutboxRepository.java @@ -0,0 +1,16 @@ +package com.michelet.inventory.domain.repository; + +import com.michelet.inventory.domain.model.InventoryOutbox; +import com.michelet.inventory.domain.model.OutboxStatus; +import java.util.List; +import java.util.Optional; +import java.util.UUID; + +public interface InventoryOutboxRepository { + InventoryOutbox save(InventoryOutbox outbox); + + Optional findById(UUID id); + + // 스케줄러에서 사용할 배치 조회 메서드 + List findTop50ByStatusOrderByCreatedAtAsc(OutboxStatus status); +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/config/KafkaConfig.java b/src/main/java/com/michelet/inventory/infrastructure/config/KafkaConfig.java index 34b2c40..1fc6c44 100644 --- a/src/main/java/com/michelet/inventory/infrastructure/config/KafkaConfig.java +++ b/src/main/java/com/michelet/inventory/infrastructure/config/KafkaConfig.java @@ -1,42 +1,95 @@ package com.michelet.inventory.infrastructure.config; +import java.util.HashMap; +import java.util.Map; +import lombok.extern.slf4j.Slf4j; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.autoconfigure.kafka.ConcurrentKafkaListenerContainerFactoryConfigurer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory; import org.springframework.kafka.core.ConsumerFactory; +import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.core.KafkaTemplate; import org.springframework.kafka.listener.DeadLetterPublishingRecoverer; import org.springframework.kafka.listener.DefaultErrorHandler; import org.springframework.util.backoff.FixedBackOff; +@Slf4j @Configuration public class KafkaConfig { /** * 카프카 컨슈머 에러 핸들러 설정 - 예외 발생 시 1초 간격으로 3번 재시도 후 DLT 토픽으로 전송 */ + @Value("${inventory.kafka.consumer.retry.interval-ms:1000}") + private long retryIntervalMs; + + @Value("${inventory.kafka.consumer.retry.max-attempts:3}") + private long retryMaxAttempts; + @Bean public DefaultErrorHandler errorHandler(KafkaTemplate kafkaTemplate) { // 1. 에러가 난 메시지를 DLT(Dead Letter Topic, 예: stock.restored.DLT)로 보내는 역할 - DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(kafkaTemplate); + // 파티션 충돌 방지 및 추적 로그 삽입 + DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer( + kafkaTemplate, + (cr, e) -> { + log.error("[Inventory 장애 감지] 에러 핸들러 작동! DLT 토픽으로 이동. 대상: {}", cr.topic() + ".DLT"); + // 원본 파티션을 고집하지 않고 가용한 파티션에 안전하게 넣도록 강제 (-1) + return new TopicPartition(cr.topic() + ".DLT", -1); + } + ); // 2. 기본 재시도 정책: 1초(1000ms) 간격으로 최대 3번 재시도 - DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, new FixedBackOff(1000L, 3L)); + DefaultErrorHandler errorHandler = new DefaultErrorHandler(recoverer, + new FixedBackOff(retryIntervalMs, retryMaxAttempts)); - // 3. 특정 예외(형식이 아예 틀린 경우)는 재시도해봤자 의미 없으므로 즉시 DLT로 직행 - errorHandler.addNotRetryableExceptions(IllegalArgumentException.class); + // 3. 재시도가 의미 없는 예외만 즉시 DLT로 보냄 + errorHandler.addNotRetryableExceptions( + org.springframework.kafka.support.serializer.DeserializationException.class, + IllegalArgumentException.class + ); return errorHandler; } @Bean - public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory( - ConsumerFactory consumerFactory, + public ConcurrentKafkaListenerContainerFactory kafkaListenerContainerFactory( + ConcurrentKafkaListenerContainerFactoryConfigurer configurer, + ConsumerFactory consumerFactory, DefaultErrorHandler errorHandler // 위에서 만든 에러 핸들러 주입 ) { - ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); - factory.setConsumerFactory(consumerFactory); - factory.setCommonErrorHandler(errorHandler); // 컨슈머 팩토리에 에러 핸들러 장착 + ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); + configurer.configure(factory, consumerFactory); + factory.setCommonErrorHandler(errorHandler); + return factory; + } + + @Bean + public ConcurrentKafkaListenerContainerFactory dltListenerContainerFactory( + ConsumerFactory consumerFactory + ) { + ConcurrentKafkaListenerContainerFactory factory = new ConcurrentKafkaListenerContainerFactory<>(); + Map props = new HashMap<>(consumerFactory.getConfigurationProperties()); + + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); + + // ErrorHandlingDeserializer/Json 관련 잔여 설정 일관성 있게 정리 + props.remove("spring.deserializer.key.delegate.class"); + props.remove("spring.deserializer.value.delegate.class"); + props.remove("spring.json.trusted.packages"); + props.remove("spring.json.type.mapping"); + + factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(props)); + + // DLT 처리 실패 시 기본 핸들러(FixedBackOff(0,9))로 폴백되지 않도록 명시적 지정 + factory.setCommonErrorHandler(new DefaultErrorHandler(new FixedBackOff(0L, 0L))); + return factory; } } diff --git a/src/main/java/com/michelet/inventory/infrastructure/repository/InventoryOutboxRepositoryImpl.java b/src/main/java/com/michelet/inventory/infrastructure/repository/InventoryOutboxRepositoryImpl.java new file mode 100644 index 0000000..c50288b --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/repository/InventoryOutboxRepositoryImpl.java @@ -0,0 +1,32 @@ +package com.michelet.inventory.infrastructure.repository; + +import com.michelet.inventory.domain.model.InventoryOutbox; +import com.michelet.inventory.domain.model.OutboxStatus; +import com.michelet.inventory.domain.repository.InventoryOutboxRepository; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Repository; + +@Repository +@RequiredArgsConstructor +public class InventoryOutboxRepositoryImpl implements InventoryOutboxRepository { + + private final JpaInventoryOutboxRepository jpaRepository; + + @Override + public InventoryOutbox save(InventoryOutbox outbox) { + return jpaRepository.save(outbox); + } + + @Override + public Optional findById(UUID id) { + return jpaRepository.findById(id); + } + + @Override + public List findTop50ByStatusOrderByCreatedAtAsc(OutboxStatus status) { + return jpaRepository.findTop50ByStatusOrderByCreatedAtAsc(status); + } +} diff --git a/src/main/java/com/michelet/inventory/infrastructure/repository/JpaInventoryOutboxRepository.java b/src/main/java/com/michelet/inventory/infrastructure/repository/JpaInventoryOutboxRepository.java new file mode 100644 index 0000000..f0e25e9 --- /dev/null +++ b/src/main/java/com/michelet/inventory/infrastructure/repository/JpaInventoryOutboxRepository.java @@ -0,0 +1,13 @@ +package com.michelet.inventory.infrastructure.repository; + +import com.michelet.inventory.domain.model.InventoryOutbox; +import com.michelet.inventory.domain.model.OutboxStatus; +import java.util.List; +import java.util.UUID; +import org.springframework.data.jpa.repository.JpaRepository; + +public interface JpaInventoryOutboxRepository extends JpaRepository { + + // OOM 방지 및 생성 시간순 처리를 위한 Top N 쿼리 + List findTop50ByStatusOrderByCreatedAtAsc(OutboxStatus status); +} diff --git a/src/main/java/com/michelet/inventory/presentation/SchedulerTriggerController.java b/src/main/java/com/michelet/inventory/presentation/SchedulerTriggerController.java index e4a2528..89b4762 100644 --- a/src/main/java/com/michelet/inventory/presentation/SchedulerTriggerController.java +++ b/src/main/java/com/michelet/inventory/presentation/SchedulerTriggerController.java @@ -2,6 +2,7 @@ import com.michelet.common.response.ApiResponse; import com.michelet.inventory.application.ExhibitionSchedulerService; +import com.michelet.inventory.application.InventoryOutboxScheduler; import com.michelet.inventory.application.StockSchedulerService; import lombok.RequiredArgsConstructor; import org.springframework.http.ResponseEntity; @@ -16,6 +17,7 @@ public class SchedulerTriggerController { private final StockSchedulerService stockSchedulerService; private final ExhibitionSchedulerService exhibitionSchedulerService; + private final InventoryOutboxScheduler inventoryOutboxScheduler; // 일일 재고 강제 리셋 버튼 - 테스트용! @PostMapping("/trigger-stock") @@ -30,4 +32,11 @@ public ResponseEntity> triggerExhibition() { exhibitionSchedulerService.updateExhibitionStatus(); return ResponseEntity.ok(ApiResponse.ok("전시 상태 갱신 트리거 작동 완료")); } + + // Outbox 스케줄러 강제 갱신 버튼 - 테스트용! + @PostMapping("/trigger-outbox") + public ResponseEntity> triggerOutbox() { + inventoryOutboxScheduler.processOutboxEvents(); + return ResponseEntity.ok(ApiResponse.ok("Outbox 스케줄러 작동 완료 (Kafka 발행됨)")); + } } diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index df685ff..eb0fe4b 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -51,6 +51,12 @@ inventory: status-changed: "product.status-changed" daily-reset: "stock.daily-reset" product-updated: "product.updated" + consumer: + retry: + interval-ms: 1000 + max-attempts: 3 + outbox: + max-retries: 3 # 최대 재시도 횟수 설정 internal: auth: diff --git a/src/test/java/com/michelet/inventory/application/ProductCommandServiceTest.java b/src/test/java/com/michelet/inventory/application/ProductCommandServiceTest.java index 97dfd34..620a477 100644 --- a/src/test/java/com/michelet/inventory/application/ProductCommandServiceTest.java +++ b/src/test/java/com/michelet/inventory/application/ProductCommandServiceTest.java @@ -3,10 +3,8 @@ import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatThrownBy; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.BDDMockito.given; -import static org.mockito.Mockito.lenient; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -27,8 +25,6 @@ import java.util.List; import java.util.Map; import java.util.UUID; -import java.util.concurrent.CompletableFuture; -import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -37,8 +33,6 @@ import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; -import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.test.util.ReflectionTestUtils; @ExtendWith(MockitoExtension.class) class ProductCommandServiceTest { @@ -54,8 +48,11 @@ class ProductCommandServiceTest { private ProductExhibitionRepository productExhibitionRepository; @Mock private StockRepository stockRepository; + + // KafkaTemplate 대신 OutboxHelper 주입 @Mock - private KafkaTemplate kafkaTemplate; + private InventoryOutboxHelper outboxHelper; + @Captor private ArgumentCaptor eventCaptor; @@ -69,13 +66,8 @@ class ProductCommandServiceTest { @Captor private ArgumentCaptor> stocksCaptor; - @BeforeEach - void setUp() { - ReflectionTestUtils.setField(productCommandService, "topicProductCreated", "product.created"); - } - @Test - @DisplayName("성공: 상품 등록 시 모든 도메인 모델(상품/전시/옵션/재고)이 올바른 값으로 저장소에 전달되어야 한다") + @DisplayName("성공: 상품 등록 시 모든 도메인 모델이 올바른 값으로 저장되고 Outbox에 적재되어야 한다") void createProductUnitTest() { // given LocalDateTime now = LocalDateTime.now(); @@ -98,7 +90,7 @@ void createProductUnitTest() { // 1. Product 저장 시 ID 자동 생성 모킹 given(productRepository.save(any(Product.class))).willAnswer(invocation -> { Product product = invocation.getArgument(0); - ReflectionTestUtils.setField(product, "id", UUID.randomUUID()); + org.springframework.test.util.ReflectionTestUtils.setField(product, "id", UUID.randomUUID()); return product; }); @@ -106,19 +98,11 @@ void createProductUnitTest() { given(productOptionRepository.saveAll(any())).willAnswer(invocation -> { List options = invocation.getArgument(0); for (ProductOption option : options) { - // JPA가 DB 삽입 후 ID를 채워주는 동작을 리플렉션으로 강제 시뮬레이션 - ReflectionTestUtils.setField(option, "id", UUID.randomUUID()); + org.springframework.test.util.ReflectionTestUtils.setField(option, "id", UUID.randomUUID()); } return options; }); - // 3. KafkaTemplate 모킹: send 호출 시 빈 가짜 영수증(CompletableFuture?) 반환 - CompletableFuture> mockFuture - = CompletableFuture.completedFuture(new org.springframework.kafka.support.SendResult<>(null, null)); - - // lenient()를 추가 - Mockito의 엄격한 Stubbing 검사(PotentialStubbingProblem) 유연하게 통과시킴 - lenient().when(kafkaTemplate.send(anyString(), anyString(), any())).thenReturn(mockFuture); - // when ProductResult result = productCommandService.createProduct(command); @@ -131,53 +115,32 @@ void createProductUnitTest() { // 1. 반환 결과(ProductResult) 및 상품 검증 Product savedProduct = productCaptor.getValue(); - // Kafka 발행 시 '파티션 키(Partition Key)'가 ProductId로 정확히 매핑되었는지 검증 - verify(kafkaTemplate, times(1)).send( - eq("product.created"), + // Kafka 직접 전송 대신 Outbox 적재 여부 검증 + verify(outboxHelper, times(1)).append( + eq("PRODUCT"), eq(savedProduct.getId().toString()), + eq("PRODUCT_CREATED"), eventCaptor.capture() ); ProductCreatedEvent capturedEvent = eventCaptor.getValue(); - // 캡처된 Kafka 이벤트의 페이로드(내용물)가 정확한지 추가 검증 + // 캡처된 Kafka 이벤트의 페이로드 전체 필드 정밀 검증 assertThat(capturedEvent.productId()).isEqualTo(savedProduct.getId()); - assertThat(capturedEvent.restaurantId()).isEqualTo(command.restaurantId()); assertThat(capturedEvent.name()).isEqualTo("미슐랭 밀키트 세트"); - assertThat(capturedEvent.category()).isEqualTo(ProductCategory.MEALKIT.name()); + assertThat(capturedEvent.category()).isEqualTo("MEALKIT"); + assertThat(capturedEvent.basePrice()).isEqualByComparingTo(new BigDecimal("45000")); + assertThat(capturedEvent.attributes()).containsEntry("servings", 2).containsEntry("cookingTime", "20min"); - assertThat(capturedEvent.basePrice()).isEqualByComparingTo(command.basePrice()); + // 옵션이 정확히 DTO로 변환되었는지 검증 + assertThat(capturedEvent.options()).hasSize(2); + assertThat(capturedEvent.options().get(0).name()).isEqualTo("맵기 보통"); + assertThat(capturedEvent.options().get(0).totalQuantity()).isEqualTo(100); + assertThat(capturedEvent.options().get(0).currentDailyStock()).isEqualTo(20); + assertThat(capturedEvent.options().get(0).dailyLimit()).isEqualTo(20); assertThat(result).isNotNull(); - assertThat(result.productId()).isNotNull(); assertThat(result.productId()).isEqualTo(savedProduct.getId()); - assertThat(result.options()).hasSize(2); // 옵션이 2개 반환되었는지 검증 - assertThat(result.options().get(0).optionId()).isNotNull(); // 발급된 ID 검증 - - assertThat(savedProduct.getName()).isEqualTo("미슐랭 밀키트 세트"); - assertThat(savedProduct.getAttributes().get("servings")).isEqualTo(2); - - // 2. 전시 정보 검증 - ProductExhibition savedExhibition = exhibitionCaptor.getValue(); - assertThat(savedExhibition.getProduct()).isEqualTo(savedProduct); - assertThat(savedExhibition.getStartAt()).isEqualTo(command.exhibition().startAt()); - - // 3. 옵션 검증 - List savedOptions = optionsCaptor.getValue(); - assertThat(savedOptions).hasSize(2); - - // 4. 재고 검증 - List savedStocks = stocksCaptor.getValue(); - assertThat(savedStocks).hasSize(2); - - Stock normalStock = savedStocks.get(0); - assertThat(normalStock.getTotalQuantity()).isEqualTo(100); - assertThat(normalStock.getCurrentDailyStock()).isEqualTo(20); - - Stock spicyStock = savedStocks.get(1); - assertThat(spicyStock.getTotalQuantity()).isEqualTo(80); - assertThat(spicyStock.getCurrentDailyStock()).isEqualTo(15); - assertThat(spicyStock.getMaxLimit()).isEqualTo(1); } @Test diff --git a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java index ecbe9c5..67532c9 100644 --- a/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java +++ b/src/test/java/com/michelet/inventory/application/StockCommandServiceTest.java @@ -23,8 +23,6 @@ import com.michelet.inventory.presentation.dto.RestoreStockRequest; import java.util.Optional; import java.util.UUID; -import java.util.concurrent.CompletableFuture; -import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; @@ -33,8 +31,6 @@ import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; -import org.springframework.kafka.core.KafkaTemplate; -import org.springframework.test.util.ReflectionTestUtils; @ExtendWith(MockitoExtension.class) class StockCommandServiceTest { @@ -53,7 +49,7 @@ class StockCommandServiceTest { private ProductRepository productRepository; @Mock - private KafkaTemplate kafkaTemplate; + private InventoryOutboxHelper outboxHelper; // 카프카로 전송된 이벤트를 낚아채서 내부 값을 검증하기 위한 Captor @Captor @@ -62,15 +58,8 @@ class StockCommandServiceTest { @Captor private ArgumentCaptor restoredEventCaptor; - @BeforeEach - void setUp() { - ReflectionTestUtils.setField(stockCommandService, "topicStockReserved", "stock.reserved"); - ReflectionTestUtils.setField(stockCommandService, "topicStockRestored", "stock.restored"); - ReflectionTestUtils.setField(stockCommandService, "topicStatusChanged", "product.status-changed"); - } - @Test - @DisplayName("성공: 재고가 충분하면 차감되고 Kafka 이벤트가 발행된다.") + @DisplayName("성공: 재고가 충분하면 차감되고 Outbox에 이벤트가 적재된다.") void reserveStock_Success() { // given UUID optionId = UUID.randomUUID(); @@ -78,16 +67,17 @@ void reserveStock_Success() { ReserveStockRequest request = new ReserveStockRequest(optionId, 2, null); given(stockRepository.findById(optionId)).willReturn(Optional.of(stock)); - given(kafkaTemplate.send(eq("stock.reserved"), eq(optionId.toString()), any())) - .willReturn(CompletableFuture.completedFuture(null)); // when stockCommandService.reserveStock(request); // then verify(stockRepository, times(1)).save(any(Stock.class)); - // 캡처를 통해 Kafka로 넘어간 페이로드(이벤트 객체)의 상세 데이터 검증 - verify(kafkaTemplate, times(1)).send(eq("stock.reserved"), eq(optionId.toString()), eventCaptor.capture()); + + // 카프카 전송 대신 OutboxHelper의 append 호출 여부 검증 + verify(outboxHelper, times(1)).append(eq("STOCK"), eq(optionId.toString()), eq("STOCK_RESERVED"), + eventCaptor.capture()); + StockReservedEvent capturedEvent = eventCaptor.getValue(); assertThat(capturedEvent.optionId()).isEqualTo(optionId); assertThat(capturedEvent.totalQuantity()).isEqualTo(98); // 100 - 2 @@ -108,8 +98,8 @@ void reserveStock_Fail_NotEnoughDailyStock() { assertThatThrownBy(() -> stockCommandService.reserveStock(request)) .isInstanceOf(OutOfStockException.class); - // 예외가 터졌으므로 Kafka 전송은 절대 일어나지 않아야 함 - verifyNoInteractions(kafkaTemplate); + // 예외가 터졌으므로 Outbox 적재는 절대 일어나지 않아야 함 + verifyNoInteractions(outboxHelper); } // 전체 재고 부족 예외 테스트 @@ -127,7 +117,7 @@ void reserveStock_Fail_NotEnoughTotalQuantity() { assertThatThrownBy(() -> stockCommandService.reserveStock(request)) .isInstanceOf(SoldOutException.class); - verifyNoInteractions(kafkaTemplate); + verifyNoInteractions(outboxHelper); } // 1인당 구매 한도 초과 예외 테스트 @@ -137,7 +127,7 @@ void reserveStock_Fail_MaxLimitExceeded() { // given UUID optionId = UUID.randomUUID(); Stock stock = Stock.create(optionId, 100, 50, 1); // 1인당 1개 제한 - ReserveStockRequest request = new ReserveStockRequest(optionId, 2, null); //2개 요청 + ReserveStockRequest request = new ReserveStockRequest(optionId, 2, null); // 2개 요청 given(stockRepository.findById(optionId)).willReturn(Optional.of(stock)); @@ -145,11 +135,11 @@ void reserveStock_Fail_MaxLimitExceeded() { assertThatThrownBy(() -> stockCommandService.reserveStock(request)) .isInstanceOf(MaxLimitExceededException.class); - verifyNoInteractions(kafkaTemplate); + verifyNoInteractions(outboxHelper); } @Test - @DisplayName("성공: 재고 복구 경로 정상 동작 및 Kafka 발행 테스트") + @DisplayName("성공: 재고 복구 경로 정상 동작 및 Outbox 적재 테스트") void restoreStock_Success() { UUID optionId = UUID.randomUUID(); // 1. 초기 재고 100개, 일일 재고 50개 생성 @@ -162,15 +152,13 @@ void restoreStock_Success() { RestoreStockRequest request = new RestoreStockRequest(optionId, 2, null); given(stockRepository.findById(optionId)).willReturn(Optional.of(stock)); - given(kafkaTemplate.send(eq("stock.restored"), eq(optionId.toString()), any())) - .willReturn(CompletableFuture.completedFuture(null)); // when stockCommandService.restoreStock(request); // then verify(stockRepository, times(1)).save(any(Stock.class)); - verify(kafkaTemplate, times(1)).send(eq("stock.restored"), eq(optionId.toString()), + verify(outboxHelper, times(1)).append(eq("STOCK"), eq(optionId.toString()), eq("STOCK_RESTORED"), restoredEventCaptor.capture()); StockRestoredEvent event = restoredEventCaptor.getValue(); @@ -192,6 +180,6 @@ void restoreStock_Fail_NotFound() { assertThatThrownBy(() -> stockCommandService.restoreStock(request)) .isInstanceOf(StockNotFoundException.class); - verifyNoInteractions(kafkaTemplate); + verifyNoInteractions(outboxHelper); } }