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 @@ -3,7 +3,6 @@
import com.michelet.inventory.application.dto.ProductStatusChangedEvent;
import com.michelet.inventory.application.dto.StockReservedEvent;
import com.michelet.inventory.application.dto.StockRestoredEvent;
import com.michelet.inventory.domain.exception.ConcurrencyFailureException;
import com.michelet.inventory.domain.exception.StockNotFoundException;
import com.michelet.inventory.domain.model.Product;
import com.michelet.inventory.domain.model.ProductOption;
Expand All @@ -14,16 +13,14 @@
import com.michelet.inventory.domain.repository.StockRepository;
import com.michelet.inventory.presentation.dto.ReserveStockRequest;
import com.michelet.inventory.presentation.dto.RestoreStockRequest;
import jakarta.annotation.PostConstruct;
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.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
import org.springframework.transaction.support.TransactionSynchronization;
import org.springframework.transaction.support.TransactionSynchronizationManager;
import org.springframework.transaction.support.TransactionTemplate;

@Slf4j
@Service
Expand All @@ -33,12 +30,7 @@ public class StockCommandService {
private final StockRepository stockRepository;
private final ProductOptionRepository productOptionRepository;
private final ProductRepository productRepository;

private final KafkaTemplate<String, Object> kafkaTemplate;
private final TransactionTemplate transactionTemplate; // 트랜잭션을 수동으로 제어하기 위해 주입

@Value("${inventory.retry.max-count:3}")
private int maxRetryCount;

@Value("${inventory.kafka.topic.reserved:stock.reserved}")
private String topicStockReserved;
Expand All @@ -49,63 +41,17 @@ public class StockCommandService {
@Value("${inventory.kafka.topic.status-changed:product.status-changed}")
private String topicStatusChanged;


// 재시도 횟수 설정 오류 방어 - 최소 1회 실행 보장
@PostConstruct
public void validateConfig() {
if (this.maxRetryCount <= 0) {
log.warn("maxRetryCount [{}] 설정이 0 이하입니다. 1로 강제 조정합니다.", this.maxRetryCount);
this.maxRetryCount = Math.max(1, this.maxRetryCount);
}
}

// @Transactional을 제거 - AOP Self-Invocation 방지
public void reserveStockWithRetry(ReserveStockRequest request) {
int retryCount = 0;
while (retryCount < maxRetryCount) {
try {
// 1. 매 재시도마다 독립된 새로운 트랜잭션을 연다
StockReservedEvent event = transactionTemplate.execute(status -> reserveStockInternal(request));

// 2. TransactionTemplate이 에러 없이 끝나면 DB 커밋이 완료된 것
// Fire-and-forget 방지. whenComplete 콜백을 달아 전송 실패 시 인지 및 후속 처리 가능하도록 수정
kafkaTemplate.send(topicStockReserved, request.optionId().toString(), event)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("Kafka 메시지 발행 실패 (데이터 불일치 위험)! optionId: {}", request.optionId(), ex);
// TODO 추후 재처리 로직이나 아웃박스 패턴으로 고도화 필요
} else {
long offset =
(result != null && result.getRecordMetadata() != null) ? result.getRecordMetadata()
.offset() : -1;
log.info("Kafka 예약 메시지 발행 성공! offset: {}", offset);
}
});
return;

} catch (ObjectOptimisticLockingFailureException e) {
retryCount++;
log.warn("재고 차감 동시성(낙관적 락) 충돌 발생. 재시도 횟수: {}/{}", retryCount, maxRetryCount);
if (retryCount >= maxRetryCount) {
throw new ConcurrencyFailureException();
}
backoff("차감");
}
}
}

