Java

[Project Reactor] 2. Reactor Core Features

noahkim_ 2026. 8. 18. 16:18

0. 'reactor-core' 모듈

  • Reactor의 핵심 모듈.
  • Publisher를 구현하면서 다양한 Operator를 제공하는 타입 제공
  • ✅ 대략 몇개의 데이터를 발행하는지 타입 자체로 표현함

 

1. Flux, an Asynchronous Sequence of 0-N Items 

  • 0~N개의 데이터를 비동기적으로 발행할 수 있는 Publisher<T>
  • ✅ Infinite Stream: Flux는 반드시 종료되는 Stream일 필요는 없음.

 

Signal

Signal 의미
onNext(T) 데이터 1개 전달
onComplete() 정상적으로 Stream 종료
onError(Throwable) Error로 Stream 종료

 

 

예시) 여러 데이터 발행

더보기
Flux.just("A", "B", "C")
    .subscribe(System.out::println);
  • onNext("A") → onNext("B") → onNext("C") → onComplete()

 

예시) Empty Flux

더보기
Flux<String> flux = Flux.empty();
  • 데이터가 없으므로 onNext 호출 안됨
  • 바로 onComplete()가 호출되므로 종료됨

 

예시) Error로 종료

더보기
Flux.just(1, 2, 3)
    .map(i -> {
        if (i == 2) {
            throw new RuntimeException("error");
        }
        return i;
    })
    .subscribe(
        System.out::println,
        System.out::println
    );
  • 1 → Error → 종료

 

예시) Infinite Stream

더보기
Flux.interval(Duration.ofSeconds(1))
    .subscribe(System.out::println);
  • onComplete()가 자동으로 발생하지 않는 Infinite Stream

 

2. Mono, an Asynchronous 01 Result

  • 0~1개의 데이터를 비동기적으로 발행하는 Publisher<T>

예시) 아무것도 발행 안됨

더보기
Mono<String> mono = Mono.never();

 

예시) 값 0개 발행

더보기
Mono<String> mono = Mono.empty();

 

예시) 값 1개 발행

더보기
Mono<String> mono = Mono.just("hello");

mono.subscribe(System.out::println);

 

예시) Error Mono

더보기
Mono<String> mono = Mono.error(new RuntimeException("조회 실패"));
  • onError가 발생하며 종료됨

 

예시) Operator에 따라 Mono → Flux로 변환

더보기
Mono<String> first = Mono.just("A");

Flux<String> result = first.concatWith(Mono.just("B"));

 

 

3. Simple Ways to Create a Flux or Mono and Subscribe to It

  • Flux/Mono는 다양한 Factory Method를 통해 쉽게 생성할 수 있음
Factory Method 의미 예시
Flux.just(...) 지정한 여러 값을 Flux로 생성 Flux.just("A", "B", "C")
Flux.fromIterable(...) Collection / IterableFlux로 변환 Flux.fromIterable(list)
Flux.range(start, count) 연속된 정수 데이터를 생성 Flux.range(5, 3)5, 6, 7
Flux.empty() 데이터 없이 정상 종료하는 Flux 생성 Flux.empty()
Flux.error(...) 즉시 Error를 발생시키는 Flux 생성 Flux.error(new RuntimeException())
Flux.never() 데이터·완료·Error Signal이 없는 Flux 생성 Flux.never()
Mono.just(value) 값 1개를 가진 Mono 생성 Mono.just("A")
Mono.justOrEmpty(value) 값이 있으면 발행, null이면 Empty Mono.justOrEmpty(user)
Mono.empty() 값 없이 정상 종료하는 Mono 생성 Mono.empty()
Mono.error(...) 즉시 Error를 발생시키는 Mono 생성 Mono.error(new RuntimeException())
Mono.never() 아무 Signal도 발생하지 않는 Mono 생성 Mono.never()

 

subscribe() Examples

  • subscribe()는 Publisher를 구독하고 Pipeline의 데이터 흐름을 시작함

 

예시) 단순 구독

더보기
Flux<Integer> ints = Flux.range(1, 3);

ints.subscribe();

 

예시) 데이터 처리

더보기
Flux.range(1, 3)
    .subscribe(i -> System.out.println(i));

 

예시) 데이터 + Error 처리

더보기
Flux<Integer> ints = Flux.range(1, 4)
    .map(i -> {
        if (i <= 3) return i;
        throw new RuntimeException("Got to 4");
    });

ints.subscribe(
    i -> System.out.println(i),
    error -> System.err.println("Error: " + error)
);

 

예시) 데이터 + Error + Complete 처리

더보기
Flux.range(1, 4)
    .subscribe(
        i -> System.out.println(i),
        error -> System.err.println("Error: " + error),
        () -> System.out.println("Done")
    );

 

