Java

[Project Reactor] 4. Which operator do I need?

noahkim_ 2026. 8. 20. 00:33

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

 

출처