Spring WebFlux Mono·Flux와 Reactor 연산자 반환 타입 정리

최병현·2026년 7월 27일

spring boot

목록 보기
34/34

Spring WebFlux를 사용하면 반환 타입으로 일반 객체나 List 대신 MonoFlux를 자주 사용한다.

처음에는 연산자 이름보다도 다음 부분이 가장 헷갈린다.

  • map을 사용해야 하는지 flatMap을 사용해야 하는지
  • 연산자를 사용한 뒤 반환 타입이 Mono인지 Flux인지
  • 데이터가 없을 때 어떻게 처리해야 하는지
  • Flux<T>Mono<List<T>>로 어떻게 변환하는지

이번 글에서는 Mono, Flux의 기본 개념과 함께 실무에서 자주 사용하는 Reactor 연산자를 입력 타입과 반환 타입 중심으로 정리한다.


1. Mono와 Flux

Reactor의 Publisher는 대표적으로 MonoFlux로 나뉜다.

Mono<T>  : 0개 또는 1개의 데이터
Flux<T>  : 0개 이상의 여러 데이터

공식 문서에서 Mono는 최대 하나의 값을 발행하는 Publisher로 정의된다. 값 없이 완료될 수도 있으며, 여러 값을 반환해야 하는 경우에는 Flux를 사용한다.


2. Mono

Mono<T>는 데이터가 없거나 최대 1개인 경우 사용한다.

Mono<User>
Mono<Product>
Mono<Order>
Mono<Void>

주로 다음과 같은 상황에서 사용한다.

Mono<User> findById(Long userId);

Mono<Product> save(Product product);

Mono<Void> deleteById(Long productId);

Mono 데이터 흐름

데이터 1개
   ↓
Mono<User>
   ↓
User 전달 후 완료

데이터가 존재하지 않으면 값 없이 완료된다.

Mono.empty()
   ↓
onComplete

3. Flux

Flux<T>는 0개 이상의 여러 데이터를 순차적으로 발행한다.

Flux<User>
Flux<Product>
Flux<Message>

주로 목록 조회나 실시간 스트림에 사용한다.

Flux<User> findAll();

Flux<Message> findByRoomId(String roomId);

Flux<ServerSentEvent<MessageResponse>> subscribe(String roomId);

Flux 데이터 흐름

User1
User2
User3
  ↓
Flux<User>

Flux는 모든 데이터를 한 번에 List로 반환하는 것이 아니라 데이터를 하나씩 발행한다.

onNext(User1)
onNext(User2)
onNext(User3)
onComplete()

4. map

map은 Publisher 내부의 일반 데이터를 다른 일반 데이터로 변환할 때 사용한다.

T → R

즉, 변환 함수의 반환값이 MonoFlux가 아닌 일반 객체일 때 사용한다.

공식 Reactor 문서에서도 map은 기존 데이터를 1대1로 변환할 때 사용하는 연산자로 설명한다.


Mono에서 map

Mono<T>
    .map(T -> R)

반환 타입은 다음과 같다.

Mono<T> → Mono<R>

예시:

Mono<User> userMono = userRepository.findById(userId);

Mono<UserResponse> responseMono = userMono
        .map(user -> new UserResponse(
                user.getId(),
                user.getName()
        ));

반환 타입:

Mono<UserResponse>

흐름:

Mono<User>
   ↓ map
User → UserResponse
   ↓
Mono<UserResponse>

Flux에서 map

Flux<T>
    .map(T -> R)

반환 타입:

Flux<T> → Flux<R>

예시:

Flux<UserResponse> responses = userRepository.findAll()
        .map(user -> new UserResponse(
                user.getId(),
                user.getName()
        ));

흐름:

User1 → UserResponse1
User2 → UserResponse2
User3 → UserResponse3

최종 반환 타입:

Flux<UserResponse>

map을 사용하는 기준

람다식 내부에서 일반 객체를 반환한다면 map을 사용한다.

.map(user -> UserResponse.from(user))

여기서 UserResponse.from(user)의 반환 타입은 다음과 같다.

UserResponse

따라서 map을 사용한다.


5. flatMap