Cancelling a subscribe() with Its Disposable

  • lambda 기반 subscribe()는 Disposable을 반환함
기능 의미 동작
dispose() 구독 취소 현재 구독을 취소하고, 이미 실행에 들어간 작업은 처리될 수 있음
swap() 현재 활성 구독을 하나만 유지 새로운 Disposable로 교체하면 기존 구독은 취소되고 새 구독으로 교체
composite() 여러 Disposable을 하나로 묶음 여러 구독을 관리하다가 한 번에 취소할 때 사용

 

예시) dispose()

더보기
Disposable disposable = flux.subscribe(System.out::println);

disposable.dispose(); // 구독 취소

 

예시) swap() - 현재 Disposable을 취소하고 다른 Disposable로 교체

더보기
Disposables.Swap swap = Disposables.swap();

// 첫 번째 구독
Disposable d1 = Flux.interval(Duration.ofSeconds(1))
        .subscribe(i -> System.out.println("A: " + i));

swap.update(d1);


// 두 번째 구독
Disposable d2 = Flux.interval(Duration.ofSeconds(1))
        .subscribe(i -> System.out.println("B: " + i));

swap.update(d2);
  • swap이라는 보관함을 만듬
  • swap.update(d1): d1을 현재 관리 대상으로 등록
  • swap.update(d2): 기존 d1 구독을 취소하고 d2를 새로운 관리 대상으로 등록

 

예시) composite()

더보기
Disposable.Composite composite = Disposables.composite();

Disposable d1 = Flux.interval(Duration.ofSeconds(1))
        .subscribe(i -> System.out.println("A: " + i));

Disposable d2 = Flux.interval(Duration.ofSeconds(1))
        .subscribe(i -> System.out.println("B: " + i));

Disposable d3 = Flux.interval(Duration.ofSeconds(1))
        .subscribe(i -> System.out.println("C: " + i));

composite.add(d1);
composite.add(d2);
composite.add(d3);

 

An Alternative to Lambdas: BaseSubscriber

  • request량, Cancellation, Lifecycle Signal을 직접 제어해야 할 때 사용하는 Subscriber 구현 방식
  • ✅ Lambda 기반 subscribe()보다 더 세밀한 기능 제공

 

예제) BaseSubscriber

더보기
public class SampleSubscriber<T> extends BaseSubscriber<T> {

    @Override
    protected void hookOnSubscribe(Subscription subscription) {
        System.out.println("Subscribed");
        request(1);
    }

    @Override
    protected void hookOnNext(T value) {
        System.out.println(value);
        request(1);
    }
}

Method 역할
hookOnSubscribe() 구독 시작
hookOnNext() 데이터 수신
hookOnError() Error
hookOnComplete() 정상 종료
hookOnCancel() Cancellation
hookFinally() 어떤 방식으로든 종료될 때 실행

 

Flux<Integer> ints = Flux.range(1, 4);

ints.subscribe(new SampleSubscriber<>());
  • subscribe → request(1) → onNext(1) → request(1) → onNext(2) → ...

 

Subscription Patterns

  • 보통 개발자는 파이프라인까지만 만들어 반환하는 역할이 주됨
  • ✅ Application Code → Mono/Flux 반환 → Spring WebFlux / Reactor Netty가 subscribe
  •  Framework가 최종 구독을 담당하는 경우가 일반적
방식 핵심 특징
Fire-and-forget 직접 subscribe()하고 반환값을 무시 취소·에러 처리를 외부 Pipeline이 제어하기 어려움
Dispatch and reference Disposable을 저장해 나중에 취소 직접 실행하되 Lifecycle 관리 가능
BaseSubscriber Subscriber를 직접 구현해 세밀하게 제어 request(n), cancel, Error/Complete 등 직접 관리

 

예시) Fire-and-forget

더보기
Mono<Void> handle(T arg) {
    sideEffectService.doSomething(arg).subscribe();
    return Mono.empty();
}
Mono<Void> handle(T arg) {
    return sideEffectService.doSomething(arg).then();
}
  • 가능하면 그대로 연결하도록 파이프라인을 반환하는 것이 좋음

 

예시) Dispatch and reference

더보기
Disposable disposable =
    sideEffectService.doSomething(arg)
        .subscribe(
            result -> System.out.println(result),
            error -> System.err.println(error)
        );
  • 작업 중 더이상 필요하지 않으면 취소하기

 

예시) BaseSubscriber

더보기
Flux.range(1, 10)
    .subscribe(new BaseSubscriber<Integer>() {

        @Override
        protected void hookOnSubscribe(Subscription subscription) {
            request(1);
        }

        @Override
        protected void hookOnNext(Integer value) {
            System.out.println(value);
            request(1);
        }

        @Override
        protected void hookOnComplete() {
            System.out.println("Done");
        }
    });
  • subscribe → request(1) → onNext(1) → request(1) → onNext(2) → ... → onComplete

 

