JAVA · 심층 가이드
Java 비동기/동시성 완전 정리
Thread 직접 생성부터 Virtual Thread와 구조화된 동시성까지, Java 동시성 코드의 실행 모델과 안전한 종료·취소 방식을 21개 주제로 정리합니다.
Java의 동시성은 오랫동안 '스레드는 비싸다'는 전제 위에 설계됐습니다. OS 스레드와 1:1로 매핑되는 플랫폼 스레드를 아끼려고 풀을 만들었고, 블로킹을 피하려고 콜백과 CompletableFuture 체인으로 코드를 뒤집었습니다. Virtual Thread가 들어오면서 이 전제가 흔들립니다. 블로킹이 싸지면 풀 크기 계산도, 억지 비동기 변환도 상당 부분 불필요해집니다. 초보자가 걸려 넘어지는 지점이 여기입니다. 두 모델이 같은 ExecutorService API 위에 공존하기 때문에, 지금 읽고 있는 코드가 어느 쪽 전제로 쓰였는지 스스로 판단해야 합니다.
읽는 순서는 실행 주체 → 결과 조합 → 수명 관리로 잡는 편이 낫습니다. Thread vs Runnable과 ExecutorService & ThreadPool에서 '무엇을 실행하는가'와 '누가 실행하는가'를 분리해 보고, Future와 Callable로 결과를 되받는 가장 단순한 형태를 확인한 뒤 CompletableFuture — 비동기 프로그래밍으로 넘어가면 조합 연산자가 왜 필요한지가 분명해집니다. Structured Concurrency (Java 21 Preview)는 그 반대 방향입니다. 흩어진 작업을 다시 부모 스코프에 묶어, 하나가 실패하면 나머지를 취소하는 규칙을 문법 차원에서 강제합니다.
가장 자주 밟는 지뢰는 실행자를 지정하지 않은 CompletableFuture입니다. supplyAsync를 executor 인자 없이 호출하면 ForkJoinPool의 공용 풀에서 돌아가는데, 이 풀은 CPU 코어 수를 기준으로 작게 잡히므로 여기에 DB 호출 같은 블로킹 작업을 던지면 무관한 다른 병렬 작업까지 함께 멈춥니다. Virtual Thread 쪽에도 별도의 함정이 있습니다. synchronized 블록 안에서 블로킹하면 캐리어 스레드에 고정(pinning)될 수 있으므로, 그 구간의 락은 ReentrantLock으로 바꿔 두는 편이 안전합니다.
01CompletableFuture — 비동기 프로그래밍
Java 8의 CompletableFuture로 비동기 작업을 조합하고 체이닝합니다.
Java code
import java.util.concurrent.*;
import java.util.List;
public class AsyncDemo {
static CompletableFuture<String> fetchUser(int id) {
return CompletableFuture.supplyAsync(() -> {
sleep(100);
return "User#" + id;
});
}
static CompletableFuture<String> fetchOrders(String userId) {
return CompletableFuture.supplyAsync(() -> {
sleep(100);
return "Orders of " + userId;
});
}
public static void main(String[] args) throws Exception {
// thenApply — 변환 (Function)
CompletableFuture<String> upper = fetchUser(1)
.thenApply(String::toUpperCase);
// thenAccept — 소비 (Consumer)
fetchUser(1).thenAccept(u ->
System.out.println("받음: " + u));
// thenCompose — 플랫맵 (다른 CF 반환)
CompletableFuture<String> orders = fetchUser(1)
.thenCompose(AsyncDemo::fetchOrders);
// thenCombine — 두 CF 결과 합치기
CompletableFuture<String> combined =
fetchUser(1).thenCombine(fetchUser(2),
(u1, u2) -> u1 + " & " + u2);
// allOf — 모두 완료 대기
CompletableFuture<Void> all = CompletableFuture.allOf(
fetchUser(1), fetchUser(2), fetchUser(3));
all.join();
// anyOf — 가장 빠른 것
CompletableFuture<Object> any = CompletableFuture.anyOf(
fetchUser(1), fetchUser(2));
// 예외 처리
CompletableFuture<String> safe = fetchUser(-1)
.exceptionally(ex -> "기본 사용자")
.handle((result, ex) -> {
if (ex != null) return "에러: " + ex.getMessage();
return result;
});
System.out.println(orders.get());
System.out.println(combined.get());
System.out.println(safe.get());
}
static void sleep(long ms) {
try { Thread.sleep(ms); } catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}thenApply()와 thenApplyAsync()의 차이: 전자는 이전 스테이지와 같은 스레드에서, 후자는 ForkJoinPool에서 실행됩니다.
CompletableFuture.get()은 체크드 예외를 던집니다. 더 간단한 .join()을 사용하면 unchecked exception으로 래핑됩니다.
02Virtual Threads (Java 21 — Project Loom)
수백만 개의 경량 스레드를 생성할 수 있는 Virtual Threads로 동시성 코드를 단순화합니다.
Java code
import java.util.concurrent.*;
import java.time.Duration;
public class VirtualThreads {
public static void main(String[] args) throws Exception {
// 기존 Platform Thread
Thread platform = new Thread(() ->
System.out.println("플랫폼 스레드: " + Thread.currentThread()));
platform.start();
platform.join();
// Virtual Thread (Java 21+) — 간단 생성
Thread virtual = Thread.ofVirtual()
.name("my-virtual")
.start(() ->
System.out.println("가상 스레드: " + Thread.currentThread()));
virtual.join();
// 수백만 개 가상 스레드
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
var futures = new java.util.ArrayList<Future<String>>();
for (int i = 0; i < 100_000; i++) {
final int id = i;
futures.add(executor.submit(() -> {
Thread.sleep(Duration.ofMillis(10)); // I/O 시뮬레이션
return "task-" + id;
}));
}
long done = futures.stream()
.filter(f -> { try { f.get(); return true; } catch (Exception e) { return false; } })
.count();
System.out.println("완료: " + done + "개");
}
// StructuredTaskScope (Java 21 Preview)
// try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
// var user = scope.fork(() -> fetchUser(1));
// var orders = scope.fork(() -> fetchOrders("user1"));
// scope.join().throwIfFailed();
// process(user.get(), orders.get());
// }
System.out.println("Virtual Thread는 플랫폼 스레드인가? " +
Thread.currentThread().isVirtual());
}
}Virtual Thread는 I/O 바운드 작업에 최적입니다. CPU 바운드 작업은 여전히 플랫폼 스레드가 적합합니다. synchronized 블록에서 블로킹 I/O를 사용하면 Virtual Thread의 이점이 줄어듭니다.
Virtual Thread에서 ThreadLocal을 남용하면 메모리 문제가 생깁니다. 수백만 스레드 × ThreadLocal 데이터 = 메모리 폭발. ScopedValue(Java 21+)를 고려하세요.
03ExecutorService & ThreadPool
스레드 풀로 자원을 효율적으로 관리하고 병렬 처리를 구현합니다.
Java code
import java.util.concurrent.*;
import java.util.*;
public class ThreadPoolDemo {
public static void main(String[] args) throws Exception {
// FixedThreadPool — 고정 크기 스레드 풀
ExecutorService fixed =
Executors.newFixedThreadPool(4);
List<Future<String>> futures = new ArrayList<>();
for (int i = 0; i < 8; i++) {
final int id = i;
futures.add(fixed.submit(() -> {
Thread.sleep(100);
return "task-" + id + " by " +
Thread.currentThread().getName();
}));
}
for (Future<String> f : futures) {
System.out.println(f.get()); // 블로킹 대기
}
fixed.shutdown();
// ScheduledExecutorService — 예약 실행
ScheduledExecutorService scheduler =
Executors.newScheduledThreadPool(2);
scheduler.schedule(
() -> System.out.println("1초 후 실행"),
1, TimeUnit.SECONDS);
scheduler.scheduleAtFixedRate(
() -> System.out.println("2초마다 실행: " + new Date()),
0, 2, TimeUnit.SECONDS);
Thread.sleep(6000);
scheduler.shutdown();
// ForkJoinPool — 분할정복 병렬 처리
var pool = ForkJoinPool.commonPool();
int[] arr = java.util.stream.IntStream.range(0, 1000).toArray();
int sum = pool.submit(() ->
java.util.Arrays.stream(arr).parallel().sum()
).get();
System.out.println("병렬 합계: " + sum);
}
}Executors.newCachedThreadPool()은 스레드 수 제한이 없습니다. 갑자기 부하가 증가하면 수천 개의 스레드가 생성되어 OOM이 발생할 수 있습니다. 프로덕션에서는 newFixedThreadPool()이나 ThreadPoolExecutor로 직접 설정하세요.
fixed.shutdown() 후 awaitTermination()을 호출하지 않으면 실행 중인 작업이 완료되기 전에 프로그램이 종료될 수 있습니다.
04Scoped Values (Java 21 Preview)
ThreadLocal 대체, 코루틴 친화적 컨텍스트 전달
Java code
<span class="cm">// Scoped Values (Java 21 Preview) 예제
// data/prompts.js의 생성 프롬프트로 상세 코드 생성 가능</span>
fun main() { println("Scoped Values (Java 21 Preview)") }JAVA 공식 문서를 함께 참고하세요.
자주 발생하는 실수에 주의하세요.
05CompletableFuture 고급 패턴
CompletableFuture의 고급 조합 패턴. 여러 비동기 작업을 병렬 실행하고, 타임아웃을 설정하고, 예외를 체계적으로 처리하는 실전 패턴을 익힙니다.
Java code
import java.util.concurrent.*;
import java.util.List;
public class CFAdvanced {
static final ExecutorService pool =
Executors.newVirtualThreadPerTaskExecutor();
// 1. 여러 작업 병렬 실행 후 결과 수집
static CompletableFuture<List<String>> fetchAll(List<String> urls) {
List<CompletableFuture<String>> futures = urls.stream()
.map(url -> CompletableFuture
.supplyAsync(() -> fetch(url), pool)
.orTimeout(3, TimeUnit.SECONDS)
.exceptionally(ex -> "FAILED: " + url))
.toList();
return CompletableFuture.allOf(futures.toArray(new CompletableFuture[0]))
.thenApply(v -> futures.stream()
.map(CompletableFuture::join)
.toList());
}
// 2. 가장 빠른 결과 사용 (anyOf)
static CompletableFuture<String> fetchFastest(String... urls) {
CompletableFuture<String>[] futures = java.util.Arrays.stream(urls)
.map(url -> CompletableFuture.supplyAsync(() -> fetch(url), pool))
.toArray(CompletableFuture[]::new);
return CompletableFuture.anyOf(futures)
.thenApply(String.class::cast);
}
// 3. 파이프라인 패턴
static CompletableFuture<String> pipeline(String input) {
return CompletableFuture.supplyAsync(() -> validate(input), pool)
.thenApplyAsync(CFAdvanced::transform, pool)
.thenApplyAsync(CFAdvanced::enrich, pool)
.whenComplete((result, ex) -> {
if (ex != null) System.err.println("파이프라인 실패: " + ex);
else System.out.println("완료: " + result);
});
}
static String fetch(String url) { return "data:" + url; }
static String validate(String s) { return s.trim(); }
static String transform(String s) { return s.toUpperCase(); }
static String enrich(String s) { return s + "+enriched"; }
}orTimeout()과 completeOnTimeout()은 Java 9+에서 사용 가능. Virtual Thread executor와 함께 사용하면 블로킹 I/O도 효율적으로 처리됩니다.
allOf()는 CompletableFuture<Void>를 반환하므로 개별 결과를 얻으려면 원래 future들에서 join()해야 합니다. 또한 예외 처리 없이 join()하면 CompletionException이 전파됩니다.
06Structured Concurrency (Java 21 Preview)
Java 21의 구조화된 동시성(Structured Concurrency)은 여러 하위 작업의 수명을 부모 작업에 종속시킵니다. 하나라도 실패하면 나머지를 자동 취소하여 리소스 누수를 방지합니다.
Java code
import java.util.concurrent.*;
import jdk.incubator.concurrent.StructuredTaskScope;
// Java 21 Preview — 구조화된 동시성
record UserProfile(String name, int orderCount, double credit) {}
UserProfile fetchProfile(long userId) throws Exception {
// ShutdownOnFailure: 하나라도 실패 시 나머지 자동 취소
try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
// 3개의 하위 작업 동시 실행
Subtask<String> nameTask =
scope.fork(() -> fetchUserName(userId));
Subtask<Integer> orderTask =
scope.fork(() -> fetchOrderCount(userId));
Subtask<Double> creditTask =
scope.fork(() -> fetchCredit(userId));
// 모든 작업 완료 대기 (또는 실패 시 즉시 반환)
scope.join()
.throwIfFailed();
// 모든 결과 수집
return new UserProfile(
nameTask.get(),
orderTask.get(),
creditTask.get()
);
}
// try 블록 종료 시 미완료 작업 자동 취소
}
// ShutdownOnSuccess: 가장 빠른 성공 결과만 사용
String fetchFromMirrors(String key) throws Exception {
try (var scope = new StructuredTaskScope.ShutdownOnSuccess<String>()) {
scope.fork(() -> fetchFromMirror1(key));
scope.fork(() -> fetchFromMirror2(key));
scope.fork(() -> fetchFromMirror3(key));
scope.join();
return scope.result(); // 가장 먼저 성공한 결과
}
}Structured Concurrency는 Virtual Thread와 함께 사용하도록 설계되었습니다. --enable-preview --add-modules jdk.incubator.concurrent 플래그가 필요합니다.
fork()된 작업은 반드시 join() 후에 get()해야 합니다. join 전에 get하면 IllegalStateException이 발생합니다. scope를 try-with-resources 밖에서 사용하지 마세요.
07Thread vs Runnable
스레드 생성 방법과 Thread 상속 vs Runnable 구현의 차이입니다.
Java code
public class ThreadBasics {
// 방법 1: Thread 상속 (권장하지 않음)
static class MyThread extends Thread {
@Override
public void run() {
System.out.printf("[%s] Thread 상속%n",
Thread.currentThread().getName());
}
}
// 방법 2: Runnable 구현 (권장)
static class MyRunnable implements Runnable {
@Override
public void run() {
System.out.printf("[%s] Runnable 구현%n",
Thread.currentThread().getName());
}
}
public static void main(String[] args) throws InterruptedException {
// Thread 상속
new MyThread().start();
// Runnable
new Thread(new MyRunnable()).start();
// 람다 (가장 간결)
Thread t = new Thread(() ->
System.out.println("람다 스레드!"));
t.start();
t.join(); // 스레드 종료 대기
// 스레드 속성
Thread current = Thread.currentThread();
System.out.println("이름: " + current.getName());
System.out.println("우선순위: " + current.getPriority());
System.out.println("데몬: " + current.isDaemon());
}
}Runnable을 사용하면 다른 클래스를 상속할 수 있고, ExecutorService에도 제출할 수 있어 유연합니다.
run()을 직접 호출하면 새 스레드가 아닌 현재 스레드에서 실행됩니다. 반드시 start()를 호출하세요.
08ExecutorService 심화
스레드 풀을 관리하는 ExecutorService의 다양한 구현과 활용법입니다.
Java code
import java.util.concurrent.*;
import java.util.*;
public class ExecutorServiceAdvanced {
public static void main(String[] args) throws Exception {
// 고정 크기 스레드 풀
ExecutorService fixed = Executors.newFixedThreadPool(4);
// 캐시 스레드 풀 (필요 시 생성, 60초 유휴 시 제거)
ExecutorService cached = Executors.newCachedThreadPool();
// 단일 스레드 (순서 보장)
ExecutorService single = Executors.newSingleThreadExecutor();
// 태스크 제출
List<Future<String>> futures = new ArrayList<>();
for (int i = 0; i < 10; i++) {
final int id = i;
futures.add(fixed.submit(() -> {
Thread.sleep(100);
return "결과-" + id;
}));
}
// 결과 수집
for (Future<String> f : futures) {
System.out.println(f.get(5, TimeUnit.SECONDS));
}
// invokeAll — 모두 완료될 때까지 대기
List<Callable<String>> tasks = List.of(
() -> "A", () -> "B", () -> "C");
List<Future<String>> results = fixed.invokeAll(tasks);
// 종료 (필수!)
fixed.shutdown();
if (!fixed.awaitTermination(10, TimeUnit.SECONDS)) {
fixed.shutdownNow();
}
cached.shutdown();
single.shutdown();
}
}try-with-resources(Java 19+ ExecutorService가 AutoCloseable)를 사용하면 자동으로 종료됩니다.
ExecutorService를 shutdown()하지 않으면 JVM이 종료되지 않습니다. finally 블록에서 반드시 종료하세요.
09Future와 Callable
Callable로 결과를 반환하는 비동기 작업과 Future의 한계를 이해합니다.
Java code
import java.util.concurrent.*;
public class FutureCallable {
// Callable — 결과 반환 + 예외 던지기 가능
static class PriceCalculator implements Callable<Double> {
private final String product;
PriceCalculator(String product) { this.product = product; }
@Override
public Double call() throws Exception {
Thread.sleep(1000); // 외부 API 호출 시뮬레이션
return Math.random() * 100;
}
}
public static void main(String[] args) throws Exception {
ExecutorService executor = Executors.newFixedThreadPool(3);
try {
Future<Double> future = executor.submit(
new PriceCalculator("노트북"));
// 결과 대기 전 다른 작업 가능
System.out.println("가격 조회 중...");
System.out.println("취소됨? " + future.isCancelled());
System.out.println("완료됨? " + future.isDone());
// get() — 블로킹 대기
Double price = future.get(5, TimeUnit.SECONDS);
System.out.printf("가격: %.2f%n", price);
// 취소
Future<Double> cancelable = executor.submit(
new PriceCalculator("키보드"));
cancelable.cancel(true); // 인터럽트로 취소
} finally {
executor.shutdown();
}
}
}Future.get()에 타임아웃을 항상 지정하세요. 무한 대기를 방지합니다. 논블로킹이 필요하면 CompletableFuture를 사용하세요.
Future.get()은 블로킹 호출입니다. 여러 Future를 순차적으로 get()하면 병렬 효과가 사라집니다.
10CompletableFuture 심화
CompletableFuture로 비동기 파이프라인을 구성하고 결합합니다.
Java code
import java.util.concurrent.*;
import java.util.*;
public class CompletableFutureAdvanced {
static CompletableFuture<String> fetchUser(String id) {
return CompletableFuture.supplyAsync(() -> {
sleep(100);
return "User-" + id;
});
}
static CompletableFuture<String> fetchOrder(String user) {
return CompletableFuture.supplyAsync(() -> {
sleep(100);
return "Order-" + user;
});
}
static void sleep(long ms) {
try { Thread.sleep(ms); } catch (InterruptedException e) {}
}
public static void main(String[] args) {
// 체이닝: thenApply, thenCompose, thenAccept
CompletableFuture<String> pipeline = fetchUser("123")
.thenCompose(user -> fetchOrder(user)) // flatMap
.thenApply(order -> order.toUpperCase()) // map
.exceptionally(ex -> "에러: " + ex.getMessage());
System.out.println(pipeline.join());
// 여러 작업 병렬 실행 후 결합
CompletableFuture<String> f1 = fetchUser("1");
CompletableFuture<String> f2 = fetchUser("2");
CompletableFuture<String> f3 = fetchUser("3");
// allOf — 모두 완료 대기
CompletableFuture.allOf(f1, f2, f3)
.thenRun(() -> System.out.println("모두 완료!"))
.join();
// anyOf — 하나라도 완료
CompletableFuture<Object> fastest =
CompletableFuture.anyOf(f1, f2, f3);
System.out.println("가장 빠른: " + fastest.join());
// thenCombine — 두 결과 결합
f1.thenCombine(f2, (a, b) -> a + " & " + b)
.thenAccept(System.out::println)
.join();
}
}thenApply은 동기 변환(map), thenCompose는 비동기 변환(flatMap)입니다. Async 접미사 메서드는 다른 스레드에서 실행합니다.
join()은 checked exception을 던지지 않아 편하지만, 예외가 CompletionException으로 래핑됩니다. 원본 예외를 얻으려면 getCause()를 호출하세요.
11ScheduledExecutorService
주기적 작업과 지연 실행을 관리하는 스케줄러입니다.
Java code
import java.util.concurrent.*;
import java.time.LocalTime;
public class SchedulerDemo {
public static void main(String[] args) throws Exception {
ScheduledExecutorService scheduler =
Executors.newScheduledThreadPool(2);
// 1. 지연 실행 (3초 후 1회)
scheduler.schedule(
() -> System.out.println(LocalTime.now() + " 지연 실행"),
3, TimeUnit.SECONDS);
// 2. 고정 주기 (시작 후 1초마다)
// 이전 실행 시작 시간 기준
ScheduledFuture<?> periodic = scheduler.scheduleAtFixedRate(
() -> System.out.println(LocalTime.now() + " 주기 실행"),
0, 1, TimeUnit.SECONDS);
// 3. 고정 지연 (이전 실행 완료 후 1초 대기)
ScheduledFuture<?> delayed = scheduler.scheduleWithFixedDelay(
() -> {
System.out.println(LocalTime.now() + " 지연 주기");
try { Thread.sleep(500); } catch (InterruptedException e) {}
},
0, 1, TimeUnit.SECONDS);
// 5초 후 중지
Thread.sleep(5000);
periodic.cancel(false);
delayed.cancel(false);
scheduler.shutdown();
scheduler.awaitTermination(10, TimeUnit.SECONDS);
}
}scheduleAtFixedRate는 실행 간격 보장, scheduleWithFixedDelay는 완료 후 대기 시간 보장입니다. 작업 시간이 간격보다 길면 차이가 큽니다.
스케줄 작업에서 예외가 발생하면 해당 태스크가 조용히 중단됩니다. 반드시 try-catch로 감싸고 로깅하세요.
12Fork/Join 프레임워크
분할-정복 알고리즘을 병렬로 실행하는 Fork/Join 프레임워크입니다.
Java code
import java.util.concurrent.*;
public class ForkJoinDemo {
// RecursiveTask — 결과 반환
static class SumTask extends RecursiveTask<Long> {
private final long[] array;
private final int start, end;
private static final int THRESHOLD = 1000;
SumTask(long[] array, int start, int end) {
this.array = array;
this.start = start;
this.end = end;
}
@Override
protected Long compute() {
int size = end - start;
if (size <= THRESHOLD) {
// 기본 케이스: 직접 계산
long sum = 0;
for (int i = start; i < end; i++) sum += array[i];
return sum;
}
// 분할
int mid = start + size / 2;
SumTask left = new SumTask(array, start, mid);
SumTask right = new SumTask(array, mid, end);
left.fork(); // 비동기 실행
long rightResult = right.compute(); // 현재 스레드
long leftResult = left.join(); // 결과 대기
return leftResult + rightResult;
}
}
public static void main(String[] args) {
long[] array = new long[10_000_000];
for (int i = 0; i < array.length; i++) array[i] = i + 1;
ForkJoinPool pool = new ForkJoinPool();
long result = pool.invoke(new SumTask(array, 0, array.length));
System.out.println("합계: " + result);
System.out.println("병렬도: " + pool.getParallelism());
pool.shutdown();
}
}오른쪽 태스크를 compute()로 현재 스레드에서 실행하고, 왼쪽만 fork()하면 불필요한 스레드 생성을 줄입니다.
두 태스크 모두 fork()하면 현재 스레드가 놀게 됩니다. 하나는 반드시 compute()로 직접 실행하세요.
13Virtual Thread 심화
Java 21의 가상 스레드로 대규모 동시성을 저렴하게 구현합니다.
Java code
import java.util.concurrent.*;
import java.time.*;
public class VirtualThreadAdvanced {
public static void main(String[] args) throws Exception {
// 가상 스레드 생성
Thread vt = Thread.ofVirtual()
.name("virtual-1")
.start(() -> System.out.println("가상 스레드!"));
vt.join();
// 가상 스레드 ExecutorService
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
// 100만 개 동시 작업도 가능!
int taskCount = 10_000;
var latch = new CountDownLatch(taskCount);
Instant start = Instant.now();
for (int i = 0; i < taskCount; i++) {
executor.submit(() -> {
try {
Thread.sleep(Duration.ofSeconds(1));
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
} finally {
latch.countDown();
}
});
}
latch.await();
Duration elapsed = Duration.between(start, Instant.now());
System.out.printf("%d 태스크 완료: %dms%n",
taskCount, elapsed.toMillis());
// 플랫폼 스레드: 수십 초, 가상 스레드: ~1초
}
// isVirtual() 확인
System.out.println("메인 가상? " +
Thread.currentThread().isVirtual()); // false
}
}가상 스레드는 I/O 바운드 작업에 최적화됩니다. CPU 집약적 작업에는 플랫폼 스레드 풀이 더 적합합니다.
가상 스레드에서 synchronized 블록은 캐리어 스레드를 고정(pinning)합니다. ReentrantLock으로 대체하세요.
14Structured Concurrency 심화
Java 21의 구조적 동시성으로 스레드 수명을 체계적으로 관리합니다.
Java code
// Java 21 Preview: Structured Concurrency
// --enable-preview 필요
import java.util.concurrent.*;
public class StructuredConcurrencyDemo {
record User(String name) {}
record Order(String id) {}
record Response(User user, Order order) {}
static User fetchUser() throws Exception {
Thread.sleep(100);
return new User("홍길동");
}
static Order fetchOrder() throws Exception {
Thread.sleep(150);
return new Order("ORD-001");
}
static Response handle() throws Exception {
// StructuredTaskScope — 하위 작업을 구조적으로 관리
// try (var scope = new StructuredTaskScope.ShutdownOnFailure()) {
// Subtask<User> userTask = scope.fork(() -> fetchUser());
// Subtask<Order> orderTask = scope.fork(() -> fetchOrder());
//
// scope.join(); // 모든 하위 작업 대기
// scope.throwIfFailed(); // 실패 시 예외
//
// return new Response(userTask.get(), orderTask.get());
// }
// 하나 실패 시 나머지 자동 취소!
// 현재 GA에서 사용 가능한 대안
try (var executor = Executors.newVirtualThreadPerTaskExecutor()) {
Future<User> userFuture = executor.submit(() -> fetchUser());
Future<Order> orderFuture = executor.submit(() -> fetchOrder());
return new Response(userFuture.get(), orderFuture.get());
}
}
public static void main(String[] args) throws Exception {
Response resp = handle();
System.out.println(resp);
}
}구조적 동시성은 스레드 누수를 방지합니다. 부모 스코프가 끝나면 모든 자식 스레드가 자동으로 종료됩니다.
기존의 비구조적 동시성(fire-and-forget)은 스레드 누수와 에러 전파 실패의 원인입니다. 구조적 동시성으로 전환하세요.
15ReentrantLock
synchronized보다 유연한 ReentrantLock의 활용법입니다.
Java code
import java.util.concurrent.locks.*;
import java.util.concurrent.*;
public class ReentrantLockDemo {
private final ReentrantLock lock = new ReentrantLock(true); // 공정성
private int balance = 1000;
void deposit(int amount) {
lock.lock();
try {
balance += amount;
System.out.printf("입금 %d -> 잔액 %d%n", amount, balance);
} finally {
lock.unlock(); // 반드시 finally에서!
}
}
boolean tryWithdraw(int amount) {
// tryLock — 논블로킹 시도
if (lock.tryLock()) {
try {
if (balance >= amount) {
balance -= amount;
System.out.printf("출금 %d -> 잔액 %d%n", amount, balance);
return true;
}
} finally {
lock.unlock();
}
}
return false;
}
boolean timedWithdraw(int amount) throws InterruptedException {
// 타임아웃 있는 시도
if (lock.tryLock(1, TimeUnit.SECONDS)) {
try {
balance -= amount;
return true;
} finally {
lock.unlock();
}
}
return false;
}
public static void main(String[] args) {
var account = new ReentrantLockDemo();
account.deposit(500);
account.tryWithdraw(200);
System.out.println("대기 스레드: " + account.lock.getQueueLength());
}
}ReentrantLock(true)로 공정(fair) 모드를 켜면 대기 순서가 보장됩니다. 다만 처리량이 약간 감소합니다.
lock() 후 unlock()을 finally에서 호출하지 않으면 데드락이 발생합니다. 이것은 synchronized 대비 단점입니다.
16ReadWriteLock
읽기-쓰기 잠금으로 읽기 동시성을 높이는 전략입니다.
Java code
import java.util.concurrent.locks.*;
import java.util.*;
public class ReadWriteLockDemo {
private final ReadWriteLock rwLock = new ReentrantReadWriteLock();
private final Lock readLock = rwLock.readLock();
private final Lock writeLock = rwLock.writeLock();
private final Map<String, String> cache = new HashMap<>();
// 읽기 — 동시 접근 허용
String get(String key) {
readLock.lock();
try {
return cache.get(key);
} finally {
readLock.unlock();
}
}
// 쓰기 — 배타적 접근
void put(String key, String value) {
writeLock.lock();
try {
cache.put(key, value);
} finally {
writeLock.unlock();
}
}
// 읽기 -> 쓰기 업그레이드 (직접 불가, 재획득 필요)
String computeIfAbsent(String key, java.util.function.Function<String, String> fn) {
readLock.lock();
try {
String value = cache.get(key);
if (value != null) return value;
} finally {
readLock.unlock();
}
// 읽기 잠금 해제 후 쓰기 잠금 획득
writeLock.lock();
try {
// 다시 확인 (다른 스레드가 먼저 넣었을 수 있음)
return cache.computeIfAbsent(key, fn);
} finally {
writeLock.unlock();
}
}
public static void main(String[] args) {
var demo = new ReadWriteLockDemo();
demo.put("key1", "value1");
System.out.println(demo.get("key1"));
System.out.println(demo.computeIfAbsent("key2", k -> "computed-" + k));
}
}읽기가 쓰기보다 훨씬 빈번한 캐시 시나리오에서 ReadWriteLock이 효과적입니다. Java 8+에서는 StampedLock도 고려하세요.
ReadWriteLock에서 읽기 잠금을 쓰기 잠금으로 직접 업그레이드할 수 없습니다. 읽기를 먼저 해제하고 쓰기를 획득해야 합니다.
17CountDownLatch와 CyclicBarrier
스레드 동기화 도구인 CountDownLatch와 CyclicBarrier를 비교합니다.
Java code
import java.util.concurrent.*;
public class SyncTools {
public static void main(String[] args) throws Exception {
// CountDownLatch — 일회용, N개 이벤트 대기
int workerCount = 3;
CountDownLatch latch = new CountDownLatch(workerCount);
for (int i = 0; i < workerCount; i++) {
final int id = i;
new Thread(() -> {
System.out.println("Worker " + id + " 완료");
latch.countDown(); // 카운트 감소
}).start();
}
latch.await(); // 0이 될 때까지 대기
System.out.println("모든 Worker 완료!");
// CyclicBarrier — 재사용 가능, 모두 도달 시 진행
int parties = 3;
CyclicBarrier barrier = new CyclicBarrier(parties,
() -> System.out.println("== 모두 도착! 다음 단계 =="));
for (int i = 0; i < parties; i++) {
final int id = i;
new Thread(() -> {
try {
System.out.println("Runner " + id + " 1단계");
barrier.await(); // 모두 대기
System.out.println("Runner " + id + " 2단계");
barrier.await(); // 재사용!
} catch (Exception e) { e.printStackTrace(); }
}).start();
}
Thread.sleep(1000); // 결과 출력 대기
}
}CountDownLatch는 일회용(카운트 리셋 불가), CyclicBarrier는 재사용 가능합니다. 단계별 병렬 처리에는 Barrier가 적합합니다.
CyclicBarrier에서 한 스레드가 예외로 중단되면 다른 스레드도 BrokenBarrierException을 받습니다. 예외 처리를 반드시 하세요.
18Phaser
동적으로 참여자를 추가/제거할 수 있는 유연한 Phaser입니다.
Java code
import java.util.concurrent.*;
public class PhaserDemo {
public static void main(String[] args) throws Exception {
// Phaser — 동적 참여자 + 여러 단계
Phaser phaser = new Phaser(1); // 메인 스레드 등록
for (int i = 0; i < 3; i++) {
phaser.register(); // 참여자 등록
final int id = i;
new Thread(() -> {
for (int phase = 0; phase < 3; phase++) {
System.out.printf("Worker %d, 단계 %d 작업 중%n",
id, phase);
try { Thread.sleep(100); } catch (InterruptedException e) {}
phaser.arriveAndAwaitAdvance(); // 단계 동기화
}
phaser.arriveAndDeregister(); // 완료 후 탈퇴
}).start();
}
// 메인 스레드: 각 단계 완료 확인
for (int phase = 0; phase < 3; phase++) {
phaser.arriveAndAwaitAdvance();
System.out.println("== 단계 " + phase + " 완료 ==");
}
phaser.arriveAndDeregister();
System.out.println("최종 단계: " + phaser.getPhase());
System.out.println("등록 수: " + phaser.getRegisteredParties());
}
}Phaser는 CyclicBarrier의 상위 호환입니다. 참여자가 동적으로 변하는 시나리오에서 사용하세요.
arriveAndDeregister()를 호출하지 않으면 Phaser가 영원히 대기합니다. 작업 완료 시 반드시 탈퇴하세요.
19BlockingQueue
생산자-소비자 패턴을 위한 BlockingQueue 구현체들입니다.
Java code
import java.util.concurrent.*;
public class BlockingQueueDemo {
// 생산자-소비자 패턴
static final BlockingQueue<String> queue =
new ArrayBlockingQueue<>(5); // 용량 제한
static class Producer implements Runnable {
@Override
public void run() {
try {
for (int i = 0; i < 10; i++) {
String item = "상품-" + i;
queue.put(item); // 큐 가득 차면 대기
System.out.println("생산: " + item);
}
queue.put("END"); // 종료 신호
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
static class Consumer implements Runnable {
@Override
public void run() {
try {
while (true) {
String item = queue.take(); // 비어있으면 대기
if ("END".equals(item)) break;
System.out.println("소비: " + item);
Thread.sleep(200); // 처리 시간
}
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
}
public static void main(String[] args) throws Exception {
Thread producer = new Thread(new Producer());
Thread consumer = new Thread(new Consumer());
producer.start();
consumer.start();
producer.join();
consumer.join();
System.out.println("완료!");
}
}ArrayBlockingQueue는 용량 고정, LinkedBlockingQueue는 가변(기본 무한)입니다. 메모리 관리를 위해 용량을 지정하세요.
offer()는 큐가 가득 차면 false를 반환하고, put()은 블로킹합니다. 목적에 맞는 메서드를 선택하세요.
20리액티브 스트림과 배압
리액티브 스트림의 배압(backpressure) 메커니즘을 이해합니다.
Java code
import java.util.concurrent.*;
import java.util.concurrent.Flow.*;
public class BackpressureDemo {
// 배압 지원 Subscriber
static class ControlledSubscriber implements Subscriber<Integer> {
private Subscription subscription;
private int received = 0;
private final int batchSize;
ControlledSubscriber(int batchSize) { this.batchSize = batchSize; }
@Override
public void onSubscribe(Subscription s) {
this.subscription = s;
s.request(batchSize); // 초기 요청량
}
@Override
public void onNext(Integer item) {
System.out.println("처리: " + item);
received++;
if (received % batchSize == 0) {
// 배치 처리 완료 후 다음 배치 요청
System.out.println("-- 다음 배치 요청 --");
subscription.request(batchSize);
}
}
@Override
public void onError(Throwable t) {
System.err.println("에러: " + t);
}
@Override
public void onComplete() {
System.out.println("스트림 완료! 총 " + received + "개");
}
}
public static void main(String[] args) throws Exception {
// SubmissionPublisher — JDK 내장 Publisher
try (var publisher = new SubmissionPublisher<Integer>()) {
publisher.subscribe(new ControlledSubscriber(3));
for (int i = 1; i <= 10; i++) {
publisher.submit(i);
}
} // close() -> onComplete() 호출
Thread.sleep(1000); // 비동기 처리 대기
}
}SubmissionPublisher는 JDK에 포함된 유일한 Publisher 구현입니다. 프로덕션에서는 Project Reactor를 사용하세요.
request(Long.MAX_VALUE)는 사실상 배압 없음입니다. 처리 능력에 맞는 요청량을 설정하지 않으면 메모리 폭주가 발생합니다.
21리액티브 스트림 연산자
SubmissionPublisher와 Processor를 사용한 데이터 변환 파이프라인입니다.
Java code
import java.util.concurrent.*;
import java.util.concurrent.Flow.*;
import java.util.function.Function;
public class ReactiveProcessor {
// Processor — Publisher + Subscriber (변환기)
static class TransformProcessor<T, R>
extends SubmissionPublisher<R>
implements Processor<T, R> {
private final Function<T, R> transform;
private Subscription subscription;
TransformProcessor(Function<T, R> transform) {
this.transform = transform;
}
@Override
public void onSubscribe(Subscription s) {
this.subscription = s;
s.request(1);
}
@Override
public void onNext(T item) {
submit(transform.apply(item)); // 변환 후 전달
subscription.request(1);
}
@Override
public void onError(Throwable t) { closeExceptionally(t); }
@Override
public void onComplete() { close(); }
}
public static void main(String[] args) throws Exception {
var publisher = new SubmissionPublisher<String>();
var upper = new TransformProcessor<String, String>(String::toUpperCase);
var lengthProc = new TransformProcessor<String, Integer>(String::length);
// 파이프라인: publisher -> upper -> lengthProc -> subscriber
publisher.subscribe(upper);
upper.subscribe(lengthProc);
lengthProc.subscribe(new Flow.Subscriber<>() {
public void onSubscribe(Subscription s) { s.request(Long.MAX_VALUE); }
public void onNext(Integer item) { System.out.println("길이: " + item); }
public void onError(Throwable t) {}
public void onComplete() { System.out.println("완료!"); }
});
java.util.List.of("hello", "world", "java").forEach(publisher::submit);
publisher.close();
Thread.sleep(500);
}
}Processor는 Publisher와 Subscriber를 모두 구현합니다. 데이터 변환, 필터링, 버퍼링 등의 중간 처리에 사용합니다.
Processor에서 request(1)을 호출하지 않으면 다음 데이터를 받을 수 없습니다. onNext에서 반드시 추가 요청을 하세요.
정리하며
- 블로킹 I/O가 지배적인 서버라면 풀 크기 튜닝보다 Virtual Thread 전환 검토가 먼저입니다
- CompletableFuture는 항상 executor를 명시해 공용 ForkJoinPool 오염을 피합니다
- 취소와 예외 전파는 Structured Concurrency로 부모 스코프에 묶어 한 번에 처리합니다
- Virtual Thread 구간의 상호 배제는 synchronized 대신 ReentrantLock으로 구현합니다
더 깊이 들어가고 싶다면 Java 학습 라이브러리에서 다른 주제 가이드를 이어서 보거나, 언어 비교에서 같은 개념이 다른 언어에서 어떻게 표현되는지 확인해 보세요.