// 트랜잭션 내부에서 실행될 순수 비즈니스 로직
private StockReservedEvent reserveStockInternal(ReserveStockRequest request) {
// while문, Thread.sleep, TransactionTemplate 모두 제거 및 순수 비즈니스 로직만 남김
@Transactional
public void reserveStock(ReserveStockRequest request) {
Stock stock = stockRepository.findById(request.optionId())
.orElseThrow(StockNotFoundException::new);

// 1. 도메인 로직을 통한 3중 검증 및 차감
stock.reserve(request.quantity());

// 2. 가독성 위해... 어차피 더티 체킹에 의해 flush 시점에 버전 체크가 발생함
stockRepository.save(stock);

// 품절 상태 자동 전이 로직
// 2. 품절 상태 자동 전이 로직
if (stock.getTotalQuantity() == 0) {
ProductOption option = productOptionRepository.findById(stock.getOptionId())
.orElseThrow(() -> new IllegalArgumentException("옵션 정보를 찾을 수 없습니다."));
Expand All @@ -121,73 +67,26 @@ private StockReservedEvent reserveStockInternal(ReserveStockRequest request) {

ProductStatusChangedEvent statusEvent = new ProductStatusChangedEvent(product.getId(),
product.getStatus().name());

if (TransactionSynchronizationManager.isSynchronizationActive()) {
TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
@Override
public void afterCommit() {
// 카프카 전송 실패 시 에러 핸들링 로깅 추가
kafkaTemplate.send(topicStatusChanged, product.getId().toString(), statusEvent)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("SOLDOUT 상태 변경 Kafka 메시지 발행 실패 (데이터 불일치 위험)! productId: {}",
product.getId(), ex);
} else {
long offset = (result != null && result.getRecordMetadata() != null)
? result.getRecordMetadata().offset() : -1;
log.info("SOLDOUT 상태 변경 Kafka 메시지 발행 성공! productId={}, offset={}",
product.getId(), offset);
}
});
}
});
}
publishKafkaEvent(
topicStatusChanged,
product.getId().toString(),
statusEvent,
"SOLDOUT 상태 변경"
);
}
}

// 3. Kafka 이벤트 발행
return new StockReservedEvent(
// 3. Kafka 이벤트 발행 (DB 커밋 이후에 실행되도록 공통 메서드 분리)
StockReservedEvent event = new StockReservedEvent(
stock.getOptionId(),
stock.getTotalQuantity(),
stock.getCurrentDailyStock()
);
publishKafkaEvent(topicStockReserved, request.optionId().toString(), event, "재고차감");
}

//복구로직
public void restoreStockWithRetry(RestoreStockRequest request) {
int retryCount = 0;
while (retryCount < maxRetryCount) {
try {
// 1. 매 재시도마다 독립된 새로운 트랜잭션을 연다
StockRestoredEvent event = transactionTemplate.execute(status -> restoreStockInternal(request));

// 2. DB 커밋 완료 후 Kafka 발행
kafkaTemplate.send(topicStockRestored, request.optionId().toString(), event)
.whenComplete((result, ex) -> {
if (ex != null) {
log.error("Kafka 복구 메시지 발행 실패 (재고 유실 위험)! optionId: {}", request.optionId(), ex);
} else {
long offset =
(result != null && result.getRecordMetadata() != null) ? result.getRecordMetadata()
.offset() : -1;
log.info("Kafka 복구 메시지 발행 성공! offset: {}", offset);
}
});
return;

} catch (ObjectOptimisticLockingFailureException e) {
retryCount++;
log.warn("재고 복구 동시성(낙관적 락) 충돌 발생. 재시도 횟수: {}/{}", retryCount, maxRetryCount);
if (retryCount >= maxRetryCount) {
throw new ConcurrencyFailureException();
}
backoff("복구");
}
}
}