flatMap은 Publisher 내부의 값을 사용해 새로운 Mono 또는 Flux를 호출할 때 사용한다.

T → Publisher<R>

즉, 람다식 내부의 반환 타입이 다음 중 하나라면 flatMap을 사용한다.

Mono<R>
Flux<R>
Publisher<R>

Reactor 공식 문서에서는 비동기 작업이나 Publisher를 반환하는 작업을 연결할 때 flatMap을 사용한다고 설명한다.


Mono에서 flatMap

Mono<T>
    .flatMap(T -> Mono<R>)

반환 타입:

Mono<T> → Mono<R>

예시:

Mono<User> userMono = userRepository.findById(userId);

Mono<Order> orderMono = userMono
        .flatMap(user -> orderRepository.findLatestByUserId(user.getId()));

orderRepository.findLatestByUserId()가 반환하는 타입은 다음과 같다.

Mono<Order>

따라서 map이 아니라 flatMap을 사용한다.


map을 잘못 사용한 경우

Mono<Mono<Order>> result = userMono
        .map(user -> orderRepository.findLatestByUserId(user.getId()));

map은 람다식의 반환값을 그대로 감싼다.

Mono<User>
   ↓ map
Mono<Order>
   ↓
Mono<Mono<Order>>

Publisher가 중첩된다.

Mono<Mono<Order>>

이 중첩 구조를 평탄화하는 연산자가 flatMap이다.

Mono<Order> result = userMono
        .flatMap(user -> orderRepository.findLatestByUserId(user.getId()));
Mono<User>
   ↓ flatMap
Mono<Order>
   ↓
Mono<Order>

핵심 비교

.map(user -> UserResponse.from(user))

람다식 반환 타입:

UserResponse

최종 반환 타입:

Mono<UserResponse>

반면 다음 코드는:

.flatMap(user -> userRepository.save(user))

람다식 반환 타입:

Mono<User>

최종 반환 타입:

Mono<User>

Flux에서 flatMap

Flux<T>
    .flatMap(T -> Publisher<R>)

반환 타입:

Flux<T> → Flux<R>

예시:

Flux<Order> orders = userRepository.findAll()
        .flatMap(user -> orderRepository.findByUserId(user.getId()));

여기서 반환 타입을 보면 다음과 같다.

userRepository.findAll()Flux<User>

orderRepository.findByUserId(...)Flux<Order>

전체 결과는 다음과 같다.

Flux<Order>

flatMap의 순서

flatMap은 내부 Publisher를 동시에 처리할 수 있기 때문에 원본 순서가 보장되지 않을 수 있다.

Flux<Integer> result = Flux.just(1, 2, 3)
        .flatMap(number -> asyncProcess(number));

비동기 처리 시간이 다르면 결과가 다음처럼 올 수 있다.

2
1
3

순서가 중요하다면 concatMap을 사용할 수 있다.

Flux<Integer> result = Flux.just(1, 2, 3)
        .concatMap(number -> asyncProcess(number));

concatMap은 이전 Publisher가 완료된 후 다음 Publisher를 처리하므로 순서를 유지한다.

1
2
3

6. flatMapMany

Mono 내부의 하나의 값을 사용해 여러 데이터를 반환해야 할 때 사용한다.

Mono<T>
    .flatMapMany(T -> Publisher<R>)

반환 타입:

Mono<T> → Flux<R>

Mono.flatMap()은 최대 하나의 값을 유지하기 때문에 반환 타입이 Mono<R>이다. 여러 값을 반환하려면 flatMapMany()를 사용한다.

예시:

Mono<User> userMono = userRepository.findById(userId);

Flux<Order> orders = userMono
        .flatMapMany(user ->
                orderRepository.findByUserId(user.getId())
        );

흐름:

Mono<User>
   ↓ flatMapMany
Flux<Order>

7. filter

filter는 조건에 맞는 데이터만 통과시킨다.

filter(T -> boolean)

조건이 true이면 데이터를 유지하고, false이면 데이터를 제거한다.


Mono에서 filter

Mono<T>
    .filter(T -> boolean)

반환 타입:

Mono<T> → Mono<T>

예시:

Mono<User> activeUser = userRepository.findById(userId)
        .filter(user -> user.getStatus() == UserStatus.ACTIVE);

사용자가 활성 상태이면:

Mono<User>

활성 상태가 아니면:

Mono.empty()

filter는 예외를 발생시키지 않는다.

조건에 맞지 않으면 단순히 빈 Publisher가 된다.


Flux에서 filter

Flux<T>
    .filter(T -> boolean)

반환 타입:

Flux<T> → Flux<T>

예시:

Flux<Product> availableProducts = productRepository.findAll()
        .filter(Product::isAvailable);

흐름:

Product1 true  → 통과
Product2 false → 제거
Product3 true  → 통과

최종 결과:

Flux<Product>

8. switchIfEmpty

switchIfEmpty는 Publisher가 비어 있을 때 다른 Publisher로 전환한다.

switchIfEmpty(Publisher<? extends T> alternative)

중요한 점은 null일 때 동작하는 것이 아니라 Mono.empty() 또는 빈 Flux일 때 동작한다는 것이다.


Mono에서 switchIfEmpty

Mono<T>
    .switchIfEmpty(Mono<T>)

반환 타입:

Mono<T> → Mono<T>

예시:

Mono<User> user = userRepository.findById(userId)
        .switchIfEmpty(
                Mono.error(new UserNotFoundException())
        );

데이터가 존재하면:

Mono<User>

데이터가 없으면:

Mono.error(UserNotFoundException)

기본값 반환

Mono<User> user = userRepository.findById(userId)
        .switchIfEmpty(
                Mono.just(User.guest())
        );

데이터가 없으면 기본 사용자 객체를 반환한다.


Flux에서 switchIfEmpty

Flux<T>
    .switchIfEmpty(Flux<T>)

반환 타입:

Flux<T> → Flux<T>

예시:

Flux<Product> products = productRepository.findByCategory(categoryId)
        .switchIfEmpty(
                productRepository.findPopularProducts()
        );

카테고리 상품이 있으면 해당 상품을 반환하고, 없다면 인기 상품 Flux로 전환한다.


9. defaultIfEmpty

defaultIfEmpty는 Publisher가 비었을 때 Publisher가 아닌 기본값 하나를 반환한다.

defaultIfEmpty(T defaultValue)

예시:

Mono<String> name = userRepository.findById(userId)
        .map(User::getName)
        .defaultIfEmpty("알 수 없는 사용자");

switchIfEmpty와 차이는 다음과 같다.

.defaultIfEmpty("기본값")

일반 값을 전달한다.

.switchIfEmpty(Mono.just("기본값"))

Publisher를 전달한다.

DB 조회나 비동기 작업으로 대체해야 한다면 switchIfEmpty를 사용한다.

.switchIfEmpty(userRepository.findGuestUser())

단순한 기본값만 필요하면 defaultIfEmpty를 사용한다.

.defaultIfEmpty("Guest")

10. collectList

collectListFlux<T>가 발행하는 모든 데이터를 모아서 List<T>로 만든다.

Flux<T>
    .collectList()

반환 타입:

Flux<T> → Mono<List<T>>

Reactor 공식 API에서도 collectList()는 Flux의 모든 요소를 List에 모은 후, 해당 List를 하나의 값으로 발행하는 Mono를 반환한다고 설명한다. Flux가 비어 있으면 빈 List를 발행한다.

예시:

Mono<List<User>> users = userRepository.findAll()
        .collectList();

흐름:

Flux<User>
 ├─ User1
 ├─ User2
 └─ User3
      ↓ collectList
Mono<List<User>>

최종적으로 하나의 List 객체를 발행하기 때문에 Mono가 된다.


collectList 이후 map 사용

Mono<UserListResponse> response = userRepository.findAll()
        .map(UserResponse::from)
        .collectList()
        .map(UserListResponse::new);

각 단계의 반환 타입은 다음과 같다.

userRepository.findAll()
→ Flux<User>

.map(UserResponse::from)
→ Flux<UserResponse>

.collectList()
→ Mono<List<UserResponse>>

.map(UserListResponse::new)
→ Mono<UserListResponse>

Flux가 비어 있는 경우

Flux<User> emptyFlux = Flux.empty();

