Mono와 Flux의 다양한 연산자

김기현·2025년 8월 1일

Spring WebFlux

목록 보기
7/28

Project Reactor의 Mono와 Flux는 다양한 연산자(operator)를 제공하여 데이터 스트림을 변환, 조합, 필터링 할 수 있도록 한다. 이 연산자들을 통해 복잡한 비동기 로직을 선언적이고
간결하게 표현할 수 있다.


subscribe(): 스트림 실행 및 결과 소비

subscribe()는 리액티브 스트림이 실제로 실행되도록 트리거하는 가장 중요한 메소드이다.
Mono나 Flux는 lazy 특성 때문에 subscribe()가 호출되기 전까지는 아무런 작업도 수행하지 않는다.

subscribe()는 다양한 오버르도된 버전을 제공한다.

  • subscribe(Consumer<T>): 간단한 데이터 처리
  • subscribe(Consumer<T>, Consumer<Throwable>): 데이터와 오류 처리
  • subscribe(Consumer<T>, Consumer<Throwable>, Runnable): 완전한 라이프사이클 제어
  • subscribe(Consumer<T>, Consumer<Throwable>, Runnable, Context) 또는 subscribe(Subscriber): 고급 제어(백프레셔, Context)

내부 코드

public class Flux {

    /*
     * * 역할: 이 오버로드는 가장 기본적인 subscribe메소드이다. 아무런 인자도 받지 않는다.
     * * 기능: 스트림을 활성화하고 작업을 시작하게 만든다. 하지만 스트림에서 발생되는 데이터, 발생할 수 있는 오류, 또는 와뇰 이벤트에 대해 어떤 특정 콜백도 제공하지 않는다.
     * * 반환값: Disposable 객체를 반환한다. 이 Disposable을 사용하여 스트림의 구독을 명시적으로 취소할 수 있다.
     *
     * 단순히 스트림의 작업을 시작하고 싶지만, 그 결과나 상태에 직접적인 관심이 없을 때(ex: 백그라운드에서 동작하고 결과를 다른 곳으로 푸쉬하는 스트림) 사용될 수 있다.
     * */
    public final Disposable subscribe() {
    }

    /*
     * * 역할: 스트림에서 성공적으로 발행되는 각 항목(onNext 이벤트)를 처리하기 위한 콜백(consumer)을 제공한다.
     * * 기능: cunsumer 람다 또는 메소드 참조를 통해 스트림이 발행하는 모든 데이터를 소비한다.
     *      * 오류 발생 시나 스트림 완료 시에 대한 처리는 명시하지 않는다.
     *      * 오류가 발생하면 기본적으로 Hook.onErrorDropped등을 통해 처리되거나, Unhandled OnError가 발생할 수 있다.
     * * 반환값: Disposable 객체를 반환하여 구독을 취소할 수 있게 한다.
     *
     * 스트림에서 오는 데이터만 필요하고 오류나 완료는 특별히 다룰 필요가 없을 때 (또는 GlobalExceptionHandler가 설정되어있을 때) 사용한다.
     *
     * */
    public final Disposable subscribe(Consumer<? super T> consumer) {
    }

    /*
     * * 역할: 스트림에서 발행되는 각 항목(cunsumer)과 오류가 발생했을 때(onError 이벤트) 처리하기 위한 콜백(errorConsumer)를 제공한다.
     * * 기능: onNext와 onError 이벤트를 모두 다룰 수 있다.
     *      * onComplete 이벤트에 대한 처리는 명시하지 않는다.
     * * 반환값: Disposable 객체를 반환한다.
     *
     * 데이터를 처리하고 발생할 수 있는 오류에 대해 구체적으로 대응해야 할 때 사용한다.
     * */
    public final Disposable subscribe(
            @Nullable Consumer<? super T> consumer,
            Consumer<? super Throwable> errorConsumer
    ) {
    }

    /*
     * * 역할: 스트림에서 발행되는 각 항목(cunsumer), 오류 발생(errorCunsumer) 그리고 스트림이 성공적으로 완료되었을 때(onComplete 이벤트) 처리하기 위한 콜백(completeConsumer)을 모두 제공한다.
     * * 기능: onNext, onError, onComplete 세 가지 핵심 이벤트를 모두 처리할 수 있는 가장 완벽한 형태의 구독 메소드이다.
     * * 반환값: Disposable 객체를 반환한다.
     *
     * 스트림 전체의 라이프사이클에 걸쳐 모든 이벤트를 완벽하게 제어하고 싶을 때 사용한다.
     * 데이터를 모두 처리한 후 특정 작업을 수행해야 할 때 completeConsumer가 유용하다.
     * */
    public final Disposable subscribe(
            @Nullable Consumer<? super T> consumer,
            @Nullable Consumer<? super Throwable> errorConsumer,
            @Nullable Runnable completeConsumer
    ) {
    }

