0. Reactive Programming이란?
- Reactive Programming은 비동기 데이터 Stream과 변화의 전파를 다루는 Programming Paradigm이다.
- ✅ 데이터가 도착하면 Subscriber에게 전달
- ✅ 값뿐 아니라 Error / Complete도 Signal로 전달
- ✅ 처리 과정을 명령형이 아니라 선언형 Chain으로 표현
Iterator vs Reactive Streams
| 구분 | Iterator | Reactive Streams |
| 데이터 요청 방식 | Consumer가 직접 가져옴 | Publisher가 Subscriber에게 전달 |
| 대표 구조 | Iterable → Iterator | Publisher → Subscriber |
| 처리 방식 | Imperative | Declarative |
| 다음 데이터 | next() 호출 | onNext() Signal |
1. Blocking Can Be Wasteful
- Java Application은 전통적으로 Blocking 방식으로 많이 작성된다.
- ⚠️ I/O 대기 중에도 Thread라는 Resource가 계속 묶여 있다는 것이다.
- ➡️ Thread 증가 → Memory 사용 증가 → Context Switching 증가 → Contention 증가
예시) 전통적인 Blocking 방식
더보기
- Thread → DB 요청 → 응답 대기 → Thread Idle → 응답 도착 → 다시 실행
2. Asynchronicity to the Rescue?
- I/O 완료를 기다리는 동안 Thread를 점유하지 않고, 해당 Thread가 다른 작업을 처리할 수 있다.
- ✅ 요청 → 즉시 반환 → 다른 작업 처리 → 완료 Event → 기존 작업 이어서 처리
- ⚠️ Java에는 비동기 Programming 방식인 Callback, Future를 제공하지만 복잡한 비동기 흐름을 조합할 때 한계가 있다.
Callback / Future와 Reactor
| 방식 | 문제점 | Reactor 해결방식 |
| Callback | - 중첩된 Callback으로 가독성 저하 - Error 처리 반복 - Callback Hell |
여러 비동기 작업을 Operator Chain으로 연결 |
| Future | - 여러 Future 조합이 복잡 - get() 시 Blocking 가능 - 여러 값 처리와 Error 조합이 불편 |
Mono / Flux와 Operator로 비동기 작업 간 관계를 Pipeline으로 표현 |
예시) Reactor - chain 형식으로 표기
더보기
Favorite ID 조회 → 상세 정보 비동기 조회 → 비어 있으면 추천 조회 → 5개 제한 → UI Thread 이동 → 출력
userService.getFavorites(userId)
.flatMap(favoriteService::getDetails)
.switchIfEmpty(suggestionService.getSuggestions())
.take(5)
.publishOn(UiUtils.uiThreadScheduler())
.subscribe(uiList::show, UiUtils::errorPopup);
| Operator | 의미 |
| flatMap() | 각 값을 이용해 다른 비동기 Publisher를 실행하고 결과를 연결 |
| switchIfEmpty() | 결과가 비어 있으면 다른 Publisher로 대체 |
| take(5) | 최대 5개의 데이터만 사용 |
| publishOn() | 이후 처리를 지정한 Scheduler / Thread에서 수행 |
| subscribe() | Stream 실행 및 최종 데이터 / Error 처리 |
- getFavorites() 호출 → Flux/Mono(Publisher) 반환 → Operator로 Pipeline 정의 → subscribe()가 트리거 → 실제 데이터 흐름 시작
예제) Reactor - timeout / fallback도 chain 처리
더보기
Favorite 조회 → 800ms 넘으면 Error → Cache 조회로 대체 → 각 ID 상세조회 → 결과 없으면 추천조회 → 5개만 사용 → UI Thread로 전환 → 화면 표시
userService.getFavorites(userId)
.timeout(Duration.ofMillis(800))
.onErrorResume(cacheService.cachedFavoritesFor(userId))
.flatMap(favoriteService::getDetails)
.switchIfEmpty(suggestionService.getSuggestions())
.take(5)
.publishOn(UiUtils.uiThreadScheduler())
.subscribe(uiList::show, UiUtils::errorPopup);
| 코드 | 의미 |
| getFavorites(userId) | Favorite ID들을 발행할 Publisher 반환 |
| timeout(800ms) | 800ms 동안 데이터가 안 오면 Timeout Error 발생 |
| onErrorResume(...) | 앞 단계에서 Error가 나면 Cache Publisher로 대체 |
| flatMap(...) | 각 Favorite ID를 상세 Favorite 조회 Publisher로 변환/연결 |
| switchIfEmpty(...) | Favorite 결과가 0개면 추천 Publisher로 대체 |
| take(5) | 최대 5개까지만 받음 |
| publishOn(...) | 이후 Signal 처리를 UI Scheduler/Thread에서 수행 |
| subscribe(...) | Pipeline 구독 시작, 값은 show, Error는 errorPopup으로 처리 |
예제) CompletableFuture 조합
더보기
ID 목록 조회 → 각 ID별 Name 조회 + Stat 조회 → 둘을 합침 → 전체 결과 List 생성
CompletableFuture<List<String>> ids = ifhIds();
CompletableFuture<List<String>> result = ids.thenComposeAsync(idList -> {
List<CompletableFuture<String>> tasks = idList.stream()
.map(id -> {
CompletableFuture<String> nameTask = ifhName(id);
CompletableFuture<Integer> statTask = ifhStat(id);
return nameTask.thenCombineAsync(
statTask,
(name, stat) -> "Name " + name + " has stats " + stat
);
})
.toList();
CompletableFuture<Void> allDone =
CompletableFuture.allOf(tasks.toArray(new CompletableFuture[0]));
return allDone.thenApply(v ->
tasks.stream().map(CompletableFuture::join).toList()
);
});
List<String> results = result.join();
- Future<List<ID>> → 각 ID마다 Future<Name> + Future<Stat> → thenCombine → List<Future> → allOf → join 필요
예제) Reactor - Future 조합 단순화
더보기
Flux → flatMap → Name Mono + Stat Mono → zipWith → 결과 Flux → collectList → Mono
Flux<String> ids = ifhrIds();
Flux<String> combinations =
ids.flatMap(id -> {
Mono<String> nameTask = ifhrName(id);
Mono<Integer> statTask = ifhrStat(id);
return nameTask.zipWith(
statTask,
(name, stat) -> "Name " + name + " has stats " + stat
);
});
Mono<List<String>> result = combinations.collectList();
| Operator | 의미 |
| flatMap() | 각 ID에 대해 비동기 작업 실행 |
| zipWith() | 두 비동기 결과를 조합 |
| collectList() | 여러 값을 하나의 List로 모음 |
Imperative vs Reactive
| 구분 | Imperative | Reactive |
| 실행 방식 | 어떻게 실행할지 직접 제어 | 데이터 흐름을 선언 |
| 핵심 구성 | Loop / if / callback 중심 | Operator Chain 중심 |
| 관점 | 작업 순서 중심 | Stream 변환 관계 중심 |
| 동작 특성 | Blocking 코드가 섞이기 쉬움 | Async / Non-Blocking 조합에 적합 |
Reactor가 해결하려는 핵심
- Callback과 Future가 비동기 작업 자체를 가능하게 해주긴 하지만 복잡한 비동기 작업을 조합하는 데 한계가 있음.
- ✅ Reactor는 Reactive Streams의 Publisher → Subscriber 모델을 기반으로 전체 과정을 Operator Chain으로 표현한다.
- ✅ 데이터 생성 → 비동기 변환 → 조합 → Error 처리 → Thread 전환 → 최종 소비
- ex) Publisher → flatMap → timeout → onErrorResume → take → publishOn → subscribe
3. From Imperative to Reactive Programming
- Reactive Library는 기존 Callback / Future 방식의 단점을 줄이면서 많은 기능을 제공함
- ➡️ 비동기 작업을 명령 순서가 아니라 Data가 흘러가는 Pipeline으로 선언함
Composability and Readability
- Composability는 여러 비동기 작업을 연결하거나 조합하여 하나의 처리 흐름으로 구성하는 능력
- ⚠️ 기존 Callback 방식은 작업이 복잡해질수록 Callback이 중첩되어 Callback Hell이 발생하기 쉬움.
- ✅ Reactor에서는 대부분 같은 Depth에서 Chain으로 표현 가능함
- ex) Publisher → flatMap → zipWith → filter → map → subscribe
예시) 여러 비동기 작업을 연결하거나 조합하는 상황
더보기
- 사용자 조회 → 사용자 ID로 주문 조회 → 주문 결과로 결제 조회
- Name 조회 + Stat 조회 → 두 결과 결합
The Assembly Line Analogy
- Reactor는 데이터 처리를 공장 조립 라인에 비유할 수 있음
- ✅ Publisher → 원재료 공급 Source
- ✅ Operator → 데이터를 가공하는 각 작업 공정
- ✅ Subscriber → 최종 결과를 받는 소비자
- 어떤 특정 단계가 너무 느리면 Backpressure를 통해 upstream에 데이터 공급량을 줄이라는 신호를 보낼 수 있음
예시) backpressure
더보기
Flux.range(1, 100)
.map(i -> i * 10)
.limitRate(5)
.subscribe(System.out::println);
- limitRate(): downstream이 데이터를 처리할 때 upstream에 한 번에 너무 많은 데이터를 요청하지 않도록 Demand를 조절
Operators
- 기존 Publisher를 감싸 새로운 Publisher를 만들고 Pipeline을 연결하는 처리 단계
- ✅ 기존 Publisher 자체를 수정하는 게 아니라 새로운 Publisher를 반환함
예시) operators
더보기
Flux.just(1, 2, 3)
.map(i -> i * 10)
.filter(i -> i >= 20);
- Flux(1,2,3) → map → 10,20,30 → filter → 20,30
Nothing Happens Until You Subscribe()
- pipeline은 구독자가 subscribe() 할 때 실행됨
- ✅ pipeline 정의 시, 바로 실행되지 않음
- ✅ subscribe() → 구독/요청 Signal이 upstream으로 전달 → Source Publisher 동작 → Data가 downstream으로 전달
Backpressure
- Subscriber가 자신이 처리할 수 있는 데이터의 양을 upstream에 전달하여 데이터 공급량을 조절하는 방식이다.
- ✅ Subscriber는 request(n)을 통해 Demand를 전달함.
- ✅ Data: Publisher → Subscriber
- ✅ Demand: Subscriber → request(n) → Publisher
- ➡️ Push + Pull이 결합된 구조로 동작함.
| 구분 | 핵심 동작 | 예시 |
| Unbounded | 제한 없이 데이터를 요청하여 Publisher가 가능한 만큼 계속 전달 | request(Long.MAX_VALUE) |
| Bounded | 처리 가능한 개수만 요청하고, 모두 전달되면 추가 request() 대기 | request(10) → 최대 10개 전달 → 추가 요청 대기 |
| Operator Demand 조절 | 중간 Operator가 자신의 데이터 단위에 맞게 upstream 요청량을 변환 | request(1 buffer) → buffer(10) → upstream request(10 items) |
| Prefetch | 필요한 데이터를 일정량 미리 요청하여 반복적인 요청 비용을 줄임 | 미리 request(n) → 내부 확보 → 순차 처리 |
예제) Bounded
더보기
Flux.range(1, 10)
.subscribe(new BaseSubscriber<Integer>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
request(10);
}
@Override
protected void hookOnNext(Integer value) {
System.out.println("처리 시작: " + value);
try {
Thread.sleep(1000);
} catch (InterruptedException e) {
}
System.out.println("처리 끝: " + value);
}
});
- 앞으로 onNext를 최대 10번 호출해도 됨
- 더 받고 싶으면 다시 request(10) 이 호출됨
예제) Operator Demand 조절
더보기
Flux.range(1, 100)
.buffer(10)
.subscribe(new BaseSubscriber<List<Integer>>() {
@Override
protected void hookOnSubscribe(Subscription subscription) {
request(1);
}
@Override
protected void hookOnNext(List<Integer> value) {
System.out.println(value);
}
});
- buffer(10): 한개의 버퍼당 10개 데이터
- request(1): 버퍼 1개 요청
예제) Prefetch
더보기
Flux.range(1, 100)
.flatMap(
i -> Mono.fromCallable(() -> {
Thread.sleep(1000); // 느린 외부 API / DB 작업이라고 가정
return "result-" + i;
}).subscribeOn(Schedulers.boundedElastic()),
4, // 동시에 최대 4개 작업
8 // 각 내부 Publisher에서 미리 확보할 양
)
.subscribe(System.out::println);
- 필요할 때마다 하나씩 요청하지 말고, 미리 일정량을 확보해두기
Hot vs Cold
- Subscriber가 구독할 때 Source가 어떻게 동작하느냐에 따라 Cold와 Hot으로 구분할 수 있다.
| 구분 | Cold | Hot |
| 실행 방식 | Subscriber마다 새로 시작 | 하나의 Stream을 공유 |
| Source 실행 | 구독할 때마다 다시 실행 | Subscriber마다 다시 실행하지 않음 |
| 늦게 구독 | 처음부터 데이터 수신 | 보통 구독 이후 데이터부터 수신 |
| 대표 예시 | HTTP 요청, DB 조회 | 실시간 Event, 주식 가격, Sensor |
예시 코드) Cold - HTTP 요청
더보기
Mono<String> httpCall = Mono.defer(() -> {
System.out.println("HTTP 요청 실행");
return Mono.just("response");
});
httpCall.subscribe(System.out::println);
httpCall.subscribe(System.out::println);
예시 코드) Hot - 실시간 Event
더보기
Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer();
Flux<String> events = sink.asFlux();
events.subscribe(data -> System.out.println("A: " + data));
sink.tryEmitNext("event-1");
sink.tryEmitNext("event-2");
events.subscribe(data -> System.out.println("B: " + data));
sink.tryEmitNext("event-3");
| 코드 | 의미 | 현재 구독자 | 결과 |
| Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer(); | 여러 Subscriber에게 데이터를 전달할 수 있는 Hot Sink 생성 | 없음 | 아직 아무 데이터도 전달되지 않음 |
| Flux<String> events = sink.asFlux(); | Sink를 Subscriber가 구독할 수 있는 Flux 형태로 노출 | 없음 | events가 Sink의 데이터 Stream 역할 |
| events.subscribe(data -> System.out.println("A: " + data)); | Subscriber A가 Stream에 참여 | A | 이후 발행되는 데이터부터 A가 수신 |
| sink.tryEmitNext("event-1"); | Sink에 "event-1" 발행 | A | A: event-1 |
| sink.tryEmitNext("event-2"); | Sink에 "event-2" 발행 | A | A: event-2 |
| events.subscribe(data -> System.out.println("B: " + data)); | Subscriber B가 뒤늦게 Stream에 참여 | A, B | B는 이전 event-1, event-2를 받지 못함 |
| sink.tryEmitNext("event-3"); | Sink에 "event-3" 발행 | A, B | A: event-3, B: event-3 |
출처
'Java' 카테고리의 다른 글
| [Project Reactor] 2. Reactor Core Features (0) | 2026.08.18 |
|---|---|
| [Reactive Streams] 1. 비동기 스트림과 Backpressure (0) | 2026.08.17 |
| [Java][Tutorial] 1-3. Learning the Java Language: Classes and Objects (0) | 2024.01.14 |
| [Java][Tutorial] 1-2. Learning the Java Language: Language Basics (0) | 2024.01.14 |
| [Java] Stream (3) | 2023.10.22 |