Mono<List<User>> result = emptyFlux.collectList();

결과는 Mono.empty()가 아니다.

Mono<List<User>>

내부 값은 빈 List다.

List.of()

따라서 다음 코드는 보통 동작하지 않는다.

userRepository.findAll()
        .collectList()
        .switchIfEmpty(Mono.error(new RuntimeException()));

collectList()는 Flux가 비어 있어도 빈 List를 발행하기 때문에 Mono 자체는 비어 있지 않다.

목록이 비어 있는지 검사하려면 다음처럼 처리해야 한다.

userRepository.findAll()
        .collectList()
        .filter(users -> !users.isEmpty())
        .switchIfEmpty(
                Mono.error(new UserNotFoundException())
        );

하지만 일반적인 목록 조회 API에서는 빈 목록을 정상 응답으로 반환하는 것이 자연스럽다.

[]

11. collectMap

collectMapFlux<T>의 데이터를 Map<K, V>로 수집한다.

Flux<T>
    .collectMap(keyMapper)

반환 타입:

Flux<T> → Mono<Map<K, T>>

예시:

Mono<Map<Long, User>> userMap = userRepository.findAll()
        .collectMap(User::getId);

결과:

{
    1: User1,
    2: User2,
    3: User3
}

값도 직접 지정할 수 있다.

Mono<Map<Long, String>> userNameMap = userRepository.findAll()
        .collectMap(
                User::getId,
                User::getName
        );

반환 타입:

Mono<Map<Long, String>>

12. next

next()Flux<T>에서 첫 번째 데이터만 가져와 Mono<T>로 변환한다.

Flux<T>
    .next()

반환 타입:

Flux<T> → Mono<T>

예시:

Mono<Message> latestMessage = messageRepository
        .findByRoomIdOrderByCreatedAtDesc(roomId)
        .next();

Flux에 데이터가 없다면:

Mono.empty()

13. then

then()은 이전 Publisher의 데이터를 무시하고 완료 신호만 전달한다.

Publisher<T>
    .then()

반환 타입:

Mono<T> → Mono<Void>
Flux<T> → Mono<Void>

예시:

Mono<Void> result = userRepository.deleteById(userId)
        .then();

또는 저장 결과를 무시하고 다음 작업을 실행할 수 있다.

Mono<Order> result = paymentRepository.save(payment)
        .then(orderRepository.save(order));

흐름:

payment 저장 완료
   ↓ 데이터 무시
order 저장 실행
   ↓
Mono<Order>

thenReturn

이전 작업이 완료되면 지정한 값을 반환한다.

Mono<T>
    .thenReturn(R)

반환 타입:

Mono<T> → Mono<R>

예시:

Mono<Boolean> result = userRepository.save(user)
        .thenReturn(true);

thenMany

이전 Publisher 완료 후 여러 데이터를 반환하는 Publisher를 실행한다.

Mono<T>
    .thenMany(Flux<R>)

반환 타입:

Mono<T> → Flux<R>

예시:

Flux<Order> result = auditRepository.save(audit)
        .thenMany(orderRepository.findAll());

14. doOnNext

doOnNext는 데이터를 변경하지 않고 중간 동작을 실행한다.

주로 로그 출력이나 디버깅에 사용한다.

Mono<User> result = userRepository.findById(userId)
        .doOnNext(user ->
                log.info("조회된 사용자: {}", user.getId())
        );

반환 타입은 변하지 않는다.

Mono<T> → Mono<T>
Flux<T> → Flux<T>

doOnNext 내부에서 객체를 반환할 수 없다.

.doOnNext(user -> UserResponse.from(user))

이 코드는 변환 결과를 사용하지 않는다.

객체 변환이 목적이면 map을 사용해야 한다.

.map(UserResponse::from)

15. doOnSuccess

doOnSuccessMono가 성공적으로 완료됐을 때 실행된다.

Mono<User> result = userRepository.save(user)
        .doOnSuccess(savedUser ->
                log.info("사용자 저장 완료")
        );

주의할 점은 Mono.empty()도 정상 완료이므로 doOnSuccess가 호출될 수 있다는 것이다.

이 경우 전달되는 값은 null일 수 있다.