    /**
     * * 역할: 세 가지 기본 콜백에 추가로, 구독 시 Context 객체를 초기화할 수 있도록 합니다.
     * * 기능: Project Reactor의 Context는 리액티브 체인 내에서 스레드 로컬과 유사하게 데이터를 전달할 수 있는 메커니즘을 제공합니다. 이 오버로드는 구독이 시작될 때 특정 Context를 주입할 수 있도록 합니다.
     * * 반환값: Disposable 객체를 반환한다.     
     *
     * 트레이싱 ID, 인증 정보 등 스트림의 각 연산자에 접근해야 하는 공통 컨텍스트 정보를 전달할 때 유용하다.
     * */
    public final Disposable subscribe(
            @Nullable Consumer<? super T> consumer,
            @Nullable Consumer<? super Throwable> errorConsumer,
            @Nullable Runnable completeConsumer,
            @Nullable Context initialContext
    ) {
    }

    /*
     * * 역할: Reactive Streams 사양의 Subscriber 인터페이스를 직접 구현한 객체를 받아 스트림을 구독한다.
     * * 기능: 이 메소드는 Mono와 Flux의 subscribe 오버로드들이 내부적으로 최종적으로 호출하는 핵심 메소드이다.
     *      * 제공된 Subscriber가 스트림에서 발생하는 모든 이벤트(onSubscribe, onNext, onError, onComplete)를 직접 처리한다.
     * * 반환값: void (Disposable을 직접 반환하지 않는다. Subscriber 내에서 Subscription을 통해 관리한다)
     * */
    public final void subscribe(Subscriber<? super T> actual) {
    }
}

예제


@Slf4j
@SpringBootTest
public class SubscribeTest {

    @Test
    void subscribe() {
        Flux.just("Apple", "Banana", "Cherry")
                .subscribe(
                        item -> log.info("Received: {}", item),                 // onNext
                        error -> log.error("error: {}", error.getMessage()),    // onError
                        () -> log.info("Flux Completed")                        // onComplete
                );

        log.info("----------------------------------");

        Mono.empty()
                .subscribe(
                        data -> log.info("Empty Mono: {}", data),
                        error -> log.error("Empty Mono error: {}", error.getMessage()),
                        () -> log.info("Empty Mono Completed")
                );
    }
}

map(): 각 항목을 동기적으로 변환

map()연산자는 스트림의 각 항목에 대해 동기적인 변환 함수를 적용하여 새로운 타입 또는 값으로 매핑한다.
입력 스트림의 각 요소가 변환되어 다음 스트림으로 전달된다.

예제


@Slf4j
@SpringBootTest
public class OperatorTest {

    @Test
    void map() {
        Flux.just("Apple", "Banana", "Cherry")
                .map(String::toUpperCase)
                .subscribe(log::info);

        log.info("--------------------------------------");

        Flux.range(1, 10)
                .map(n -> String.valueOf(n * 2))
                .subscribe(log::info);
    }
}

flatMap(): 각 항목을 비동기적으로 변환(스트림을 평면화)

flatMap()연산자는 map()과 유사하게 각 항목에 변환 함수를 적용하지만 이 변환 함수는 Mono 또는 Flux와 같은 새로운 리액티브 스트림을 반환해야 한다.
flatMap()은 이러한 내부 스트림들을 병합(flatten)하여 단일 스트림으로 만든다. 이는 비동기 작업의 체이닝에 매우 유용하다.

주의점

  • flatMap()은 내부 스트림 발생 순서를 보장하지 않는다. 병렬적으로 처리되므로 순서가 섞일 수 있다.
  • 순서가 중요한 경우 concatMap()을 고려해야 한다.

예제


@Slf4j
@SpringBootTest
public class OperatorTest {

    @Test
    void flatMap() throws InterruptedException {
        Flux.just(1, 2, 3)
                .flatMap(OperatorTest::fetchUserById)
                .subscribe(
                        user -> log.info("Fetched user: {}", user),
                        error -> log.error("Error: {}", error.getMessage()),
                        () -> log.info("All users fetched")
                );

        Thread.sleep(500);
    }

    private static Mono<String> fetchUserById(int id) {
        if (id == 1) return Mono.just("Kim").delayElement(Duration.ofMillis(100));   // 0.1초
        if (id == 2) return Mono.just("Lee").delayElement(Duration.ofMillis(50));    // 0.05초
        return Mono.empty();
    }
}

filter(): 조건에 따라 항목을 필터링

filter()연산자는 주어진 Predicate(조건 함수)를 만족하는 항목만 다음 스트림으로 전달한다.
조건에 맞지 않는 항목은 스트림에서 제외된다.

예제


@Slf4j
@SpringBootTest
public class OperatorTest {

    @Test
    void filter() {
        Flux.range(1, 10)
                .filter(n -> n % 2 == 0)    // 짝수만 통과
                .map(String::valueOf)
                .subscribe(log::info);

        log.info("-----------------------------");

        // 길이가 5보다 긴 단어만 필터링
        Flux.just("apple", "banana", "kiwi", "grapefruit")
                .filter(word -> word.length() > 5)
                .subscribe(log::info);
    }
}

