[Spring] order-service 실습 : MSA 아키텍처 변경

이지연·2026년 2월 20일

실습 개요

  1. order-system 모놀리식 아키텍처 구조 파일들을 서비스별로 분리해서 생성

    • api-gateway, eureka, member, product, ordering
  2. 불필요한 의존성과 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 반환
        → 이렇게 하면 토큰 없이 서버 직격 불가능

api-gateway

주요 역할: 라우팅, CORS, 토큰 검증

1. build.gradle 설정

  • spring-cloud-starter-gateway 의존성 추가
  • spring-cloud-starter-netflix-eureka-client 의존성 추가
  • 인증 공통 처리 시 jwt 토큰 관련 의존성 추가
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'
    }
}

2. application.yml 설정

  • url 패턴별 라우팅 대상과 경로 지정
  • 라우팅 시 참고할 eureka 서버 엔드포인트 등록
  • CORS 설정 (각 서비스 모듈은 브라우저와 직접 통신 X → CORS 필요 없음, gateway에서만 처리)
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 시크릿 키
  • port는 8080
  • eureka 위치 정보 필요
  • name은 크게 중요 X (eureka 등록용 이름 정도)

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에서 처리.

3. JwtAuthFilter 클래스

@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 토큰 인증 처리

    • 특징
      • Spring Security 대신 GlobalFilter 사용
      • GlobalFilter는 Netty 기반 비동기 구조에 최적화 → 성능 향상
    • 장점
      • 한 군데에서 인증 처리 → 코드 중복 제거, 유지보수성 향상
      • 인증 실패 시 gateway에서 바로 차단 → 불필요한 내부 호출 줄어들어 성능 향상
    • 단점
      • id, role 등 필요한 header 값을 gateway에서 직접 세팅해줘야 함
      • 내부 네트워크(Private Network) 아키텍처 전제 필요

eureka

1. build.gradle 설정

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'
    }
}

2. 진입점 설정

  • 애플리케이션 진입점에 @EnableEurekaServer 어노테이션 추가

3. application.yml 설정

server:
  port: 8761

spring:
  application:
    name: eureka

eureka:
  client:
    # eureka 서버가 자기 자신은 eureka에 등록하지 않겠다는 설정
    register-with-eureka: false
    fetch-registry: false
  • MEMBER-SERVICE(마이크로서비스)에서 Eureka 등록 성공 로그 예시
    • 204 = 등록 완료 (UP 상태)
    • localhost:member-service:0 이 Eureka(http://localhost:8761)에 살아있음
    • 서비스 디스커버리 정상 동작
    • 문제 발생 시 status 404/500 나오면 네트워크/포트/Eureka 설정 확인

각 서비스 모듈 공통 개념

  • 서비스 간 상호 의존성 제거 후 각 서비스 독립 실행
  • DB 설계
    • entity 변수로 타 모듈 객체 직접 참조 제거
    • 필요 시 반정규화로 모듈 간 통신 최소화
  • 서비스 간 통신
    • 기존 모놀리식에서 엔티티 참조로 처리하던 부분을 모듈 간 통신 방식으로 변경
    • 예: 주문 상황에서 order 모듈 ↔ product 모듈
      • 동기 요청: RestTemplate(RestClient) 또는 FeignClient
      • 비동기 이벤트 기반 요청: RabbitMQ 또는 Kafka

1. member

build.gradle

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'
    }
}

application.yml

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 서버 실행 시 유레카에 포트번호가 등록되기 때문에, 게이트웨이가 알아서 그 포트로 라우팅한다.

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인 경우 관례적으로 사용.


2. ordering

build.gradle

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'
    }
}

application.yml

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

RestTemplateConfig

@Configuration
public class RestTemplateConfig {

    @Bean
    @LoadBalanced
    // eureka에 등록된 서비스명을 사용하여 내부 서비스 호출(내부통신)하는 어노테이션
    public RestTemplate makeRestTemplate() {
        return new RestTemplate();
    }
}

이 클래스를 정의한 이유는 단순히 빈 생성을 위한 목적이 아니라,
@LoadBalanced를 통해 eureka에 질의하여 서비스 모듈 간 통신을 하기 위해서다.
결국 uri: lb://product-service와 동일한 로직.

기존 vs 변경 (동기 HTTP 예시)