.doOnSuccess(user -> {
    if (user != null) {
        log.info("사용자: {}", user.getId());
    }
});

16. doOnError

에러가 발생했을 때 부수 효과를 실행한다.

Mono<User> result = userRepository.findById(userId)
        .doOnError(error ->
                log.error("사용자 조회 실패", error)
        );

doOnError는 에러를 처리하거나 제거하지 않는다.

로그만 실행한 뒤 에러는 그대로 다음 단계로 전달된다.


17. onErrorReturn

에러가 발생하면 지정한 기본값을 반환한다.

Mono<User> result = userRepository.findById(userId)
        .onErrorReturn(User.guest());

반환 타입:

Mono<T> → Mono<T>
Flux<T> → Flux<T>

간단한 기본값을 반환할 때 사용한다.


18. onErrorResume

에러가 발생했을 때 다른 Publisher로 전환한다.

Mono<User> result = userRepository.findById(userId)
        .onErrorResume(error ->
                fallbackUserRepository.findGuestUser()
        );

람다식에서 Mono 또는 Flux를 반환한다.

Throwable → Publisher<T>

예외 종류에 따라 분기할 수도 있다.

Mono<User> result = userRepository.findById(userId)
        .onErrorResume(UserNotFoundException.class,
                error -> Mono.just(User.guest())
        );

19. zip

서로 독립적인 여러 Publisher의 결과를 결합할 때 사용한다.

Mono<User> userMono = userRepository.findById(userId);
Mono<Profile> profileMono = profileRepository.findByUserId(userId);

Mono<UserDetailResponse> response = Mono.zip(
        userMono,
        profileMono
).map(tuple -> new UserDetailResponse(
        tuple.getT1(),
        tuple.getT2()
));

반환 타입:

Mono<User> + Mono<Profile>
→ Mono<Tuple2<User, Profile>>
→ Mono<UserDetailResponse>

두 작업이 서로 의존하지 않는다면 순차적인 flatMap보다 zip을 통해 함께 실행하는 구조를 고려할 수 있다.

userRepository.findById(userId)
        .flatMap(user ->
                profileRepository.findByUserId(userId)
                        .map(profile ->
                                new UserDetailResponse(user, profile)
                        )
        );

위 코드는 사용자 조회 후 프로필 조회가 실행되는 순차 구조다.

반면:

Mono.zip(userMono, profileMono)

두 조회가 서로 독립적이면 동시에 구독해 결과를 결합할 수 있다.


20. hasElement와 hasElements

Publisher에 데이터가 존재하는지 확인한다.

Mono.hasElement

Mono<T>
    .hasElement()

반환 타입:

Mono<T> → Mono<Boolean>

예시:

Mono<Boolean> exists = userRepository.findById(userId)
        .hasElement();

Flux.hasElements

Flux<T>
    .hasElements()

반환 타입:

Flux<T> → Mono<Boolean>

예시:

Mono<Boolean> hasMessages = messageRepository.findByRoomId(roomId)
        .hasElements();

21. count

Flux가 발행한 데이터 개수를 반환한다.

Flux<T>
    .count()

반환 타입:

Flux<T> → Mono<Long>

예시:

Mono<Long> messageCount = messageRepository.findByRoomId(roomId)
        .count();

22. distinct

중복 데이터를 제거한다.

Flux<String> names = Flux.just(
        "Kim",
        "Lee",
        "Kim",
        "Park"
).distinct();

결과:

Kim
Lee
Park

반환 타입:

Flux<T> → Flux<T>

특정 필드를 기준으로 중복을 제거할 수도 있다.

Flux<User> users = userFlux
        .distinct(User::getId);

23. sort

Flux 데이터를 정렬한다.

Flux<Integer> numbers = Flux.just(3, 1, 2)
        .sort();

결과:

1
2
3

반환 타입:

Flux<T> → Flux<T>

Comparator도 사용할 수 있다.

Flux<User> users = userRepository.findAll()
        .sort(Comparator.comparing(User::getName));

다만 모든 데이터를 모아 정렬해야 하므로 무한 스트림이나 매우 큰 데이터에서는 주의해야 한다.

