Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
6 changes: 6 additions & 0 deletions apps/commerce-api/build.gradle.kts
Original file line number Diff line number Diff line change
@@ -1,7 +1,10 @@
dependencies {
implementation(project(":apps:commerce-core"))

// add-ons
implementation(project(":modules:jpa"))
implementation(project(":modules:redis"))
implementation(project(":modules:kafka"))
implementation(project(":supports:jackson"))
implementation(project(":supports:logging"))
implementation(project(":supports:monitoring"))
Expand Down Expand Up @@ -29,7 +32,10 @@ dependencies {
// test-fixtures
testImplementation(testFixtures(project(":modules:jpa")))
testImplementation(testFixtures(project(":modules:redis")))
testImplementation(testFixtures(project(":modules:kafka")))
testImplementation("com.github.javafaker:javafaker:1.0.2") {
exclude(group = "org.yaml", module = "snakeyaml")
}
// Kafka 테스트용
testImplementation("org.awaitility:awaitility:4.2.0")
}
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,26 @@
import jakarta.annotation.PostConstruct;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
import org.springframework.boot.context.properties.ConfigurationPropertiesScan;
import org.springframework.cloud.openfeign.EnableFeignClients;
import org.springframework.scheduling.annotation.EnableScheduling;

import java.util.TimeZone;

@ConfigurationPropertiesScan
@SpringBootApplication
@SpringBootApplication(
scanBasePackages = {
"com.loopers.application", // api 앱의 application 레이어
"com.loopers.infrastructure", // api 앱의 infrastructure (OutboxPublisher, PG 등)
"com.loopers.interfaces", // api 앱의 인터페이스 (Controller 등)
"com.loopers.domain", // 🆕 core 앱의 도메인 스캔 (Service, Repository 인터페이스 등)
"com.loopers.config" // JPA 모듈의 JpaConfig, DataSourceConfig 등 설정 클래스 스캔
},
exclude = {DataSourceAutoConfiguration.class} // 커스텀 DataSource 설정 사용을 위해 자동 설정 제외
)
// @EntityScan은 JpaConfig에서 이미 com.loopers 전체를 스캔하므로 불필요
// @EnableJpaRepositories도 JpaConfig에서 이미 com.loopers.infrastructure를 스캔하므로 불필요
@EnableFeignClients
@EnableScheduling
public class CommerceApiApplication {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -166,10 +166,14 @@ public void handlePgPaymentCallback(String orderId, PaymentV1Dto.PaymentCallback
log.info("PG 상태 확인 성공 - orderId: {}, matched status: {}", orderId, matchedTransaction.status());

// 결제 상태 업데이트 (이벤트 발행은 OrderService.updateStatusByPgResponse에서 처리)
orderService.updateStatusByPgResponse(
order.getId(),
matchedTransaction
);
// PgV1Dto를 도메인 VO로 변환
com.loopers.domain.order.vo.PgTransactionResponse domainResponse =
new com.loopers.domain.order.vo.PgTransactionResponse(
matchedTransaction.transactionKey(),
matchedTransaction.status(),
matchedTransaction.reason()
);
orderService.updateStatusByPgResponse(order.getId(), domainResponse);
}

@Transactional
Expand Down Expand Up @@ -207,7 +211,14 @@ public OrderInfo payOrder(String userId, Long orderId, OrderRequest.PaymentReque

pgPaymentExecutor.requestPaymentAsync(pgRequest)
.thenAccept(pgResponse -> {
orderService.updateStatusByPgResponse(orderId, pgResponse);
// PgV1Dto를 도메인 VO로 변환
com.loopers.domain.order.vo.PgTransactionResponse domainResponse =
new com.loopers.domain.order.vo.PgTransactionResponse(
pgResponse.transactionKey(),
pgResponse.status(),
pgResponse.reason()
);
orderService.updateStatusByPgResponse(orderId, domainResponse);
log.info("재결제 요청 성공 - orderId: {}, pgResponse: {}", orderId, pgResponse);
})
.exceptionally(throwable -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,14 @@ private void retryPaymentRequest(Order order) {
.thenAccept(response -> {
if (response.transactionKey() != null) {
// 재결제 성공
orderService.updateStatusByPgResponse(order.getId(), response);
// PgV1Dto를 도메인 VO로 변환
com.loopers.domain.order.vo.PgTransactionResponse domainResponse =
new com.loopers.domain.order.vo.PgTransactionResponse(
response.transactionKey(),
response.status(),
response.reason()
);
orderService.updateStatusByPgResponse(order.getId(), domainResponse);
Comment on lines +109 to +116

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🛠️ Refactor suggestion | 🟠 Major

중복된 변환 로직을 헬퍼 메서드로 추출하세요.

PgV1Dto.PgTransactionResponse를 도메인 VO로 변환하는 동일한 로직이 두 곳(retryPaymentRequest와 syncOrderStatusFromPgResponse)에 중복되어 있습니다. DRY 원칙을 위반하며, 향후 변환 로직 변경 시 두 곳을 모두 수정해야 하는 유지보수 리스크가 있습니다.

🔎 변환 로직을 헬퍼 메서드로 추출하는 제안
+    /**
+     * PgV1Dto를 도메인 VO로 변환
+     */
+    private com.loopers.domain.order.vo.PgTransactionResponse convertToDomainVo(
+            PgV1Dto.PgTransactionResponse response) {
+        return new com.loopers.domain.order.vo.PgTransactionResponse(
+                response.transactionKey(),
+                response.status(),
+                response.reason()
+        );
+    }
+
     private void retryPaymentRequest(Order order) {
         log.info("재결제 시도 시작 - orderId: {}", order.getOrderId());
         PgV1Dto.PgPaymentRequest request = buildPaymentRequestFromOrder(order);
 
         pgPaymentExecutor.requestPaymentAsync(request)
                 .thenAccept(response -> {
                     if (response.transactionKey() != null) {
-                        // PgV1Dto를 도메인 VO로 변환
-                        com.loopers.domain.order.vo.PgTransactionResponse domainResponse = 
-                                new com.loopers.domain.order.vo.PgTransactionResponse(
-                                        response.transactionKey(),
-                                        response.status(),
-                                        response.reason()
-                                );
+                        com.loopers.domain.order.vo.PgTransactionResponse domainResponse = 
+                                convertToDomainVo(response);
                         orderService.updateStatusByPgResponse(order.getId(), domainResponse);
                         ...
 
     private void syncOrderStatusFromPgResponse(Order order, PgV1Dto.PgTransactionResponse response) {
-        // PgV1Dto를 도메인 VO로 변환
-        com.loopers.domain.order.vo.PgTransactionResponse domainResponse = 
-                new com.loopers.domain.order.vo.PgTransactionResponse(
-                        response.transactionKey(),
-                        response.status(),
-                        response.reason()
-                );
+        com.loopers.domain.order.vo.PgTransactionResponse domainResponse = 
+                convertToDomainVo(response);
         orderService.updateStatusByPgResponse(order.getId(), domainResponse);
         ...

Also applies to: 148-156

log.info("재결제 성공 - orderId: {}, transactionKey: {}",
order.getOrderId(), response.transactionKey());
} else {
Expand Down Expand Up @@ -139,7 +146,14 @@ private PgV1Dto.PgPaymentRequest buildPaymentRequestFromOrder(Order order) {
* PG 응답으로 주문 상태 동기화
*/
private void syncOrderStatusFromPgResponse(Order order, PgV1Dto.PgTransactionResponse response) {
orderService.updateStatusByPgResponse(order.getId(), response);
// PgV1Dto를 도메인 VO로 변환
com.loopers.domain.order.vo.PgTransactionResponse domainResponse =
new com.loopers.domain.order.vo.PgTransactionResponse(
response.transactionKey(),
response.status(),
response.reason()
);
orderService.updateStatusByPgResponse(order.getId(), domainResponse);
log.info("PG 상태 동기화 성공 - orderId: {}, transactionKey: {}, status: {}",
order.getOrderId(), response.transactionKey(), response.status());
}
Expand All @@ -153,13 +167,13 @@ private void recoverOrderPaymentByTransactionKey(Order order) {

pgPaymentExecutor.getPaymentDetailByTransactionKeyAsync(order.getPgTransactionKey())
.thenAccept(transactionDetail -> {
// PgTransactionDetailResponse를 PgTransactionResponse로 변환
PgV1Dto.PgTransactionResponse transactionResponse = new PgV1Dto.PgTransactionResponse(
// PgTransactionDetailResponse를 PgV1Dto로 변환 후 도메인 VO로 변환
PgV1Dto.PgTransactionResponse pgResponse = new PgV1Dto.PgTransactionResponse(
transactionDetail.transactionKey(),
transactionDetail.status(),
transactionDetail.pgPaymentReason()
);
syncOrderStatusFromPgResponse(order, transactionResponse);
syncOrderStatusFromPgResponse(order, pgResponse);
})
.exceptionally(e -> {
log.warn("결제 상태 복구 실패 - orderId: {}, transactionKey: {}",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,8 +72,15 @@ public void handleOrderCreated(OrderEvent.OrderCreatedEvent event) {
// PG 결제 요청 (비동기)
pgPaymentExecutor.requestPaymentAsync(pgRequest)
.thenAccept(pgResponse -> {
// PgV1Dto를 도메인 VO로 변환
com.loopers.domain.order.vo.PgTransactionResponse domainResponse =
new com.loopers.domain.order.vo.PgTransactionResponse(
pgResponse.transactionKey(),
pgResponse.status(),
pgResponse.reason()
);
// 주문 상태 업데이트
orderService.updateStatusByPgResponse(event.getOrderId(), pgResponse);
orderService.updateStatusByPgResponse(event.getOrderId(), domainResponse);
log.info("PG 결제 요청 성공 - orderId: {}, transactionKey: {}", event.getOrderId(), pgResponse.transactionKey());
})
.exceptionally(throwable -> {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
package com.loopers.infrastructure.event.outbox;

import com.loopers.domain.event.outbox.OutboxEvent;
import com.loopers.domain.event.outbox.OutboxEventRepository;
import com.loopers.domain.event.outbox.OutboxStatus;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.slf4j.Marker;
import org.slf4j.MarkerFactory;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.scheduling.annotation.Scheduled;
import org.springframework.stereotype.Component;
import org.springframework.transaction.annotation.Transactional;

import java.util.List;

@Slf4j
@Component
@RequiredArgsConstructor
public class OutboxPublisher {
// fatal marker
private static Marker FATAL = MarkerFactory.getMarker("FATAL");
private final OutboxEventRepository outboxEventRepository;
private final KafkaTemplate<Object, Object> kafkaTemplate; // KafkaConfig에서 제공하는 Bean 사용

private static final int BATCH_SIZE = 100;
private static final int MAX_RETRY_COUNT = 5;

/**
* 주기적으로 PENDING 상태의 Outbox 이벤트를 Kafka로 발행
* 5초마다 실행 (필요에 따라 조정)
*
* ⚠️ 토픽별로 Relay를 분리하는 것을 권장합니다:
* - order-events: 1초 주기 (빠른 전달 필요)
* - catalog-events: 5초 주기 (느슨해도 괜찮음)
*/
@Scheduled(fixedDelay = 5000)
@Transactional
public void publishPendingEvents() {
List<OutboxEvent> pendingEvents = outboxEventRepository.findByStatusOrderByCreatedAtAsc(OutboxStatus.PENDING);

if (pendingEvents.isEmpty()) {
return;
}

log.info("Found {} pending events to publish", pendingEvents.size());

int processedCount = 0;
for (OutboxEvent event : pendingEvents) {
if (processedCount >= BATCH_SIZE) {
break;
}

String topic = determineTopic(event.getAggregateType());
String key = event.getAggregateId();

// Kafka로 메시지 발행
kafkaTemplate.send(topic, key, event.getPayload()).whenComplete((result, exception) -> {
if (exception == null) {
event.markAsPublished();
outboxEventRepository.save(event);
log.debug("Published event: eventId={}, topic={}, key={}", event.getId(), topic, key);
} else {
handlePublishFailure(event, exception);
}
});

processedCount++;
}
}
Comment on lines +37 to +70

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🔴 Critical

동시성 제어 부재로 인한 중복 발행 가능성

여러 애플리케이션 인스턴스가 동시에 실행될 때, 동일한 PENDING 이벤트를 여러 인스턴스가 동시에 조회하고 발행할 수 있습니다.

findByStatusOrderByCreatedAtAsc()는 비관적 락을 사용하지 않으므로:

  • 인스턴스 A와 B가 동시에 같은 이벤트 조회
  • 둘 다 Kafka로 전송
  • 중복 메시지 발생 (At Least Once를 넘어 중복 발생)

해결 방법:

  1. findPendingEventsForUpdate()처럼 비관적 락 사용
  2. 또는 이벤트 조회 시 인스턴스 식별자로 필터링 (샤딩)
  3. 또는 Kafka에서 멱등성으로 중복 제거 (producer idempotence + deduplication)
🔎 제안하는 수정안
     @Scheduled(fixedDelay = 5000)
     @Transactional
     public void publishPendingEvents() {
-        List<OutboxEvent> pendingEvents = outboxEventRepository.findByStatusOrderByCreatedAtAsc(OutboxStatus.PENDING);
+        // 비관적 락을 사용하여 동시성 제어
+        List<OutboxEvent> pendingEvents = outboxEventRepository.findPendingEventsForUpdate(OutboxStatus.PENDING);

또는 Kafka Producer에 멱등성 설정 추가:

spring:
  kafka:
    producer:
      enable-idempotence: true
      acks: all

Committable suggestion skipped: line range outside the PR's diff.


⚠️ Potential issue | 🔴 Critical

🧩 Analysis chain

🏁 Script executed:

find . -type f -name "*OutboxEventRepository*" -o -name "*OutboxEvent*" | grep -E "\.(java|kt)$" | head -20

Repository: Loopers-dev-lab/loopers-spring-java-template

Length of output: 567


🏁 Script executed:

cat -n ./apps/commerce-core/src/main/java/com/loopers/domain/event/outbox/OutboxEventRepository.java

Repository: Loopers-dev-lab/loopers-spring-java-template

Length of output: 477


🏁 Script executed:

cat -n ./apps/commerce-core/src/main/java/com/loopers/infrastructure/event/outbox/OutboxEventRepositoryImpl.java

Repository: Loopers-dev-lab/loopers-spring-java-template

Length of output: 1269


🏁 Script executed:

cat -n ./apps/commerce-core/src/main/java/com/loopers/infrastructure/event/outbox/OutboxEventJpaRepository.java

Repository: Loopers-dev-lab/loopers-spring-java-template

Length of output: 1310


비동기 Kafka 콜백에서의 트랜잭션 경계 문제

publishPendingEvents() 메서드의 @Transactional이 트랜잭션을 시작하지만, kafkaTemplate.send()의 whenComplete 콜백은 원본 트랜잭션 커밋 후에 실행됩니다. 콜백 내에서 outboxEventRepository.save(event)를 호출하면, JpaRepository의 기본 동작(PROPAGATION.REQUIRED)에 따라 새로운 독립 트랜잭션에서 실행되어 원자성이 깨집니다.

문제:

  1. 이벤트 발행 성공 후 상태 업데이트가 실패할 수 있음
  2. 상태 업데이트 성공 후 이벤트 발행 실패 시에도 상태는 업데이트됨
  3. Outbox 패턴의 핵심인 발행과 상태 업데이트의 원자성 보장 불가

해결 방법:

  1. 콜백 내에서 TransactionTemplate을 사용해 새 트랜잭션 명시적 시작
  2. 또는 TransactionSynchronizationManager.registerSynchronization()으로 커밋 후 실행 등록
  3. 또는 상태 업데이트를 별도 스케줄러/폴링 메커니즘으로 분리


private String determineTopic(String aggregateType) {
return switch (aggregateType) {
case "PRODUCT" -> "catalog-events";
case "ORDER" -> "order-events";
default -> throw new IllegalArgumentException("Unknown aggregate type: " + aggregateType);
};
}

private void handlePublishFailure(OutboxEvent event, Throwable exception) {
log.error("Failed to publish event: eventId={}, retryCount={}",
event.getId(), event.getRetryCount(), exception);

if (event.getRetryCount() >= MAX_RETRY_COUNT) {
event.markAsFailed(exception.getMessage());
outboxEventRepository.save(event);
log.error(FATAL, "Event exceeded max retry count, marked as FAILED: eventId={}", event.getId());
// TODO: DLQ로 전송하거나 알림 발송
} else {
event.incrementRetry();
outboxEventRepository.save(event);
}
}
}

Original file line number Diff line number Diff line change
@@ -1,25 +1,78 @@
package com.loopers.infrastructure.like;

import com.loopers.domain.event.outbox.OutboxEventService;
import com.loopers.domain.like.product.LikeProductEvent;
import com.loopers.domain.like.product.LikeProductEventPublisher;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.context.ApplicationEventPublisher;
import org.springframework.stereotype.Component;

import java.util.HashMap;
import java.util.Map;
import java.util.UUID;

/**
* 좋아요 이벤트를 Kafka로 발행하기 위한 Publisher
*
* 변경 사항:
* - 기존: ApplicationEventPublisher를 통한 동기식 이벤트 발행
* - 변경: Outbox 패턴을 통한 비동기 Kafka 이벤트 발행
*
* 처리 플로우:
* 1. LikeProductEvent를 받아서 Outbox 테이블에 저장 (같은 트랜잭션)
* 2. OutboxPublisher가 주기적으로 Kafka로 발행
* 3. Consumer가 Kafka에서 수신하여 Metrics 집계
*/
@Component
@RequiredArgsConstructor
@Slf4j
public class LikeProductEventPublisherImpl implements LikeProductEventPublisher {
private final ApplicationEventPublisher applicationEventPublisher;
private final OutboxEventService outboxEventService;

@Override
public void publishLikeEvent(LikeProductEvent event) {
if (event == null) {
log.warn("LikeProductEvent가 null입니다. 이벤트 발행을 건너뜁니다.");
return;
}
log.info("LikeProductEvent 발행: {}", event);
applicationEventPublisher.publishEvent(event);

try {
// 이벤트 페이로드 생성
Map<String, Object> eventPayload = createEventPayload(event);

// 이벤트 타입 결정
String eventType = event.getLiked() ? "ProductLiked" : "ProductUnliked";

// Outbox 테이블에 저장 (도메인 트랜잭션 내에서 호출되어야 함)
outboxEventService.saveEvent(
"PRODUCT",
event.getProductId().toString(),
eventType,
eventPayload
);

log.info("LikeProductEvent saved to outbox: productId={}, liked={}, eventType={}",
event.getProductId(), event.getLiked(), eventType);
} catch (Exception e) {
log.error("Failed to save LikeProductEvent to outbox: productId={}, liked={}",
event.getProductId(), event.getLiked(), e);
throw e;
}
}

/**
* LikeProductEvent를 Kafka 이벤트 페이로드로 변환
*/
private Map<String, Object> createEventPayload(LikeProductEvent event) {
Map<String, Object> payload = new HashMap<>();
payload.put("eventId", UUID.randomUUID().toString());
payload.put("eventType", event.getLiked() ? "ProductLiked" : "ProductUnliked");
payload.put("aggregateId", event.getProductId().toString());
payload.put("productId", event.getProductId());
payload.put("userId", event.getUserId());
payload.put("brandId", event.getBrandId());
payload.put("liked", event.getLiked());
payload.put("createdAt", java.time.Instant.now().toString());
return payload;
}
}

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ public interface PgClient {
ApiResponse<PgV1Dto.PgTransactionResponse> requestPayment(@RequestBody PgV1Dto.PgPaymentRequest request);

@GetMapping("/api/v1/payments/{transactionKey}")
ApiResponse<PgV1Dto.PgTransactionDetailResponse> getTransactionDetail(@PathVariable String transactionKey);
ApiResponse<PgV1Dto.PgTransactionDetailResponse> getTransactionDetail(@PathVariable("transactionKey") String transactionKey);

@GetMapping(path = "/api/v1/payments", params = "orderId")
ApiResponse<PgV1Dto.PgOrderResponse> getPaymentByOrderId(@RequestParam("orderId") String orderId);
Expand Down
3 changes: 3 additions & 0 deletions apps/commerce-api/src/main/resources/application.yml
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ spring:
- redis.yml
- logging.yml
- monitoring.yml
- kafka.yml

springdoc:
use-fqn: true
Expand Down Expand Up @@ -90,6 +91,8 @@ feign:
connect-timeout: 1000
read-timeout: 3000
logger-level: basic


---
spring:
config:
Expand Down
Loading