1. Creating a New Sequence…
- Lazy하게 생성하는 방식을 제공함
- ✅ Subscription 시점에 실제 값을 생성함
API
| API | 설명 | 핵심 |
| Mono.fromCallable() | 반환값이 있는 작업을 Subscription 시점에 실행하여 Mono<T>로 생성 | 반환값 → onNext, Exception → onError |
| Mono.fromRunnable() | 반환값이 없는 작업을 Reactive Sequence로 생성 | 값 없이 완료 여부만 표현 |
예제) Mono.fromCallable()
더보기
Mono<User> userMono =
Mono.fromCallable(() -> userRepository.findById(1L))
.subscribeOn(Schedulers.boundedElastic());
- Blocking 작업이나 기존 동기 메서드를 Mono로 감쌀 때 유용함
2. Transforming an Existing Sequence
- Operator는 기존 Publisher를 수정하는 것이 아닌 새로운 Publisher를 반환함
API
| API | 설명 | 핵심 |
| map() | 기존 데이터를 1:1로 동기 변환 | T → R |
| mapNotNull() | 변환 결과가 null이면 해당 값을 제외 | 변환 + null 제거 |
| flatMap() | 각 값을 Publisher로 변환한 뒤 결과를 합침 | 비동기 작업 조합, 순서 보장 X |
| concatMap() | 각 Publisher를 하나씩 순차적으로 처리 | 순서 보장 O |
| then() | 앞 Sequence의 값을 무시하고 완료만 전달 | Mono<Void> |
| thenReturn(T) | 앞 Sequence가 정상 완료되면 지정 값 하나 반환 | 완료 후 고정 값 반환 |
| thenMany(Publisher) | 앞 Sequence가 완료되면 다른 Sequence로 전환 | 완료 후 새로운 Flux 등 실행 |
| repeatWhen() | 정상 완료 후 조건에 따라 다시 Subscribe | onComplete 기반 반복 |
예제) map()
더보기
Mono<UserDto> result =
userService.findById(1L)
.map(user -> new UserDto(user.getId(), user.getName()));
- 값 → 값 변환
- Mono<User> → map → Mono<UserDto>
예제) flatMap()
더보기
Mono<Order> result =
userService.findById(1L)
.flatMap(user -> orderService.findLatestOrder(user.getId()));
- Reactive 작업 → Reactive 작업 연결
- Mono<User> → User로 다른 Mono 생성 → Mono<Order>
예제) concatMap()
더보기
Flux<Integer> result =
Flux.just(1, 2, 3)
.concatMap(id ->
webClient.get()
.uri("/users/{id}", id)
.retrieve()
.bodyToMono(Integer.class)
);
- 여러 비동기 작업을 순서대로 하나씩 처리해야 할 때 유용 (순서 보장)
- 1번 요청 완료 → 2번 요청 → 완료 → 3번 요청
예제) then()
더보기
Mono<Void> result =
userService.deleteUser(1L)
.then();
- 앞 작업의 반환값은 필요 없고 완료됐는지만 중요할 때 유용
- ex) Delete, 저장 완료 처리 등
- 작업 완료 → Mono<Void>
예제) thenReturn()
더보기
Mono<String> result =
userService.deleteUser(1L)
.thenReturn("삭제 완료");
- 앞 작업이 끝난 다음 정해진 결과값 하나를 반환하고 싶을 때.
3. Peeking into a Sequence
- doOn* 계열은 원래 데이터 흐름을 변경하지 않고 중간 상태를 관찰하는 Side Effect 용도임.
- ex) Logging, Metric, Debugging 등에 주로 사용함.
- doFinally()는 complete, error, cancel 등 모든 종료 형태를 관찰할 수 있음.
4. Filtering a Sequence
- 원본 Sequence에서 조건에 맞는 일부 데이터만 downstream으로 전달함.
- 값의 조건뿐 아니라 개수나 시간 기준으로도 흐름을 제한할 수 있음. (ex. take, skip, sample)
| API | 설명 | 핵심 |
| sample(Duration) | 일정 시간 구간마다 해당 구간의 가장 최근 값 하나만 전달 | 빠른 데이터 흐름의 빈도 감소 |
예제) sample()
더보기
Flux<Price> sampled =
priceStream
.sample(Duration.ofSeconds(1));
- 실시간 데이터가 너무 많이 들어올 때 일부만 내려보낼 때
- 실시간 가격/센서/UI 갱신 빈도 제한에 유용
- 100 → 101 → 105 → 103이 1초 동안 들어오면 103만 전달됨
5. Handling Errors
- retry는 Error 이후 기존 흐름을 이어가는 것이 아니라 Upstream에 다시 Subscribe하는 방식임.
- Retry.backoff()는 재시도 간격을 점차 늘려 외부 시스템에 반복적으로 부하를 주는 것을 방지함.
- Backpressure 상황에서는 Buffer, Drop, Latest, Error 등의 전략을 선택할 수 있음.
| API | 설명 | 핵심 |
| Retry.backoff() | 재시도할수록 대기시간을 증가시키는 Exponential Backoff 정책 생성 | onError → 대기 → 재구독 |
| maxBackoff() | Backoff 대기시간이 지나치게 증가하지 않도록 최대 대기시간 설정 | 최대 Retry 간격 제한 |
| jitter() | Backoff 시간에 일정한 무작위 편차를 추가 | 여러 요청이 동시에 재시도하는 현상 완화 |
| doBeforeRetry() | 실제 재시도가 실행되기 직전에 부가 작업 수행 | Retry Logging / Metric 등에 사용 |
| onBackpressureLatest() | Downstream이 처리하지 못한 데이터는 버리고 가장 최근 값만 유지 | 오래된 값 제거 + 최신 값 유지 |
예제) Retry.backoff()
더보기
webClient.get()
.uri("/users/1")
.retrieve()
.bodyToMono(User.class)
.retryWhen(
Retry.backoff(3, Duration.ofSeconds(1))
);
- 실패 → 1초 후 재시도 → 또 실패 → 더 긴 시간 후 재시도 → 최대 3회
예제) maxBackoff()
더보기
.retryWhen(
Retry.backoff(5, Duration.ofSeconds(1))
.maxBackoff(Duration.ofSeconds(10))
)
- 1초 → 2초 → 4초 → 8초 → 최대 10초
- ➡️ Backoff 시간이 계속 커지는 것을 제한
예제) Retry.backoff()
더보기
webClient.get()
.uri("/users/1")
.retrieve()
.bodyToMono(User.class)
.retryWhen(
Retry.backoff(3, Duration.ofSeconds(1))
);
- 실패 → 1초 후 재시도 → 또 실패 → 더 긴 시간 후 재시도 → 최대 3회
예제) jitter()
더보기
.retryWhen(
Retry.backoff(5, Duration.ofSeconds(1))
.jitter(0.5)
)
- 원래 재시도 시간이 정확히 1초 → 2초 → 4초 였다면 약간의 랜덤 편차를 줌.
- ➡️ 여러 요청이 같은 순간에 한꺼번에 Retry하는 것을 완화.
예제) doBeforeRetry()
더보기
.retryWhen(
Retry.backoff(3, Duration.ofSeconds(1))
.doBeforeRetry(signal ->
log.warn("Retry 횟수: {}", signal.totalRetries() + 1)
)
)
- Error → doBeforeRetry() 실행 → Retry
- ➡️ 재시도 직전에 Logging, Metric 등을 넣을 때 사용.
예제) onBackpressureLatest()
더보기
Flux<Integer> flux =
Flux.range(1, 1000)
.onBackpressureLatest();
- Producer는 1 → 2 → 3 → 4 → 5 → ... 처럼 빠르게 보내는 상황
- Consumer가 느리면, 오래된 대기값을 계속 쌓는 대신 가장 최근 값만 유지함.
- ex) 2, 3, 4, 5 버림 → 최신 6 유지
- ➡️ 실시간 가격, 센서값처럼 중간값보다 최신 상태가 중요한 경우에 적합.
6. Working with Time
- Reactor는 데이터 흐름의 시간도 Operator로 제어할 수 있음.
- delayElements()는 Subscription 자체를 늦추는 게 아니라 각 onNext 데이터가 downstream으로 전달되는 시간을 지연시킴.
| API | 설명핵 | 심 |
| delayElements(Duration) | 각각의 onNext Signal 사이에 지정한 시간만큼 지연을 적용 | 데이터 방출 간격 조절 |
예제) delayElements()
더보기
Flux.just("A", "B", "C")
.delayElements(Duration.ofSeconds(1))
.subscribe(System.out::println);
- 1초 → A → 1초 → B → 1초 → C
- ➡️ 각 데이터의 전달을 일정 시간씩 늦춤.
7. Splitting a Flux
8. Going Back to the Synchronous World
9. Multicasting a Flux to several Subscribers
출처
'Java' 카테고리의 다른 글
| [Project Reactor] 3. Three Sorts of Batching (0) | 2026.08.20 |
|---|---|
| [Project Reactor] 2. Reactor Core Features (0) | 2026.08.18 |
| [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 |