DB 조회라면 가능하면 Database 계층에서 ORDER BY를 적용하는 것이 효율적이다.


24. take

처음부터 지정한 개수만큼 데이터를 가져온다.

Flux<Product> products = productRepository.findAll()
        .take(10);

반환 타입:

Flux<T> → Flux<T>

처음 10개만 발행한 뒤 완료한다.


25. skip

처음부터 지정한 개수만큼 데이터를 건너뛴다.

Flux<Product> products = productRepository.findAll()
        .skip(10);

반환 타입:

Flux<T> → Flux<T>

처음 10개를 제외한 나머지를 발행한다.

다만 DB 페이징을 구현할 때 모든 데이터를 조회한 후 skip하는 방식은 비효율적일 수 있다.

가능하면 Database 쿼리 단계에서 페이징을 적용해야 한다.


26. delayElements

각 데이터 발행 사이에 지연 시간을 적용한다.

Flux<Long> stream = Flux.interval(Duration.ofSeconds(1));

또는:

Flux<String> stream = Flux.just("A", "B", "C")
        .delayElements(Duration.ofSeconds(1));

SSE 테스트 코드에서 자주 볼 수 있다.

@GetMapping(
        value = "/stream",
        produces = MediaType.TEXT_EVENT_STREAM_VALUE
)
public Flux<String> stream() {
    return Flux.just("A", "B", "C")
            .delayElements(Duration.ofSeconds(1));
}

27. timeout

지정된 시간 동안 데이터가 오지 않으면 Timeout 예외를 발생시킨다.

Mono<User> result = userRepository.findById(userId)
        .timeout(Duration.ofSeconds(3));

fallback Publisher를 지정할 수도 있다.

Mono<User> result = userRepository.findById(userId)
        .timeout(
                Duration.ofSeconds(3),
                Mono.just(User.guest())
        );

외부 API 호출이나 서비스 간 통신에서 무한정 응답을 기다리지 않도록 사용할 수 있다.


28. retry

에러가 발생했을 때 다시 구독하여 작업을 재시도한다.

Mono<User> result = externalApiClient.getUser(userId)
        .retry(3);

최초 요청을 제외하고 최대 3번 재시도한다.

하지만 모든 작업에 무조건 retry를 적용하면 안 된다.

특히 다음과 같은 쓰기 작업은 중복 실행될 수 있다.

paymentRepository.save(payment)
        .retry(3);

결제, 주문, 메시지 전송처럼 중복 실행이 문제가 되는 작업은 멱등성 보장 없이 단순 재시도하면 안 된다.


29. subscribe

subscribe()는 Reactive Stream 실행을 시작한다.

userRepository.findById(userId)
        .subscribe(user ->
                System.out.println(user.getName())
        );

Reactor Publisher는 기본적으로 구독되기 전까지 실행되지 않는다.

Publisher 생성
   ↓
연산자 연결
   ↓
subscribe
   ↓
실제 실행

하지만 Spring WebFlux Controller나 Service 내부에서 직접 subscribe()를 호출하는 것은 일반적으로 피해야 한다.

@GetMapping("/{id}")
public Mono<UserResponse> getUser(@PathVariable Long id) {
    userService.getUser(id).subscribe(); // 직접 구독하지 않음

    return userService.getUser(id);
}

올바른 방식은 Publisher를 그대로 반환하는 것이다.

@GetMapping("/{id}")
public Mono<UserResponse> getUser(@PathVariable Long id) {
    return userService.getUser(id);
}

Spring WebFlux가 최종적으로 Publisher를 구독하고 HTTP 응답을 처리한다.


30. 자주 사용하는 반환 타입 정리