On Backpressure and Ways to Reshape Requests

  • 중간 Operator는 downstream에서 전달된 Demand를 그대로 upstream에 전달하지 않고, 자신의 동작 방식에 맞게 변환 가능
Operator 핵심 역할 예시 흐름
buffer(10) 원본 10개를 1개 묶음으로 변환 request(2 buffers) → upstream request(20 items)
prefetch 데이터를 미리 일정량 요청 prefetch(8) → 8개 미리 확보 → 소비하며 보충
limitRate(10) 큰 Demand를 작은 Batch로 나눠 전달 request(100) → upstream에는 10개 단위 요청
limitRequest(10) Stream 전체에서 최대 10개까지만 받음 1~10 수신 → onComplete → source cancel

 

예시) buffer(10)

더보기
Flux.range(1, 100)
    .buffer(10)
    .subscribe(System.out::println);
  • 원본 데이터 10개를 List 하나로 묶는 Operator
  • [1~10] → [11~20] → [21~30] → ...

 

예시) prefetch

더보기
Flux.just("A", "B", "C")
    .flatMap(
        name -> getItems(name),
        2,   // concurrency
        8    // prefetch
    )
    .subscribe(System.out::println);

 

  • Inner Publisher → 8개 미리 요청 → 내부 Queue → Downstream에 전달
  • Queue를 소비하다 일정량 줄어들면 다시 upstream에 추가 요청해서 보충함.

 

예시) limitRate(10)

더보기
Flux.range(1, 100)
    .limitRate(10)
    .subscribe(System.out::println);
  • 100개를 모두 받을 수 있는데, upstream에 한꺼번에 100개 Demand를 주지 않고 작은 Batch로 조절하기
  • Downstream request(100) → limitRate(10) → request(10) → 보충 request → ..

 

예시) limitRequest(10)

더보기
Flux.range(1, 100)
    .limitRequest(10)
    .subscribe(System.out::println);
  • 원본은 100개지만 10개를 받은 뒤 종료함.
  • Source 1~100 → limitRequest(10) → 1~10 전달 → onComplete → Source cancel

 

4. Programmatically creating a sequence

5. Threading and Schedulers

  • 실행 Thread를 바꾸고 싶다면 Scheduler를 사용함
  • Reactor의 비동기/Reactive 구조와 Thread 전환은 별개의 개념임

 

Scheduler

  • 작업을 어떤 Thread/실행 환경에서 수행할지 결정하는 추상화
Scheduler 특징 대표 용도
Schedulers.immediate() 현재 Thread에서 그대로 실행 별도 Thread 전환 불필요
Schedulers.single() 하나의 공유 Thread 사용 순차적인 단일 Thread 작업
Schedulers.boundedElastic() 제한된 Elastic 실행 환경 Blocking I/O, Legacy Blocking 코드
Schedulers.parallel() CPU Core 수에 맞춘 Worker Pool CPU 연산 중심 작업
Schedulers.fromExecutorService() 기존 ExecutorService를 Scheduler로 사용 Custom Thread Pool 연동

 

예시) boundedElastic()

더보기
Mono.fromCallable(() -> blockingDatabaseCall())
    .subscribeOn(Schedulers.boundedElastic());
  • EventLoop Thread에서는 Blocking하면 안 됨
  • 피할 수 없는 Blocking 작업은 boundedElastic() 같은 별도 Scheduler로 이동

 

예시) parallel()

더보기
Schedulers.parallel()

 

작업 스레드 바꾸기

구분 publishOn() subscribeOn()
의미 중간부터 Thread 변경 처음부터 Thread 변경
기준 데이터가 downstream으로 흘러가는 지점 subscribe()가 시작되는 지점
영향 이후 operator Source + 전체 구독 과정
위치 Pipeline 중간에 사용 보통 어디에 놓아도 Source 쪽에 영향

 

예시) publishOn()

더보기
Flux.range(1, 2)
    .map(i -> i * 10)
    .publishOn(Schedulers.parallel())
    .map(i -> "value " + i)
    .subscribe(System.out::println);
기존 Thread → map①

              ↓ publishOn

parallel Thread → map② → Subscriber
  • downstream의 실행 위치를 변경

 

예시) subscribeOn()

더보기
Flux.range(1, 2)
    .map(i -> i * 10)
    .subscribeOn(Schedulers.parallel())
    .map(i -> "value " + i)
    .subscribe(System.out::println);
호출 Thread → subscribe()

                ↓

parallel Thread → Source → map① → map② → Subscriber
  • subscribe() → subscribeOn → Scheduler Thread로 구독 이동

 

