๑ `⌃´ ๑

목표

초기 재고를 100이라 한다. `stock = 100`

100번 차감하면 정상적인 결과는 `100 - 100 = 0`이다.

그런데 여러 스레드가 동시에 `stockService.decrease(productId);`를 호출하게 한 뒤 최종 재고를 확인한다.

기대값 = 0

실제값 = ?

예: 73
    81
    64
    ...

 

일반 테스트로는 순차 실행되기 때문에 동시성 문제를 재현하기 어렵다.

Thread A ──────────┐
Thread B ──────────┤
Thread C ──────────┤ 동시에 같은 재고 접근
Thread D ──────────┤
...                │
Thread N ──────────┘

위와 같은 상황을 만들어야하기 때문에, `ExecutorService`, `CountDownLatch`를 사용한다.

 

 

1. Thread

프로그램에서 실행 흐름 하나를 Thread라고 생각할 수 있는데, 일반 테스트에서는 하나의 흐름이다.

Test Thread
재고 차감

재고 차감

재고 차감

 

우리가 원하는 것은, 여러 실행 흐름을 만드는 것이다.

Test Thread
   │
   ├── Worker Thread 1 → decrease()
   ├── Worker Thread 2 → decrease()
   ├── Worker Thread 3 → decrease()
   └── Worker Thread 4 → decrease()

 

 

2. `ExecutorService`

Java에서 직접 `new Thread(...)`를 100개 만들고 관리할 수도 있다.

하지만 굉장히 번거롭다. `ExecutorService`는 쉽게 말해서, "여러 작업을 여러 Thread에게 나눠 실행하도록 관리해주는 도구"이다.

ExecutorService executorService = Executors.newFixedThreadPool(32);
Thread Pool

Worker 1
Worker 2
Worker 3
...
Worker 32

위 코드 실행 시, 총 32개의 Worker Thread를 준비한다.

 

 

3. Thread Pool

매 요청마다 `Thread 생성 → 사용 → Thread 제거`를 반복하면 비용이 든다.

그래서 미리 일정한 수의 Thread를 만들어 놓고 재사용할 수 있다.

            작업 100개
                │
                ▼
        ┌─────────────────┐
        │   Thread Pool   │
        │                 │
        │ Worker 1        │
        │ Worker 2        │
        │ ...             │
        │ Worker 32       │
        └─────────────────┘
Executors.newFixedThreadPool(32)
// 최대 32개의 Worker Thread를 사용해서 작업을 처리한다.

위 코드의 경우, 100개의 작업을 등록하면 32개까지 동시에 실행되고, 나머지는 자리가 생길 때 차례로 처리된다.

100개가 정확히 같은 나노초에 실행될 필요는 없다.

여러 작업이 충분히 겹쳐 실행되면 동시성 문제를 재현하는 데 충분하다.

 

 

4. `ExecutorService`만 사용했을 때의 문제

executorService.submit(() -> stockService.decrease(productId));
executorService.submit(() -> stockService.decrease(productId));
executorService.submit(() -> stockService.decrease(productId));

작업을 등록했다고 해서 정확히 동시에 출발하는 것은 아니다.

Worker A → 바로 시작
Worker B → 조금 뒤 시작
Worker C → 더 뒤 시작

그래서 출발 시점을 조금 더 겹치게 만들고 싶은데, 이럴 때 사용하는 것이 `CountDownLatch`이다.

 

 

5. `CountDownLatch`

Count Down 숫자를 하나씩 감소
Latch 문을 잠가 놓는 장치

숫자가 0이 될 때까지 기다리게 만드는 문이라고 생각하면 된다.

 

CountDownLatch startLatch = new CountDownLatch(1);
startLatch

count = 1

🚪 문 닫힘

Worker Thread에서 `startLatch.await();`이면 "숫자가 0이 될 때까지 여기서 기다려"라는 뜻이다.

 

 

6. Test Thread가 문을 연다

Worker들이 기다리고 있을 때 Test Thread가 `startLatch.countDown();`를 실행하면 `1 → 0`이 되고 문이 열린다.

               startLatch

                   🔒
                    │
       ┌────────────┼────────────┐
       │            │            │
    Worker A     Worker B     Worker C
     대기         대기         대기


Test Thread

startLatch.countDown()

          ↓

count = 0

          ↓

          🔓

       ┌────────────┼────────────┐
       ↓            ↓            ↓
    Worker A     Worker B     Worker C
      실행          실행          실행

이런 식으로 여러 Worker의 시작 시점을 겹치게 만들 수 있다.

 

 

7. 그런데 테스트는 언제 결과를 확인해야 하는가?

Test Thread가 `startLatch.countDown();`한 직후 바로 `productRepository.findById(...)`를 해버리면 어떻게 될까?

Worker들은 아직 일하고 있을 수 있다.

Test Thread

결과 조회 ← 너무 빠름!

Worker A → 아직 실행 중
Worker B → 아직 실행 중
Worker C → 아직 실행 중

그래서 100개의 작업이 전부 끝날 때까지 기다리는 장치도 필요하다.

 

또 하나의 `CountDownLatch`를 사용한다.

CountDownLatch doneLatch = new CountDownLatch(100);

 

 

8. `doneLatch`는 반대로 사용한다

각 Worker가 일을 끝내면 `doneLatch.countDown();`한다.

초기

doneLatch = 100

작업 하나 완료
→ 99

또 완료
→ 98

...

마지막 작업 완료
→ 0

그리고 Test Thread는 `doneLatch.await();`한다.

"100개의 작업이 모두 완료될 때까지 기다릴게"라는 뜻이다.

 

 

9. 전체 흐름

                        Test Thread
                            │
             작업 100개 ExecutorService에 등록
                            │
                            ▼
                 ┌────────────────────┐
                 │    Thread Pool     │
                 │                    │
                 │ Worker 1           │
                 │ Worker 2           │
                 │ ...                │
                 │ Worker 32          │
                 └────────────────────┘
                            │
                       startLatch
                        await()
                            │
                         🔒 대기
                            │
Test Thread
startLatch.countDown()
                            │
                         🔓 시작
                            │
           ┌────────────────┼────────────────┐
           ↓                ↓                ↓
       decrease()       decrease()       decrease()
           │                │                │
           └────────────────┼────────────────┘
                            │
                 각 작업이 끝날 때마다
                 doneLatch.countDown()
                            │
                         100 → 0
                            │
Test Thread
doneLatch.await() ───────────┘
       │
       ▼
모든 작업 완료
       │
       ▼
최종 stock 조회

 

 

 

10. Spring/JPA와 연결

@Transactional
public void decrease(Long productId) {
    Product product = productRepository.findById(productId).orElseThrow();

    product.decreaseStock();
}

 

Worker 32개가 동시에 실행되면 각각 서로 별개의 트랜잭션이 실행된다.

Transaction A
SELECT
↓
Product
↓
decreaseStock()
↓
UPDATE
↓
COMMIT


Transaction B
SELECT
↓
Product
↓
decreaseStock()
↓
UPDATE
↓
COMMIT


Transaction C
...

 

 

11. `Product.java`

순수 Java 객체로 사용한 `Product`를 수정한다.

이제는 같은 객체를 DB에도 저장해서 여러 트랜잭션이 같은 Product 행을 조회할 수 있게 해야 한다.

따라서 JPA Entity로 변경한다.

 

`Product.java`

package com.bblackbean.jpa.stock;

import jakarta.persistence.Entity;
import jakarta.persistence.GeneratedValue;
import jakarta.persistence.GenerationType;
import jakarta.persistence.Id;
import jakarta.persistence.Table;

@Entity // [추가] JPA가 관리하는 Entity임을 표시
@Table(name = "products") // [추가] 매핑할 테이블 이름
public class Product {

    @Id // [추가] PK
    @GeneratedValue(strategy = GenerationType.IDENTITY) // [추가] PK 자동 생성
    private Long id;

    private int stock;

    // [추가] JPA가 Entity를 생성할 때 필요한 기본 생성자
    protected Product() {
    }

    public Product(int stock) {
        this.stock = stock;
    }

    public void decreaseStock() {
        if (stock <= 0) {
            throw new IllegalStateException("재고가 부족합니다.");
        }

        stock--;
    }

    public Long getId() {
        return id;
    }

    public int getStock() {
        return stock;
    }
}

 

 

12. 기본 생성자는 왜 필요한가?

기존에는 `new Product(10);`만 사용했다.

하지만 JPA는 DB 조회 결과를 이용해서 Entity 객체를 만들어야 한다.

그래서 JPA Entity에는 기본 생성자가 필요하다.

protected Product() {
}

`public`이 아니라 `protected`로 두는 이유는, "애플리케이션 코드가 의미 없이 `new Product()`를 만드는 것은 막으면서 JPA는 사용할 수 있게 하자" 정도의 의미이다.

 

 

13. `ProductRepository.java`

이 파일의 역할

` ProductRepository.java `는 DB에서 Product를 `저장`, `조회`하기 위해 필요하다.

package com.bblackbean.jpa.stock;

import org.springframework.data.jpa.repository.JpaRepository;

public interface ProductRepository extends JpaRepository<Product, Long> {
}

Spirng Data JPA가 실행 시점에 구현체를 만든다.

우리가 직접 `class ProductRepositoryImpl`을 만들 필요는 없다.

 

 

14. `StockService.java`

이 파일의 역할

이 파일은 각 재고 차감 요청마다 별도의 트랜잭션을 생성하는 진입점이다.

Worker Thread가 호출할 대상이 바로 `stockService.decrease(productId);`이다.

Worker Thread
     ↓
StockService.decrease()
     ↓
@Transactional
     ↓
Product 조회
     ↓
decreaseStock()
     ↓
Dirty Checking
     ↓
UPDATE

 

`StockService.java`

package com.bblackbean.jpa.stock;

import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

@Service
public class StockService {

    private final ProductRepository productRepository;

    public StockService(ProductRepository productRepository) {
        this.productRepository = productRepository;
    }

    @Transactional
    public void decrease(Long productId) {
        Product product = productRepository.findById(productId)
                .orElseThrow();

        product.decreaseStock();
    }
}

여기서 중요한 것은 일부러 `productRepository.save(product);`를 추가하지 않았다는 점이다.

findById()
↓
영속 Entity

decreaseStock()
↓
상태 변경

트랜잭션 종료
↓
Dirty Checking
↓
UPDATE

 

 

15. `StockConcurrencyTest.java`

이 파일의 역할

하나의 테스트 스레드에서 100번 순서대로 호출하는 것이 아니라, 순서대로 실행해서 서로 다른 트랜잭션을 겹치게 만든다.

100개의 작업
↓
32개의 Worker Thread
↓
같은 productId
↓
StockService.decrease()

 

`StockConcurrencyTest.java`

package com.bblackbean.jpa.stock;

import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;

import static org.assertj.core.api.Assertions.assertThat;

@SpringBootTest
class StockConcurrencyTest {

    @Autowired
    private ProductRepository productRepository;

    @Autowired
    private StockService stockService;

    @BeforeEach
    void setUp() {
        productRepository.deleteAll();
    }

    @Test
    void 동시에_100번_재고를_차감한다() throws InterruptedException {
        // Given
        Product product = productRepository.saveAndFlush(new Product(100));
        Long productId = product.getId();

        int requestCount = 100;

        // 최대 32개의 작업을 동시에 실행할 Worker Thread Pool
        ExecutorService executorService = Executors.newFixedThreadPool(32);

        // Worker들의 출발 시점을 최대한 맞추기 위한 Latch
        CountDownLatch startLatch = new CountDownLatch(1);

        // 100개의 작업이 모두 끝났는지 확인하기 위한 Latch
        CountDownLatch doneLatch = new CountDownLatch(requestCount);

        // When
        for (int i = 0; i < requestCount; i++) {
            executorService.submit(() -> {
                try {
                    // startLatch가 0이 될 때까지 기다린다.
                    startLatch.await();

                    stockService.decrease(productId);
                } catch (InterruptedException e) {
                    Thread.currentThread().interrupt();
                } finally {
                    // 성공/실패와 관계없이 이 작업이 종료됐음을 알린다.
                    doneLatch.countDown();
                }
            });
        }

        // 기다리고 있던 Worker Thread들을 시작시킨다.
        startLatch.countDown();

        // 100개의 작업이 모두 끝날 때까지 Test Thread가 기다린다.
        doneLatch.await();

        executorService.shutdown();

        // Then
        Product result = productRepository.findById(productId)
                .orElseThrow();

        System.out.println("기대 재고 = 0");
        System.out.println("실제 재고 = " + result.getStock());

        assertThat(result.getStock()).isEqualTo(0);
    }
}

 

 

 

16. `saveAndFlush()`를 사용한 이유

Product product = productRepository.saveAndFlush(new Product(100));

먼저 테스트옹 Product를 DB에 확실하게 준비해야 한다.

Product
stock = 100

↓ saveAndFlush()

DB에 INSERT

그 다음 여러 Worker가 이 `productId`를 조회한다.

 

 

17. `for`문

for (int i = 0; i < requestCount; i++) {
    executorService.submit(() -> {
        ...
    });
}

100번 반복하면서 작업 1, 작업 2, 작업 2, ..., 작업 100을 `ExecutorService`에 등록한다.

처음 최대 32개 → 동시에 실행
나머지 작업 → Queue에서 대기
Worker가 비면 → 다음 작업 실행

 

 

18. `startLatch.await()`는 왜 필요한가?

각 Worker는 처음에 `startLatch.awiat();`에서 기다린다.

초기에는 `startLatch = 1`이므로 문이 닫혀있다.

그러다가 Test Thread가 `startLatch.countDown();`하면, `1 → 0`이 되어 기다리던 Worker들이 실행된다.

A SELECT
B SELECT
C SELECT
D SELECT
...

결과적으로 겹칠 가능성이 높아진다.

 

 

19. `doneLatch.countDown()`을 `finally`에 넣은 이유

finally {
    doneLatch.countDown();
}

만약 `stockService.decrease(productId);` 중 문제가 생겨도 작업 하나는 끝난 것이다.

따라서 `성공 → countDown` `실패 → countDown`이 모두 이루어져야 Test Thread가 영원히 기다리지 않는다.

그래서 `finally`에 넣는다.

 

 

20. 테스트를 실행하면 무엇을 기대해야 하는가?

정상적인 순차 처리라면, `기대 재고 = 0` `실제 재고 = 0`이어야 한다.

하지만 동시성 문제가 발생하면 `기대 재고 = 0` `실제 재고 = 73` 같은 결과를 볼 수도 있다.

그럴 경우 `isEqualTo(0)` 부분에서 테스트가 실패한다.

assertThat(result.getStock()).isEqualTo(0);

 

 

21. 왜 실행할 때마다 결과가 다를 수 있을까?

첫 실행에서는 `actual`이 74일 수도 있고, 두 번째 실행에서는 68, 세 번째 실행에서는 81 처럼 달라질 수도 있다.

심지어 타이밍에 따라 테스트가 우연히 통과할 수도 있다.

왜냐하면 Race Condition은 말 그대로 Thread의 실행 순서와 타이밍에 결과가 영향을 받는 문제이기 때문이다.

따라서 동시성 테스트에서는 한 번 실행했는데 성공했다고 동시성 문제가 없다고 판단하면 안 된다.

 

 

22. 프레임워크가 해주는 것과 우리가 하는 것

우리가 하는 것

100개의 작업 생성
Thread Pool 크기 결정
동시에 출발시키기
모든 작업 완료까지 기다리기
최종 재고 확인

 

ExecutorService가 하는 것

Worker Thread 관리
등록된 작업을 Thread에게 배분
Thread 재사용

 

CountDownLatch가 하는 것

특정 조건이 될 때까지 Thread를 기다리게 함

 

Spring이 하는 것

각 worker가 `stockService.decrease(productId);`를 호출할 때 `@Transactional`을 통해 트랜잭션 경계 관리

 

JPA가 하는 것

Product 조회
영속 Entity 관리
Dirty Checking
UPDATE

 

 

22. 테스트 실행 결과

Test Thread
↓
100개 작업 등록

↓

ExecutorService
32개 Worker Thread 사용

↓

startLatch
동시에 출발

↓

각 Worker

StockService.decrease()

↓

각자 별도 Transaction

↓

같은 Product 조회 및 수정

↓

Lost Update 발생

↓

doneLatch
100개 작업 완료 대기

↓

최종 조회

Expected = 0
Actual = 88

`@Transactional`이 각 요청의 트랜잭션은 정상적으로 관리했지만, 여러 트랜잭션이 같은 재고를 동시에 수정하면서 Lost Update가 발생했다.

 

 

정리

  • 동시성 테스트는 여러 Worker Thread가 같은 공유 데이터에 접근하도록 만들어 순차 테스트에서 드러나지 않는 문제를 재현한다.
  • `ExecutorService`는 작업과 Worker Thread를 관리하고, `CountDownLatch`는 여러 Thread의 시작 또는 완료 시점을 맞추는 데 사용할 수 있다.
  • 동시성 테스트가 실패했다면 바로 해결책을 넣기보다 먼저 어떤 실행 순서 때문에 정합성이 깨졌는지 확인해야 한다.

 

🎵 Playlist
loading...
00:00 / 00:00