Java

[Project Reactor] 1. Introduction to Reactive Programming

noahkim_ 2026. 8. 17. 21:16

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

 

 

출처