6. Handling Errors

  • Error는 Terminal Signal임
  • ✅ 스트림이 더이상 계속되지 않음
  • ➡️ upstream은 종료되고 fallback Stream으로 대체됨

 

Error Handling Operators

구분 Reactor Operator 핵심
Static Fallback Value onErrorReturn() Error 발생 시 정해진 기본값으로 대체
Catch and swallow the error onErrorComplete() onErroronComplete로 바꿔 Error를 무시하고 종료
Fallback Method onErrorResume() Error 발생 시 다른 Publisher로 전환
Dynamic Fallback Value onErrorResume() Error 내용을 이용해 동적으로 fallback 값/Publisher 생성
Catch and Rethrow onErrorMap() 기존 Error를 다른 Exception으로 변환해 다시 전달
Log or React on the Side doOnError() Error를 변경하지 않고 Logging, Metric 등 Side Effect 수행
Using Resources and Finally doFinally(), using() 종료 시 Cleanup 수행 / Resource Lifecycle 관리
Terminal Aspect of onError Error Handling 전반 Error 발생 시 원래 upstream Sequence는 종료됨
Retrying retry(), retryWhen() Error 발생 시 종료된 upstream을 다시 Subscribe하여 재시도

 

예시) Static Fallback Value

더보기
Flux.just(10, 0, 5)
    .map(i -> 100 / i)
    .onErrorReturn(-1)
    .subscribe(System.out::println);

 

예시) Catch and swallow the error

더보기
Flux.just(10, 20, 0, 30)
    .map(i -> 100 / i)
    .onErrorComplete()
    .subscribe(
        System.out::println,
        System.err::println,
        () -> System.out.println("Done")
    );

 

예시) Fallback Method

더보기
Mono<User> user =
    callExternalApi(userId)
        .onErrorResume(error ->
            findUserFromCache(userId)
        );

 

예시) Dynamic Fallback Value

더보기
Mono<Result> result =
    callApi()
        .onErrorResume(error ->
            Mono.just(
                new Result("FAIL", error.getMessage())
            )
        );

 

예시) Catch and Rethrow

더보기
Mono<User> user =
    callExternalApi(userId)
        .onErrorMap(error ->
            new UserServiceException("사용자 조회 실패", error)
        );

 

예시) Log or React on the Side

더보기
Mono<User> user =
    callExternalApi(userId)
        .doOnError(error ->
            log.error("사용자 조회 실패", error)
        );

 

예시) Using Resources and Finally

더보기
Flux<String> flux =
    Flux.just("A", "B", "C")
        .doFinally(signalType ->
            System.out.println("종료: " + signalType)
        );
Flux<String> flux =
    Flux.using(
        () -> openResource(),
        resource -> readData(resource),
        resource -> resource.close()
    );
  • Resource의 생성 → 사용 → 해제를 하나의 Reactive 흐름으로 관리함.
  • Resource 생성 → Flux에서 사용 → Stream 종료 → Resource close

 

예시) Terminal Aspect of onError

더보기
Flux.just(1, 2, 0, 3, 4)
    .map(i -> 10 / i)
    .onErrorReturn(-1)
    .subscribe(System.out::println);
  • 0에서 에러 이후 스트림 종료 (3,4 처리 안함)

 

예시) Retrying

더보기
callExternalApi()
    .retry(2)
    .subscribe(System.out::println);

 

Handling Exceptions in Operators or Functions

  • Operator 내부에서 발생한 Runtime Exception은 자동으로 onError Signal로 변환됨
  • Checked Exception이 던져질 때는 try-catch로 직접 Runtime Exception 형태로 감싸서 onError에 전달해야 함

 

예시) 예외 전달

더보기
Flux.just("foo")
    .map(value -> {
        throw new IllegalArgumentException(value);
    })
    .subscribe(
        value -> System.out.println(value),
        error -> System.out.println("ERROR: " + error)
    );
ERROR: java.lang.IllegalArgumentException: foo
  • map 내부 throw → Reactor가 onError로 변환 → Subscriber의 error Consumer

 

예시) 예외 포장해서 onError에 전달

더보기
Flux<String> converted =
    Flux.range(1, 10)
        .map(i -> {
            try {
                return convert(i);
            } catch (IOException e) {
                throw Exceptions.propagate(e);
            }
        });

 

원래 예외 확인 하기

converted.subscribe(
    System.out::println,
    error -> {
        Throwable original = Exceptions.unwrap(error);

        if (original instanceof IOException) {
            System.out.println("I/O Error");
        }
    }
);
Method 역할
Exceptions.propagate(e) Checked Exception을 Reactor에서 전달 가능한 형태로 변환
Exceptions.unwrap(e) Reactor Wrapper에서 원래 Exception 추출

 

7. Sinks

 

출처