원본 타입연산자람다식 반환 타입최종 반환 타입
Mono<T>mapRMono<R>
Flux<T>mapRFlux<R>
Mono<T>flatMapMono<R>Mono<R>
Flux<T>flatMapPublisher<R>Flux<R>
Mono<T>flatMapManyPublisher<R>Flux<R>
Mono<T>filterbooleanMono<T>
Flux<T>filterbooleanFlux<T>
Mono<T>switchIfEmptyMono<T>Mono<T>
Flux<T>switchIfEmptyPublisher<T>Flux<T>
Flux<T>collectList없음Mono<List<T>>
Flux<T>collectMapKey 또는 Key/ValueMono<Map<K, V>>
Flux<T>next없음Mono<T>
Mono<T>then없음Mono<Void>
Flux<T>then없음Mono<Void>
Mono<T>thenReturnRMono<R>
Mono<T>thenManyFlux<R>Flux<R>
Mono<T>hasElement없음Mono<Boolean>
Flux<T>hasElements없음Mono<Boolean>
Flux<T>count없음Mono<Long>

31. map과 flatMap 판단 방법

가장 쉽게 판단하려면 람다식의 반환 타입을 확인하면 된다.

일반 객체 반환

.map(user -> UserResponse.from(user))

반환값:

UserResponse

따라서:

map

Mono 반환

.flatMap(user -> userRepository.save(user))

반환값:

Mono<User>

따라서:

flatMap

Flux 반환

Mono에서 여러 데이터를 반환하면:

.flatMapMany(user ->
        orderRepository.findByUserId(user.getId())
)

반환값:

Flux<Order>

따라서:

flatMapMany

Flux 내부 각 요소에서 Publisher를 반환하면:

.flatMap(user ->
        orderRepository.findByUserId(user.getId())
)

최종 반환 타입:

Flux<Order>

32. 실전 예제

사용자 ID로 사용자를 조회하고, 활성 회원인지 검사한 다음 해당 회원의 채팅방 목록을 조회한다고 가정한다.

public Flux<ChatRoomResponse> getChatRooms(Long userId) {
    return userRepository.findById(userId)
            .switchIfEmpty(
                    Mono.error(new UserNotFoundException())
            )
            .filter(user ->
                    user.getStatus() == UserStatus.ACTIVE
            )
            .switchIfEmpty(
                    Mono.error(new InactiveUserException())
            )
            .flatMapMany(user ->
                    chatRoomRepository.findByUserId(user.getId())
            )
            .map(ChatRoomResponse::from);
}

각 단계의 반환 타입은 다음과 같다.

userRepository.findById(userId)
→ Mono<User>

switchIfEmpty(...)
→ Mono<User>

filter(...)
→ Mono<User>

switchIfEmpty(...)
→ Mono<User>

flatMapMany(...)
→ Flux<ChatRoom>

map(...)
→ Flux<ChatRoomResponse>

전체 흐름:

UI
 ↓ GET /users/{userId}/chat-rooms
Controller
 ↓
Service
 ↓
UserRepository 조회
 ↓ Mono<User>
활성 상태 검사
 ↓
ChatRoomRepository 조회
 ↓ Flux<ChatRoom>
Response 변환
 ↓ Flux<ChatRoomResponse>
Controller
 ↓
JSON Response
 ↓
UI

33. 정리

Reactor 연산자를 사용할 때는 연산자 이름을 무작정 외우기보다 현재 타입과 변환 후 타입을 먼저 확인해야 한다.

일반 객체 변환
→ map

Mono 또는 Flux를 반환하는 비동기 작업 연결
→ flatMap

Mono에서 Flux로 변환
→ flatMapMany

조건에 맞는 값만 통과
→ filter

데이터가 없을 때 대체 Publisher 실행
→ switchIfEmpty

Flux를 List로 수집
→ collectList

첫 번째 데이터만 Mono로 변환
→ next

결과값을 무시하고 완료 여부만 전달
→ then

가장 중요한 기준은 다음 한 줄이다.

람다식이 일반 객체를 반환하면 map,
람다식이 Publisher를 반환하면 flatMap

그리고 연산자를 연결할 때마다 반환 타입을 직접 적어 보면 Reactor 흐름을 훨씬 쉽게 이해할 수 있다.

Mono<User>
    .flatMap(...)      // Mono<Order>
    .map(...)          // Mono<OrderResponse>
    .switchIfEmpty(...) // Mono<OrderResponse>

Reactive Programming에서는 현재 데이터가 무엇인지뿐만 아니라, 현재 데이터가 어떤 Publisher에 감싸져 있는지를 함께 확인하는 습관이 중요하다.

profile
Develop

0개의 댓글