zip(): 여러 스트림의 항목을 결함(인덱스 기반)

zip()연산자는 여러 Publisher(Mono 또는 Flux)의 항목을 인덱스 기반으로 결합한다.
각 Publisher에서 같은 인덱스에 있는 항목들을 가져와 하나의 튜플(Tuple)로 묶거나 주어진 BiFunction또는 Function을 사용하여 새로운 단일 항목으로 변환한다.
가장 짧은 스트림이 완료되면 zip도 완료된다.

예제


@Slf4j
@SpringBootTest
public class OperatorTest {

    @Test
    void zip() {
        Flux<String> names = Flux.just("Kim", "Lee", "Park");
        Flux<Integer> ages = Flux.just(30, 25, 35, 40); // ages가 names보다 길다

        // 이름과 나이를 결합
        Flux.zip(names, ages, (name, age) -> name + "는 " + age + "살 입니다.")
                .subscribe(log::info);
        // 출력은 가잘 짧은 스트림인 names에 맞춰 3개만 나온다

        log.info("------------------------------------------------");

        Flux<String> colors = Flux.just("Red", "Green", "Blue");
        Flux<String> fruits = Flux.just("Apple", "Grape", "Orange");
        Flux<Integer> counts = Flux.just(1, 2, 3);

        Flux.zip(colors, fruits, counts)
                .map(tuple -> fruitFormat(tuple.getT1(), tuple.getT2(), tuple.getT3()))
                .subscribe(log::info);
    }

    private String fruitFormat(String color, String fruit, int count) {
        return String.format("to: %s, Fruit: %s, Count: %d", color, fruit, count);

    }
}

merge(): 여러 스트림을 병학(시간 기반)

merge()연산자는 여러 Publisher의 항목을 시간 순서에 따라 병합하여 하나의 스트림으로 만든다.
각 소스 스트림에서 항목이 발행되는 즉시 결과 스트림으로 전달된다.
zip()과 달리 순서를 보장하지 않고 발행되는 대로 병합한다.

주의점

merge()는 병렬 처리가 가능하며 입력 스트림 중 하나라도 오류가 발생하면 전체 merge스트림도 오류를 발생시킨다.

예제


@Slf4j
@SpringBootTest
public class OperatorTest {

    @Test
    void merge() throws InterruptedException {
        // Flux 1: 0.1초마다 "A" 발행 (3번)
        Flux<String> flux1 = Flux.interval(Duration.ofMillis(100))
                .map(i -> "A" + i)
                .take(3); // 3개만 가져옴

        // Flux 2: 0.15초마다 "B" 발행 (3번)
        Flux<String> flux2 = Flux.interval(Duration.ofMillis(150))
                .map(i -> "B" + i)
                .take(3); // 3개만 가져옴

        // 두 Flux를 병합
        Flux.merge(flux1, flux2)
                .subscribe(
                        item -> log.info("Merged: {}", item),
                        error -> log.info("Error: {}", String.valueOf(error)),
                        () -> log.info("Merge Completed!")
                );

        // 비동기 작업이 완료될 때까지 메인 스레드 대기
        Thread.sleep(1000); // 1초 대기
    }
}

기타 중요한 연산자들

concatMap()

  • flatMap()과 유사하지만 내부 스트림의 결과를 순서를 보장하여 병합한다.
  • 즉 이전 내부 스트림이 완료될 때까지 다음 내부 스트림의 발행을 기다린다.
  • 순서가 중요할 때 사용한다.

onErrorReturn(T fallbackValue)

  • 스트림에서 오류가 발생했을 때 지정된 대체 값으로 스트림을 완료한다.

onErrorResume(Function<Throwable, Mono> fallbackMono)

  • 스트림에서 오류가 발생했을 때 다른 Mono또는 Flux로 대체하여 스트림을 계속 진행한다.

retry(long numRetries)

  • 스트림에서 오류가 발생했을 때 지정된 횟수만큼 스트림을 재구독하여 재시도한다.

delayElement(Duration duration), delayElements(Duration duration)

  • 복수형인 메소드는 Flux이고 단수형인 것은 Mono의 것이다.
  • 항목을 다음 스트림으로 전달하기 전에 지정된 시간만큼 지연시킨다.

then(), thenMany()

  • then(): Mono
  • thenMany(): Mono, Flux
  • 현재 스트림의 완료를 기다린 후 다른 Mono또는 Flux를 실행한다.
  • 이전 스트림의 결과값은 무시된다.

block()

  • 리액티브 스트림을 블로킹 방식으로 실행하여 결과를 동기적으로 얻는다.
  • 주로 테스트나 간단한 예제에서 사용하며 실제 운영에서는 비동기의 이점을 잃으므로 사용하면 안 된다.
profile
백엔드 개발자를 목표로 공부하는 대학생

0개의 댓글