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
24 changes: 24 additions & 0 deletions api-test/inventory.http
Original file line number Diff line number Diff line change
Expand Up @@ -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("=======================================");
});
%}
Original file line number Diff line number Diff line change
Expand Up @@ -9,25 +9,21 @@
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
@RequiredArgsConstructor
public class ExhibitionSchedulerService {

private final ProductRepository productRepository;
private final KafkaTemplate<String, Object> kafkaTemplate;
private final InventoryOutboxHelper outboxHelper;

// Spring AOP의 프록시를 타기 위한 자기 자신 주입
@Lazy
Expand All @@ -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);
Expand Down Expand Up @@ -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();
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
@@ -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);
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}
Original file line number Diff line number Diff line change
@@ -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<String, Object> 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<InventoryOutbox> 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());
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
}

// 재시도 횟수 처리 및 상태 변경을 돕는 실패 처리 메서드
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);
}
};
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

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";
};
}
}
Loading