order-system 모놀리식 아키텍처 구조 파일들을 서비스별로 분리해서 생성
api-gateway, eureka, member, product, ordering불필요한 의존성과 config 파일 제거 (RabbitMQ, Security 등)
cf) Security 의존성 역할
Spring Security가 Authentication 객체를 요구하며, 토큰 검증 후 사용자 정보(이메일, 권한 등)를 담은 객체를 생성
- 예정 설계 변경
- 기존: 게이트웨이 + Eureka에서 토큰 검증 → Authentication 객체 생성
- 변경:
- 게이트웨이에서 토큰 유효성만 판별 (정상/비정상)
- 서버 측에서는 이메일, role 추출 불가 (Authentication 객체 미생성)
- 주요 문제점
- 보안 취약: 사용자 토큰 없이 직접 서버 API 호출 가능 (게이트웨이 우회)
- 운영 환경에서의 대응
환경 접근 제어 방법 이유 개발 직접 API 호출 허용 편의성 운영 외부 접근 완전 차단 (게이트웨이 강제 경유) 보안 강화
- 운영 배포 세팅 추천
# application-prod.yml 예시 (Nginx/Cloud 설정) server: address: 127.0.0.1 # 로컬호스트만 접근 허용 # 또는 클라우드 보안그룹: 게이트웨이 IP만 화이트리스트
- 방화벽/클라우드 보안그룹: 외부 IP 차단, 게이트웨이만 허용
- Spring Security 추가 레이어: 게이트웨이 우회 시 401/403 반환
→ 이렇게 하면 토큰 없이 서버 직격 불가능
주요 역할: 라우팅, CORS, 토큰 검증
dependencies {
implementation 'io.jsonwebtoken:jjwt-api:0.11.5'
runtimeOnly 'io.jsonwebtoken:jjwt-impl:0.11.5'
runtimeOnly 'io.jsonwebtoken:jjwt-jackson:0.11.5'
// eureka client 의존성 추가
implementation 'org.springframework.cloud:spring-cloud-starter-netflix-eureka-client'
// gateway 의존성 추가
implementation 'org.springframework.cloud:spring-cloud-starter-gateway'
}
// Spring Cloud Version에 대한 별도 설정 필요
dependencyManagement {
imports {
mavenBom 'org.springframework.cloud:spring-cloud-dependencies:2024.0.0'
}
}
server:
port: 8080
eureka:
client:
serviceUrl:
defaultZone: http://localhost:8761/eureka/
spring:
application:
name: api-gateway
cloud:
gateway:
# CORS 처리. 서비스 모듈에서는 웹브라우저와 직접 통신하지 않으므로, CORS 처리 불필요.
globalcors:
cors-configurations:
'[/**]':
allowedOrigins:
- 'http://localhost:3000'
# - 'https://www.example.com'
# - 'https://www.jiyean199.shop'
allowedMethods: '*'
allowedHeaders: '*'
allowedCredentials: true
routes:
- id: member-service
predicates:
- Path=/member-service/**
filters:
# StripPrefix : 첫 번째 접두어를 제거 후 member-service로 HTTP 요청 전달
- StripPrefix=1
# eureka에 등록된 member-service라는 이름으로 HTTP 요청 라우팅
uri: lb://member-service
- id: ordering-service
predicates:
- Path=/ordering-service/**
filters:
- StripPrefix=1
uri: lb://ordering-service
- id: product-service
predicates:
- Path=/product-service/**
filters:
- StripPrefix=1
uri: lb://product-service
jwt:
secretKey: jwt 시크릿 키
CORS 검증 관련해서
api-gateway는 톰캣 기반 스프링부트가 아니고 Netty 엔진 기반 스프링부트이기 때문에, 일반적으로 사용하는 필터 계층에서 CORS 처리하지 않는다.
- Netty 엔진 특징: 비동기 프레임워크
- api-gateway는 모든 프로젝트의 진입점이라 성능 향상을 위해 비동기 프레임워크(Netty)를 사용
routes 설계:
member-service로 시작하는 URL을 받으면 eureka에게 member-service를 질의하겠다는 세팅
jwt 세팅:
멤버 서버에서 토큰을 생성하고, api-gateway에서 토큰을 검증한다. 이 때 시크릿키가 두 위치 모두에 필요하다.
RT는 gateway에서 검증하지 않는다. RT는 member에서 그대로 검증하고, RT로 AT 갱신받는 로직도 member에 있는 refresh 토큰 API에서 처리.
@Component
public class JwtAuthFilter implements GlobalFilter {
@Value("${jwt.secretKeyAt}")
private String secretKey;
private static final List<String> ALLOWED_PATH = List.of(
"/member/create",
"/member/doLogin",
"/member/refresh-at",
"/product/list"
);
private static final List<String> ADMIN_ONLY_PATH = List.of(
"/member/list",
"/member/detail/**",
"/ordering/list",
"/product/create"
);
@Override
public Mono<Void> filter(ServerWebExchange exchange, GatewayFilterChain chain) {
String bearerToken = exchange.getRequest().getHeaders().getFirst(HttpHeaders.AUTHORIZATION);
// StripPrefix로 접두어를 떼어냈기 때문에 /member/create 형태로 들어옴
String urlPath = exchange.getRequest().getURI().getRawPath();
// 인증이 필요 없는 경로는 필터 통과
if (ALLOWED_PATH.contains(urlPath)) {
return chain.filter(exchange);
}
try {
if (bearerToken == null || !bearerToken.startsWith("Bearer ")) {
throw new IllegalArgumentException("token이 없거나, 형식이 잘못되었습니다.");
}
String token = bearerToken.substring(7);
// token 검증 및 payload 추출
Claims claims = Jwts.parserBuilder()
.setSigningKey(secretKey)
.build()
.parseClaimsJws(token)
.getBody();
String email = claims.getSubject();
String role = claims.get("role", String.class);
// admin 권한 있어야 하는 url 검증
if (ADMIN_ONLY_PATH.contains(urlPath) && !role.equals("ADMIN")) {
exchange.getResponse().setStatusCode(HttpStatus.FORBIDDEN);
return exchange.getResponse().setComplete();
}
// header에 email, role 등 payload 값 세팅
// X- 접두어는 custom header라는 관례
// 서비스 모듈에서 @RequestHeader로 꺼내 사용 가능
ServerWebExchange serverWebExchange = exchange.mutate()
.request(r -> r.header("X-User-Email", email)
.header("X-User-Role", role))
.build();
return chain.filter(serverWebExchange);
} catch (Exception e) {
e.printStackTrace();
exchange.getResponse().setStatusCode(HttpStatus.UNAUTHORIZED);
return exchange.getResponse().setComplete();
}
}
}
CORS(CorsWebFilter) → 토큰 검증(GlobalFilter) → 라우팅 처리(GatewayFilter)
CORS 처리
웹브라우저와 직접 통신하지 않는 서비스 단에서 CORS를 처리하는 게 아니라, api-gateway에서 한 번에 처리
JWT 토큰 인증 처리
dependencies {
// eureka-server 의존성 추가
implementation 'org.springframework.cloud:spring-cloud-starter-netflix-eureka-server'
}
// Spring Cloud Version에 대한 별도 설정 필요
dependencyManagement {
imports {
mavenBom 'org.springframework.cloud:spring-cloud-dependencies:2024.0.0'
}
}
@EnableEurekaServer 어노테이션 추가server:
port: 8761
spring:
application:
name: eureka
eureka:
client:
# eureka 서버가 자기 자신은 eureka에 등록하지 않겠다는 설정
register-with-eureka: false
fetch-registry: false
dependencies {
// Spring Boot 기본 기능
implementation 'org.springframework.boot:spring-boot-starter-data-jpa'
implementation 'org.springframework.boot:spring-boot-starter-web'
implementation 'org.springframework.boot:spring-boot-starter-validation'
// DB 연결
runtimeOnly 'org.mariadb.jdbc:mariadb-java-client'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
// 보안 (JWT)
implementation 'io.jsonwebtoken:jjwt-api:0.11.5'
runtimeOnly 'io.jsonwebtoken:jjwt-impl:0.11.5'
runtimeOnly 'io.jsonwebtoken:jjwt-jackson:0.11.5'
// 캐싱
implementation 'org.springframework.boot:spring-boot-starter-data-redis'
// 마이크로서비스
implementation 'org.springframework.cloud:spring-cloud-starter-netflix-eureka-client'
// 테스트
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
}
dependencyManagement {
imports {
mavenBom 'org.springframework.cloud:spring-cloud-dependencies:2024.0.0'
}
}
server:
# port 0 → 사용 가능한 포트를 임의로 할당 (유레카에 등록될 포트)
port: 0
# 유레카 주소
eureka:
client:
serviceUrl:
defaultZone: http://localhost:8761/eureka/
spring:
# 유레카에 등록되는 이름
application:
name: member-service
datasource:
driver-class-name: org.mariadb.jdbc.Driver
url: jdbc:mariadb://localhost:3306/order_system_msa
username: root
password: RDB 비밀번호
jpa:
database: mysql
database-platform: org.hibernate.dialect.MariaDBDialect
generate-ddl: true
hibernate:
ddl-auto: create
show_sql: true
redis:
host: localhost
port: 6379
jwt:
secretKey: 엑세스 토큰 시크릿 키
expiration: 엑세스 토큰 만료 시간
secretKeyRt: 리프레시 토큰 시크릿 키
expirationRt: 리프레시 토큰 만료 시간
개발자는 각 서비스 포트번호를 알 필요가 없다.
게이트웨이가 유레카에 물어보고, member 서버 실행 시 유레카에 포트번호가 등록되기 때문에, 게이트웨이가 알아서 그 포트로 라우팅한다.
기존 코드: @AuthenticationPrincipal 사용 불가.
// 내 정보 조회
@GetMapping("/myinfo")
public ResponseEntity<?> myInfo(@AuthenticationPrincipal String principal) {
MyInfoResDto dto = memberService.myInfo(principal); // principal = "memberId"
return ResponseEntity.status(HttpStatus.OK).body(dto);
}
이제는 서버가 아니라 게이트웨이에서 토큰을 검증하고 payload에서 role, email을 꺼낸다.
api-gateway가 member로 HTTP 요청을 넘겨줄 때, header를 조작해서 role, email을 담은 커스텀 헤더를 내려보낸다.
개선 코드:
// 내 정보 조회
@GetMapping("/myinfo")
public ResponseEntity<?> myInfo(@RequestHeader("X-User-Email") String email) {
MyInfoResDto dto = memberService.myInfo(email);
return ResponseEntity.status(HttpStatus.OK).body(dto);
}
“X-”로 시작하는 헤더명은 개발자가 인위적으로 만든 Header인 경우 관례적으로 사용.
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-data-jpa'
implementation 'org.springframework.boot:spring-boot-starter-web'
implementation 'org.springframework.boot:spring-boot-starter-validation'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
runtimeOnly 'org.mariadb.jdbc:mariadb-java-client'
implementation 'software.amazon.awssdk:s3:2.29.50'
// redis
implementation 'org.springframework.boot:spring-boot-starter-data-redis'
// 마이크로서비스
implementation 'org.springframework.cloud:spring-cloud-starter-netflix-eureka-client'
// msa 간 통신을 위한 openfeign 의존성
implementation 'org.springframework.cloud:spring-cloud-starter-openfeign'
}
dependencyManagement {
imports {
mavenBom 'org.springframework.cloud:spring-cloud-dependencies:2024.0.0'
}
}
server:
port: 0
eureka:
client:
serviceUrl:
defaultZone: http://localhost:8761/eureka/
spring:
application:
name: ordering-service
datasource:
driver-class-name: org.mariadb.jdbc.Driver
url: jdbc:mariadb://localhost:3306/order_system_msa
username: root
password: RDB 비밀번호
jpa:
database: mysql
database-platform: org.hibernate.dialect.MariaDBDialect
generate-ddl: true
hibernate:
ddl-auto: create
show_sql: true
redis:
host: localhost
port: 6379
@Configuration
public class RestTemplateConfig {
@Bean
@LoadBalanced
// eureka에 등록된 서비스명을 사용하여 내부 서비스 호출(내부통신)하는 어노테이션
public RestTemplate makeRestTemplate() {
return new RestTemplate();
}
}
이 클래스를 정의한 이유는 단순히 빈 생성을 위한 목적이 아니라,
@LoadBalanced를 통해 eureka에 질의하여 서비스 모듈 간 통신을 하기 위해서다.
결국 uri: lb://product-service와 동일한 로직.
기존: api-gateway를 통한 호출
// (1) 재고 조회 요청 (동기요청: HTTP요청)
String endPoint = "http://localhost:8080/product-service/product/detail/" + itemDto.getProductId();
RestTemplate restTemplate = new RestTemplate();
ProductResDto product = restTemplate.exchange(endPoint, HttpMethod.GET, 헤더, ProductResDto.class);
if (product.getStockQuantity() < itemDto.getProductCount()) {
throw new IllegalArgumentException("재고가 부족합니다.");
}
변경: eureka에게 질의 후 product-service 직접 호출
// 상단에서 restTemplate DI
// (1) 재고 조회 요청 (동기요청: HTTP요청)
String endPoint = "http://product-service/product/detail/" + itemDto.getProductId();
ProductResDto product = restTemplate.exchange(endPoint, HttpMethod.GET, 헤더, ProductResDto.class);
if (product.getStockQuantity() < itemDto.getProductCount()) {
throw new IllegalArgumentException("재고가 부족합니다.");
}
이 때 서비스 모듈들은 같은 내부 네트워크(사설 네트워크)에 존재하고, 외부에서는 접근 불가능하지만 서비스끼리는 통신 가능한 구조.
public Long create(List<OrderItemCreateReqDto> items, String email) {
Ordering order = Ordering.builder()
.memberEmail(email)
.orderStatus(OrderStatus.ORDERED)
.build();
orderingRepository.save(order);
for (OrderItemCreateReqDto itemDto : items) {
// (1) 재고 조회 요청 (동기요청: HTTP요청)
// api gateway를 통한 호출 : http://localhost:8080/product-service/..
// eureka에게 질의 후 product-service 직접 호출 : http://product-service/ ..
String productStockEndPoint = "http://product-service/product/detail/" + itemDto.getProductId();
HttpHeaders productStockHeaders = new HttpHeaders();
HttpEntity<String> productStockHttpEntity = new HttpEntity<>(productStockHeaders);
ResponseEntity<ProductResDto> responseEntity =
restTemplate.exchange(productStockEndPoint, HttpMethod.GET, productStockHttpEntity, ProductResDto.class);
ProductResDto product = responseEntity.getBody();
if (product.getStockQuantity() < itemDto.getProductCount()) {
throw new IllegalArgumentException("재고가 부족합니다.");
}
// (2) 주문 발생
OrderingDetails orderingDetails = OrderingDetails.builder()
.ordering(order)
.productId(itemDto.getProductId())
.productName(product.getName())
.quantity(itemDto.getProductCount())
.build();
orderingDetailRepository.save(orderingDetails);
// (3) 재고 감소 요청 (동기/비동기 모두 가능)
String productStockDecreaseEndPoint = "http://product-service/product/decrease-stock";
HttpHeaders productStockDecreaseHeaders = new HttpHeaders();
productStockDecreaseHeaders.setContentType(MediaType.APPLICATION_JSON);
HttpEntity<OrderItemCreateReqDto> productStockDecreaseHttpEntity =
new HttpEntity<>(itemDto, productStockDecreaseHeaders);
// 재고 감소 요청 중 에러 발생 시 전체 롤백
restTemplate.exchange(productStockDecreaseEndPoint, HttpMethod.PUT, productStockDecreaseHttpEntity, Void.class);
}
return order.getId();
}
3단계 핵심:
Eureka + @LoadBalanced RestTemplate로
"product-service"라는 서비스명만 알고 통신하고 실제 포트/인스턴스는 eureka가 반환해준다.
서버 확장/축소 시 코드 수정 없이 자동 대응.
이 로직의 (3) 재고 감소 부분을 Kafka 비동기 처리로 개선해볼 것.
@FeignClient(name = "product-service")
public interface ProductFeignClient {
@GetMapping("/product/detail/{id}")
ProductResDto getProductById(@PathVariable("id") Long id);
@PutMapping("/product/decrease-stock")
void decreaseStockQuantity(@RequestBody OrderItemCreateReqDto dto);
}
FeignClient의 주요 장점
1) 재사용성 ↑: 다른 서비스에서도 동일 인터페이스 활용 가능
2) 가독성 ↑: 서비스명 + 엔드포인트 + HTTP 메서드가 한 눈에 보임
3) 편의성 ↑: 자동 응답 형변환 + 예외 처리 간편화
public Long createFeign(List<OrderItemCreateReqDto> items, String email) {
Ordering order = Ordering.builder()
.memberEmail(email)
.orderStatus(OrderStatus.ORDERED)
.build();
orderingRepository.save(order);
for (OrderItemCreateReqDto itemDto : items) {
// (1) 재고 조회 요청 (동기)
ProductResDto product = productFeignClient.getProductById(itemDto.getProductId());
if (product.getStockQuantity() < itemDto.getProductCount()) {
throw new IllegalArgumentException("재고가 부족합니다.");
}
// (2) 주문 발생
OrderingDetails orderingDetails = OrderingDetails.builder()
.ordering(order)
.productId(itemDto.getProductId())
.productName(product.getName())
.quantity(itemDto.getProductCount())
.build();
orderingDetailRepository.save(orderingDetails);
// (3) 재고 감소 요청
productFeignClient.decreaseStockQuantity(itemDto);
}
return order.getId();
}
기존 (동기 HTTP)
Order Service → Product Controller 호출
↓
Product Service 감소 로직 → DB
변경 (비동기 Kafka)
Order Service → Kafka 토픽으로 메시지 발행
↓
Product KafkaListener(Consumer) 수신 → Product Service 감소 로직 → DB
핵심:
Product Controller의 동기 재고 감소 API 제거,
Kafka Listener(Consumer)가 재고 감소 메서드를 트리거.
Order는 즉시 응답, Product는 비동기적으로 재고 감소 처리.
메시지 생산(발행)은 Ordering, 메시지 소비는 Product.
카프카에는 두 가지 요소:
생산자, 소비자 쪽 각각 설정이 필요하고, 보통 consumer 쪽 설정이 (Redis 리스너와 비슷하게) 조금 더 복잡하다.
루트 디렉토리 구조 예시:
order-system-msa/
├── docker-compose.yml (kafka)
├── ordering/ (producer)
│ └── application.yml
└── product/ (consumer)
└── application.yml
docker-compose.yml (단일 브로커, KRaft 모드)
version: '3.8'
services:
kafka:
image: apache/kafka:3.7.0
container_name: kafka
ports:
- "9092:9092"
environment:
# 단일 브로커
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
KAFKA_LOG_DIRS: /tmp/kraft-combined-logs
# 자동 토픽 생성 허용 (기본 true)
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
# 새 토픽 생성 시 기본 복제수 (단일 브로커이므로 1)
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
restart: unless-stopped
카프카는 원래 클러스터 구조로, 서버 여러 대를 묶어서 고가용성을 지향한다.
우리 실습에선 복잡도 줄이기 위해 단일 브로커로 세팅.
토픽은 메시지의 “주제” 개념 (Redis 채널과 비슷)
예: member 토픽, order 토픽 등
새 토픽 생성 시 기본 복제수(replication factor) 설정 필요.
단일 브로커이므로 replication_factor: 1로 설정.
(ordering, product 쪽 build.gradle에 spring-kafka 추가했다고 가정)
ordering (producer)
spring:
kafka:
kafka-server: localhost:9092
product (consumer)
spring:
kafka:
kafka-server: localhost:9092
consumer:
groud-id: product-group
(groud-id 오타는 실제 코드에서 group-id로 수정)
@Configuration
public class KafkaProducerConfig {
@Value("${spring.kafka.kafka-server}")
private String kafkaServer;
@Bean
public ProducerFactory<String, Object> producerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaServer);
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
return new DefaultKafkaProducerFactory<>(config);
}
@Bean
public KafkaTemplate<String, Object> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
여기서 producerFactory는 카프카 서버 연결 위치가 중요한 “연결 객체”.
Kafka 메시지는 기본적으로 key-value 쌍으로 구성된다.
key는 써도 되고 안 써도 된다.
우리 실습에선 key를 null로 두고 value만 보낸다.
실전에서 채팅 서비스를 카프카로 붙일 경우, roomId 같은 값을 key로 사용해서 동일 채팅방 메시지를 같은 파티션에 보내야 하므로 key를 반드시 쓴다.
우리는 value만 보고, 여기에 JSON으로 넣기 때문에 JsonSerializer로 설정.
StringSerializer를 쓰면 JSON을 직접 ObjectMapper로 직렬화해서 String으로 만들어야 하는 차이가 있다.
Redis는 JSON이 아니라 String으로 쓴 이유가, Redis 측 직렬화 유연성이 썩 좋지 않아서 그냥 String으로 넣고 꺼내는 방식으로 사용한 것.
KafkaTemplate은 실제로 주입받아서 메시지 발행에 사용할 예정.
@Configuration
public class KafkaConsumerConfig {
@Value("${spring.kafka.kafka-server}")
private String kafkaServer;
@Value("${spring.kafka.consumer.groud-id}")
private String groupId;
@Bean
public ConsumerFactory<String, Object> consumerFactory() {
Map<String, Object> config = new HashMap<>();
// 카프카 서버 주입
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaServer);
config.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
// AUTO_OFFSET_RESET_CONFIG 등은 생략
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
return new DefaultKafkaConsumerFactory<>(config);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListener() {
ConcurrentKafkaListenerContainerFactory<String, String> listener
= new ConcurrentKafkaListenerContainerFactory<>();
listener.setConsumerFactory(consumerFactory());
return listener;
}
}
카프카에서 프로듀서가 특정 토픽에 메시지 1, 2, 3을 발행하면, 그 토픽을 바라보고 있는 컨슈머(컨슈머 그룹)가 메시지를 소비한다.
Product 서버가 여러 대 있을 때 컨슈머1, 컨슈머2, 컨슈머3를 하나의 컨슈머 그룹으로 묶을 수 있고, 이걸 컨슈머 그룹이라고 부른다.
컨슈머 그룹은 특정 토픽을 바라보는 “그룹” (예: member 토픽용 그룹, product 토픽용 그룹 등)
우리 실습에서는 컨슈머가 한 대지만, group-id는 반드시 설정해주는 것이 기본.
@Component
public class StockKafkaListener {
private final ProductService productService;
private final ObjectMapper objectMapper;
@Autowired
public StockKafkaListener(ProductService productService, ObjectMapper objectMapper) {
this.productService = productService;
this.objectMapper = objectMapper;
}
// 아래 리스너는 토픽을 바라보고 있다가, message 매개변수로 메시지를 받아오는 것
@KafkaListener(topics = "stock-decrease-topic", containerFactory = "kafkaListener")
public void stockConsumer(String message) throws JsonProcessingException {
ProductStockDecreaseReqDto dto =
objectMapper.readValue(message, ProductStockDecreaseReqDto.class);
productService.decreaseStock(dto);
}
}
public Long createFeign(List<OrderItemCreateReqDto> items, String email) {
Ordering order = Ordering.builder()
.memberEmail(email)
.orderStatus(OrderStatus.ORDERED)
.build();
orderingRepository.save(order);
for (OrderItemCreateReqDto itemDto : items) {
// (1) 재고 조회 요청 (동기, Feign 사용)
ProductResDto product = productFeignClient.getProductById(itemDto.getProductId());
if (product.getStockQuantity() < itemDto.getProductCount()) {
throw new IllegalArgumentException("재고가 부족합니다.");
}
// (2) 주문 발생 (DB 저장)
OrderingDetails orderingDetails = OrderingDetails.builder()
.ordering(order)
.productId(itemDto.getProductId())
.productName(product.getName())
.quantity(itemDto.getProductCount())
.build();
orderingDetailRepository.save(orderingDetails);
// (3) 재고 감소 요청 (비동기 Kafka)
kafkaTemplate.send("stock-decrease-topic", itemDto);
}
return order.getId();
}
여기서 (3) 단계는 더 이상 REST 호출이 아니라, Kafka 토픽으로 메시지를 발행하는 구조다.
Product 컨슈머가 메시지를 읽어서 재고 감소 로직을 실행하고 DB에 반영한다.
정리하면:
프로듀서 서비스는 즉시 응답하고, 컨슈머 서비스가 독립적으로 비동기 처리 + DB 저장을 진행하는 구조다.
dependencies {
implementation 'org.springframework.boot:spring-boot-starter-data-jpa'
implementation 'org.springframework.boot:spring-boot-starter-web'
implementation 'org.springframework.boot:spring-boot-starter-validation'
testImplementation 'org.springframework.boot:spring-boot-starter-test'
testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
compileOnly 'org.projectlombok:lombok'
annotationProcessor 'org.projectlombok:lombok'
runtimeOnly 'org.mariadb.jdbc:mariadb-java-client'
implementation 'software.amazon.awssdk:s3:2.29.50'
// redis
implementation 'org.springframework.boot:spring-boot-starter-data-redis'
// 마이크로서비스
implementation 'org.springframework.cloud:spring-cloud-starter-netflix-eureka-client'
}
dependencyManagement {
imports {
mavenBom 'org.springframework.cloud:spring-cloud-dependencies:2024.0.0'
}
}
server:
port: 0
eureka:
client:
serviceUrl:
defaultZone: http://localhost:8761/eureka/
spring:
application:
name: product-service
datasource:
driver-class-name: org.mariadb.jdbc.Driver
url: jdbc:mariadb://localhost:3306/order_system_msa
username: root
password: RDB 비밀번호
jpa:
database: mysql
database-platform: org.hibernate.dialect.MariaDBDialect
generate-ddl: true
hibernate:
ddl-auto: create
show_sql: true
servlet:
multipart:
max-file-size: 1000MB
max-request-size: 1000MB
redis:
host: localhost
port: 6379
aws:
credentials:
access-key: 사용자 엑세스 키
secret-key: 사용자 시크릿 키
region: 위치
s3:
bucket: AWS S3 버킷 이름
컨트롤러 내, Order 서버가 호출하는 엔드포인트:
@PutMapping("/decrease-stock")
public ResponseEntity<?> decreaseStock(@RequestBody ProductStockDecreaseReqDto dto) {
productService.decreaseStock(dto);
return ResponseEntity.status(HttpStatus.OK).body(dto);
}
서비스 – 재고 감소 로직:
public void decreaseStock(ProductStockDecreaseReqDto dto) {
Product product = productRepository.findById(dto.getProductId())
.orElseThrow(() -> new EntityNotFoundException("상품 없음"));
if (product.getStockQuantity() < dto.getProductCount()) {
throw new IllegalArgumentException("상품재고가 부족합니다.");
} else {
product.decreaseStockQuantity(dto.getProductCount());
}
}
비동기 Kafka 구조로 바꾼 후에는, 이 decreaseStock 메서드는 그대로 두고,
컨트롤러 대신 Kafka Listener가 dto를 받아서 이 메서드를 호출하는 형태로 전환한 것.