// 트랜잭션 내부에서 실행될 순수 비즈니스 로직
private StockRestoredEvent restoreStockInternal(RestoreStockRequest request) {
@Transactional
public void restoreStock(RestoreStockRequest request) {
Stock stock = stockRepository.findById(request.optionId())
.orElseThrow(StockNotFoundException::new);

Expand All @@ -197,20 +96,40 @@ private StockRestoredEvent restoreStockInternal(RestoreStockRequest request) {
// 2. DB 업데이트 (더티 체킹 후 flush)
stockRepository.save(stock);

// 3. Kafka 이벤트 발행용 객체 리턴
return new StockRestoredEvent(
// 2. Kafka 이벤트 발행
StockRestoredEvent event = new StockRestoredEvent(
stock.getOptionId(),
stock.getTotalQuantity(),
stock.getCurrentDailyStock()
);
publishKafkaEvent(topicStockRestored, request.optionId().toString(), event, "재고복구");
}

private void backoff(String operationType) {
try {
Thread.sleep(50);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new IllegalStateException("재고 " + operationType + " 대기 중 쓰레드가 강제 종료되었습니다.", ie);
// 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);
}
});
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
package com.michelet.inventory.application;

import com.michelet.inventory.presentation.dto.ReserveStockRequest;
import com.michelet.inventory.presentation.dto.RestoreStockRequest;
import java.util.concurrent.TimeUnit;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.redisson.api.RLock;
import org.redisson.api.RedissonClient;
import org.springframework.stereotype.Component;

@Slf4j
@Component
@RequiredArgsConstructor
public class StockLockFacade {

private final RedissonClient redissonClient;
private final StockCommandService stockCommandService; // 기존 트랜잭션 서비스

// 1. 재고 선점 (Order 생성 시 호출)
public void reserveStockWithLock(ReserveStockRequest request) {
RLock lock = redissonClient.getLock("stock:" + request.optionId());

try {
// 락 획득 시도
// 매개변수: 최대 5초 대기, 로직이 끝날 때까지 5초마다 락 만료 시간을 계속 연장
boolean available = lock.tryLock(5, TimeUnit.SECONDS);
if (!available) {
log.error("[StockLockFacade] 재고 선점 락 획득 실패 - OptionId: {}", request.optionId());
throw new IllegalStateException("재고 선점 처리 중 락을 획득하지 못했습니다.");
}

// 락 획득 성공 시 실제 비즈니스 로직(트랜잭션) 실행
stockCommandService.reserveStock(request);

} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("재고 선점 락 대기 중 인터럽트 발생");
} finally {
if (lock != null && lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}

// 2. 재고 복구 (Order 실패/취소 시 보상 트랜잭션으로 호출)
public void restoreStockWithLock(RestoreStockRequest request) {
RLock lock = redissonClient.getLock("stock:" + request.optionId());

try {
// 복구 로직은 선점보다 더 중요하므로 대기 시간을 넉넉히(일단 10초) 줌
boolean available = lock.tryLock(10, TimeUnit.SECONDS);
if (!available) {
log.error("[CRITICAL] 재고 복구 락 획득 실패 (수동 복구 필요) - OptionId: {}", request.optionId());
throw new IllegalStateException("재고 복구 처리 중 락을 획득하지 못했습니다.");
}

// 락 획득 성공 시 복구 비즈니스 로직 실행
stockCommandService.restoreStock(request);

} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("재고 복구 락 대기 중 인터럽트 발생");
} finally {
if (lock != null && lock.isHeldByCurrentThread()) {
lock.unlock();
}
}
}
}
18 changes: 2 additions & 16 deletions src/main/java/com/michelet/inventory/domain/model/Stock.java
Original file line number Diff line number Diff line change
Expand Up @@ -9,19 +9,17 @@
import jakarta.persistence.Entity;
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.springframework.data.domain.Persistable;

@Entity(name = "p_stocks")
@Table(name = "p_stocks")
@Getter
@NoArgsConstructor(access = AccessLevel.PROTECTED)
public class Stock extends BaseEntity implements Persistable<UUID> {
public class Stock extends BaseEntity {

@Id
private UUID optionId; // ProductOption의 ID를 PK로 사용
Expand All @@ -38,8 +36,7 @@ public class Stock extends BaseEntity implements Persistable<UUID> {
@Column(nullable = false)
private Integer maxLimit;

@Version // 낙관적 락용 버전
private Long version;
// Redisson 분산 락을 사용하므로 @Version(낙관적 락) 필드 삭제

@Builder(access = AccessLevel.PRIVATE)
private Stock(UUID optionId, Integer totalQuantity, Integer dailyLimit, Integer maxLimit) {
Expand Down Expand Up @@ -94,17 +91,6 @@ private void validateReserve(Integer requestQuantity) {
}
}

@Override
public UUID getId() {
return optionId;
}

@Override
public boolean isNew() {
// createdAt은 DB 저장 후에 채워지므로, 객체 생성 직후엔 version(null)으로 판단하는 것이 더 정확함
return version == null;
}

// 복구 로직
public void restore(int quantity) {
if (quantity <= 0) {
Expand Down
Loading