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 / Iterable을 Flux로 변환 | 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() | onError를 onComplete로 바꿔 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
출처
'Java' 카테고리의 다른 글
| [Project Reactor] 1. Introduction to Reactive Programming (0) | 2026.08.17 |
|---|---|
| [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 |