기존: 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("재고가 부족합니다.");
}

이 때 서비스 모듈들은 같은 내부 네트워크(사설 네트워크)에 존재하고, 외부에서는 접근 불가능하지만 서비스끼리는 통신 가능한 구조.

OrderService 주문 생성 로직 (RestTemplate, 동기)

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단계 핵심:

  1. Ordering 생성 & 저장 (ID 생성)
  2. 각 아이템별로 재고 조회 → 주문 상세 저장 → 재고 감소
  3. @Transactional로 전체 트랜잭션을 묶어 재고 감소 실패 시 전체 롤백

Eureka + @LoadBalanced RestTemplate로
"product-service"라는 서비스명만 알고 통신하고 실제 포트/인스턴스는 eureka가 반환해준다.
서버 확장/축소 시 코드 수정 없이 자동 대응.

이 로직의 (3) 재고 감소 부분을 Kafka 비동기 처리로 개선해볼 것.

ProductFeignClient 정의

@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 메서드가 한 눈에 보임

  • product-service: eureka에 등록된 서비스명
  • /product/detail/{id}: 엔드포인트
  • @GetMapping / @PutMapping: HTTP 메서드

3) 편의성 ↑: 자동 응답 형변환 + 예외 처리 간편화

OrderService – Feign 사용 (동기)

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();
}

Kafka 도입 – 비동기 재고 감소

개념 변화 (동기 → 비동기)

기존 (동기 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.

카프카에는 두 가지 요소:

  • producer(생산자): 메시지 발행 → “재고 감소” 메시지 발행
  • consumer(소비자): 메시지 수신 → “재고 감소” 메시지 받아서 재고 감소

생산자, 소비자 쪽 각각 설정이 필요하고, 보통 consumer 쪽 설정이 (Redis 리스너와 비슷하게) 조금 더 복잡하다.


Kafka 설정

1. Kafka 도커 서버 실행

루트 디렉토리 구조 예시:

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로 설정.

2. Kafka 관련 의존성

(ordering, product 쪽 build.gradle에 spring-kafka 추가했다고 가정)

3. application.yml – 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로 수정)

4. Kafka 연결 빈 – Ordering (Producer)

@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: 메시지 분배 기준(파티션 결정)
  • value: 실제 데이터(JSON 객체 등)

key는 써도 되고 안 써도 된다.
우리 실습에선 key를 null로 두고 value만 보낸다.
실전에서 채팅 서비스를 카프카로 붙일 경우, roomId 같은 값을 key로 사용해서 동일 채팅방 메시지를 같은 파티션에 보내야 하므로 key를 반드시 쓴다.

우리는 value만 보고, 여기에 JSON으로 넣기 때문에 JsonSerializer로 설정.
StringSerializer를 쓰면 JSON을 직접 ObjectMapper로 직렬화해서 String으로 만들어야 하는 차이가 있다.

Redis는 JSON이 아니라 String으로 쓴 이유가, Redis 측 직렬화 유연성이 썩 좋지 않아서 그냥 String으로 넣고 꺼내는 방식으로 사용한 것.

KafkaTemplate은 실제로 주입받아서 메시지 발행에 사용할 예정.

5. Kafka 연결 빈 – Product (Consumer)

@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: 컨슈머 그룹을 식별하는 ID
  • 같은 group-id를 가진 컨슈머들은 메시지를 분담해서 consume

우리 실습에서는 컨슈머가 한 대지만, group-id는 반드시 설정해주는 것이 기본.

6. Product 쪽 Kafka Listener

@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);
    }
}

OrderService – Feign + Kafka 사용 (비동기 재고 감소)

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에는 이미 주문/주문 상세 저장
  • 컨슈머(상품 서비스): 메시지 수신 → 재고 감소 로직 수행 → DB 반영

프로듀서 서비스는 즉시 응답하고, 컨슈머 서비스가 독립적으로 비동기 처리 + DB 저장을 진행하는 구조다.


3. product

build.gradle

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'
    }
}

application.yml

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 버킷 이름

Product 서버 – Order와의 통신 (동기 버전에서 사용했던 부분)

컨트롤러 내, 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를 받아서 이 메서드를 호출하는 형태로 전환한 것.


+) Kafka


profile
Eazy하게

0개의 댓글