Skip to content

Reactive Programming ‐ Mastering Reactor Operators in WebFlux & A Deep Dive into WebFlux Operators

woojin edited this page Jul 22, 2026 · 2 revisions

Reactor에서 제공하는 데이터 변환 연산자 메서드 - map()

  • map()은 각 원소를 1:1로 동기 변환하는 연산자이다.
  • onNext()로 들어온 값 하나를 함수에 통과시켜 다른 값 하나로 바꿔서 그대로 아래로 흘려보낸다.
// 단건 조회: Mono<User> → Mono<UserResponse>
public Mono<UserResponse> getUser(Long id) {
    return userRepository.findById(id)
            .map(UserResponse::from);   // User 하나 → UserResponse 하나 (1:1 동기 변환)
}

Reactor에서 제공하는 데이터 변환 연산자 메서드 - flatMap()

  • flatMap()은 각 원소를 또 다른 비동기 스트림(Mono/Flux)으로 바꾼 뒤, 그것들을 구독해서 하나의 스트림으로 펼치는(flatten) 연산자이다.
  • 리액티브에서 DB·HTTP 호출 결과는 값이 아니라 Mono/Flux이다. 그래서 "각 원소마다 또 다른 비동기 호출"을 하면 필연적으로 스트림이 중첩되고, 이걸 풀어주는 게 flatMap() 메서드의 역할이다.
// flatMap: 함수가 'Publisher'를 반환 → 자동으로 펼쳐서 Flux<Order>
users.flatMap(user -> orderRepository.findByUserId(user.getId())) // ✅ Flux<Order>

Reactor에서 제공하는 데이터 변환 연산자 메서드 - concatMap()

  • concatMap()flatMap()과 똑같이 각 원소를 Mono/Flux로 바꿔 펼치는데, 안쪽 스트림을 하나씩 순서대로(순차적으로) 구독한다.
  • concatMap()은 무조건 들어온 순서대로 하나씩 차례대로 처리한다.
  • 만약 A, B, C가 순서대로 들어오면, A의 비동기 처리가 끝날 때까지 B, C는 시작도 하지 않고 기다린다.
  • 비동기 처리 속도가 제각각이더라도 최종 결과물은 반드시 A ➔ B ➔ C라는 원래의 순서가 엄격하게 보장된다.
users.concatMap(user ->
        orderRepository.findByUserId(user.getId())   // 반환은 flatMap과 동일 (Publisher)
);

Reactor에서 제공하는 데이터 변환 연산자 메서드 - flatMapSequential()

  • flatMapSequential은 비동기로 동시에 쏘면서 성능을 챙기면서도 최종 결과물의 순서는 원본 순서대로 정렬한다.
  • 순서를 맞추려고 버퍼를 쓰기 때문에, 앞 순서 원소가 느리면 뒤 결과들이 계속 쌓이게 된다.
users.flatMapSequential(user ->
        orderRepository.findByUserId(user.getId())   // 반환은 셋 다 동일 (Publisher)
);

Reactor에서 제공하는 데이터 변환 연산자 메서드 - flatMapMany()

  • flatMapMany()는 Mono에서 쓰는 연산자로, 하나의 값을 Flux(여러 개)로 펼쳐주는 변환이다. 즉 Mono<T>Flux<R>, 카디널리티가 1 → N으로 바뀐다.
Mono<User> user = userRepository.findById(id);   // 유저 1명

Flux<Order> orders = user.flatMapMany(u ->
        orderRepository.findByUserId(u.getId())   // 그 유저의 주문 N건 → Flux
);                                                // Mono<User> → Flux<Order>

데이터 스트림 환경에서 가공 및 조율을 위한 Reactor의 필터링 연산자 메서드 - filter()

  • filter()는 조건(Predicate)을 만족하는 원소만 통과시키고, 나머지는 버리는 연산자이다.
  • 값을 바꾸지 않고(변환 아님), 개수만 줄인다.
  • filter()의 Predicate도 map()처럼 동기로 즉시 실행된다. 그래서 그 안에서 블로킹 호출을 하면 이벤트 루프를 막게 된다.
users
    .filter(User::isActive)        // 활성 유저만 골라서 (개수 ↓)
    .map(UserResponse::from);      // 그걸 DTO로 변환 (값 바꿈)

데이터 스트림 환경에서 가공 및 조율을 위한 Reactor의 필터링 연산자 메서드 - distinct()

  • distinct()는 이미 지나간 값과 중복되는 원소를 걸러내고, 처음 보는 값만 통과시키는 연산자이다. 스트림 버전의 "중복 제거"이다.
  • 그러나 대용량 아키텍처 관점에서 치명적인 트레이드 오프를 가지고 있다.
    • 메모리 누수 위험 : Flux.interval처럼 종료되지 않고 무한히 흐르는 스트림이거나 하루에 수억 건씩 쏟아지는 금융 트랜잭션 스트림에 걸게 되면 스트림이 유지되는 내내 내부 Set에 데이터가 계속 누적되면서 힙 메모리를 잡아먹게 되고 결국 서버가 OutOfMemoryError(OOM)을 뱉을 수 있다.
    • distinctUntilChanged() : 해당 메서드는 전체 데이터가 아니라 바로 직전에 통과한 데이터와만 비교해 연속으로 중복되는 데이터만 쳐내는 오퍼레이터이다.
Flux.just(1, 2, 2, 3, 1, 3, 4)
    .distinct()
    .subscribe(System.out::println);   // 1, 2, 3, 4 (중복은 버려짐)

데이터 스트림 환경에서 가공 및 조율을 위한 Reactor의 필터링 연산자 메서드 - elementAt()

  • elementAt()은 Flux에서 특정 인덱스(N번째) 원소 하나만 꺼내는 연산자이다. 결과가 1개이므로 Flux<T>Mono<T>로 바뀐다.
  • 그러나 기본값 없는 인덱스 초과가 발생할 위험이 높아 치명적인 안티패턴으로 잘 쓰이지 않는다.
Flux.just("a", "b", "c", "d")
    .elementAt(2)                      // 인덱스 2 = 세 번째
    .subscribe(System.out::println);   // "c"

다중 스트림 구조에서 데이터 결합을 위한 Reactor 결합 연산자 메서드 - concat()

  • concat()은 여러 개의 Publisher(스트림)를 순서대로 이어붙여 하나로 만드는 연산자이다.
  • 앞 스트림이 완전히 끝나야 다음 스트림을 구독한다.
Flux<Integer> flux1 = Flux.just(1, 2, 3);
Flux<Integer> flux2 = Flux.just(4, 5, 6);

Flux.concat(flux1, flux2)
    .subscribe(System.out::println);   // 1, 2, 3, 4, 5, 6

다중 스트림 구조에서 데이터 결합을 위한 Reactor 결합 연산자 메서드 - concatWith()

  • concatWith()concat()의 인스턴스 메서드 버전이다.
  • 기존 스트림 뒤에 다른 스트림을 이어붙인다.
// concat: 스트림들을 '나란히 나열'하는 느낌 — 여러 개를 한자리에 모을 때
Flux.concat(header, body, footer);

// concatWith: 기존 체이닝에 '이어서 덧붙이는' 느낌 — 파이프라인 끝에 자연스럽게
userRepository.findVip()
    .map(UserResponse::from)
    .concatWith(userRepository.findNormal().map(UserResponse::from));    //  ↑ 앞 파이프라인에 매끄럽게 이어짐

다중 스트림 구조에서 데이터 결합을 위한 Reactor 결합 연산자 메서드 - merge()

  • merge()는 여러 스트림을 동시에 구독해서, 각 스트림에서 원소가 도착하는 대로 하나로 합치는 결합 연산자이다.
Flux<Integer> flux1 = Flux.just(1, 2, 3);
Flux<Integer> flux2 = Flux.just(4, 5, 6);

Flux.merge(flux1, flux2)
    .subscribe(System.out::println);   // 순서 보장 안 됨 (도착 순)

다중 스트림 구조에서 데이터 결합을 위한 Reactor 결합 연산자 메서드 - zip()

  • zip()은 여러 스트림에서 원소를 하나씩 짝지어(같은 순번끼리) 결합하는 연산자이다.
Flux<String> names = Flux.just("김", "이", "박");
Flux<Integer> ages  = Flux.just(20, 30, 40);

Flux.zip(names, ages, (name, age) -> name + ":" + age)
    .subscribe(System.out::println);   // 김:20, 이:30, 박:40

무한한 데이터 스트림 환경에서 조건을 검증하기 위한 연산자 메서드 - defaultIfEmpty(), switchIfEmpty()

  • defaultIfEmpty()는 스트림이 비어 있으면 미리 정해둔 값 하나를 대신 방출한다.
  • switchIfEmpty()는 스트림이 비어 있으면 다른 스트림으로 전환한다.
userRepository.findByEmail(email)     // Mono<User> — 없으면 빈 Mono
    .defaultIfEmpty(User.guest());    // 비었으면 게스트 유저로 대체
cacheRepository.findById(id)              // 1차: 캐시 조회 (없으면 빈 Mono)
    .switchIfEmpty(
        dbRepository.findById(id)         // 2차: 캐시에 없으면 DB에서 조회
    );

무한한 데이터 스트림 환경에서 조건을 검증하기 위한 연산자 메서드 - hasElement(), hasElements()

  • hasElement() / hasElements() 둘 다 스트림에 원소가 있는지 없는지를 boolean으로 알려주는 연산자이다. 값 자체가 아니라 존재 여부만 뽑아낸다.
userRepository.findById(id)      // Mono<User>
    .hasElement();               // Mono<Boolean> — 있으면 true, 없으면 false
orderRepository.findByUserId(id)   // Flux<Order>
    .hasElements();                // Mono<Boolean> — 하나라도 있으면 true

무한한 데이터 스트림 환경에서 조건을 검증하기 위한 연산자 메서드 - any(), all()

  • any() / all() 둘 다 스트림 전체를 조건(Predicate)으로 판정해서 Mono<Boolean> 하나로 만드는 연산자이다.
  • any() : 하나라도 만족하면 true
  • all() : 전부 만족해야 true
orderRepository.findByUserId(id)   // Flux<Order>
    .any(Order::isPaid);           // 결제된 주문이 하나라도 있나? → Mono<Boolean>
orderRepository.findByUserId(id)   // Flux<Order>
    .all(Order::isPaid);           // 모든 주문이 결제됐나? → Mono<Boolean>

무한한 데이터 스트림 환경에서 조건을 검증하기 위한 연산자 메서드 - then(), thenMany()

  • 둘 다 앞 작업의 결과값은 버리고, 앞이 끝난 뒤 다음 작업으로 넘어가는 연산자이다.
  • then() — 앞이 끝나면 → Mono 하나
  • thenMany() — 앞이 끝나면 → Flux
// 1) then() — 앞 값 버리고 '완료 신호'만 (Mono<Void>)
saveUser(user)
    .then();                        // User 저장 끝나면 Mono<Void>로 완료만 알림

// 2) then(Mono) — 앞 끝나면 다음 Mono 실행
saveUser(user)                      // Mono<User> (결과값은 버림)
    .then(sendWelcomeMail(user));   // 저장 끝난 뒤 메일 발송
deleteAllOldOrders(userId)          // Mono<Void> — 정리 작업 (값 없음)
    .thenMany(
        loadFreshOrders(userId)     // Flux<Order> — 끝난 뒤 새로 조회
    );   // → Flux<Order>

다중 스트림 환경에서 데이터를 모아서 하나의 산출물을 만드는 집계 연산자 메서드 - reduce(), scan()

  • reduce() / scan() 둘 다 원소들을 하나의 누적값으로 접어나가는 연산자이다. 앞 원소들을 계속 합쳐가는 누적 연산이다.
  • reduce() : 모든 원소를 누적한 마지막 값 하나만 방출한다.
  • scan() : 누적 과정을 매 단계마다 방출한다.
Flux.just(1, 2, 3, 4)
    .reduce((acc, next) -> acc + next)   // 누적 합
    .subscribe(System.out::println);     // 10 (최종값 하나만)

다중 스트림 환경에서 데이터를 모아서 하나의 산출물을 만드는 집계 연산자 메서드 - groupBy()

  • groupBy()는 하나의 Flux를 키(key) 기준으로 여러 개의 하위 스트림으로 쪼개는 연산자이다.
  • SQL의 GROUP BY, Java Stream의 Collectors.groupingBy를 리액티브 스트림에서 하는 것이다.
Flux.just(1, 2, 3, 4, 5, 6)
    .groupBy(n -> n % 2 == 0 ? "짝수" : "홀수")   // 키로 분류
    .flatMap(group ->
        group.collectList()
             .map(list -> group.key() + ": " + list)   // group.key()로 키 접근
    )
    .subscribe(System.out::println);
    // 홀수: [1, 3, 5]
    // 짝수: [2, 4, 6]

📖 Code Philosophy

📖 Java

📖 Kotlin

📖 Coroutine

📖 Spring

📖 Spring Security

📖 Spring Batch

📖 Database

📖 MySQL

📖 Redis

📖 JPA

📖 QueryDsl

📖 MSA

📖 Kafka

📖 Apache Flink

  • [Apache Flink - Apache Flink Architecture]
  • [Apache Flink - Stream Processing]
  • [Apache Flink - Data Stream API & Window]
  • [Apache Flink - State Management]

📖 HTTP

📖 AWS

📖 Docker

📖 Kubernetes

📖 Github Actions

📖 Jenkins

📖 Nginx

📖 Monitoring

📖 Test(feat. Load Testing)

  • [Test - Load Testing Fundamentals]
  • [Test - Identifying Bottlenecks with Load Testing]
  • [Test - Resolving Bottlenecks and Improving Performance]

📖 Test(feat. Java)

📖 Spring AI

📖 gRPC

  • [gRPC - Writing .proto Files with Protocol Buffers]
  • [gRPC - Various Communication Patterns in gRPC]
  • [gRPC - gRPC Optimization Techniques and Advanced Features]

📖 TDD(Test-Driven-Development)

📖 PostgreSQL

  • [PostgreSQL - Docker만을 사용하는 경량화된 환경 구성 방법]
  • [PostgreSQL - PostgreSQL에서 제공하는 데이터 타입]
  • [PostgreSQL - PostgreSQI의 JSONB, 역인덱싱과 활용 방법]
  • [PostgreSQL - 데이터베이스 성능을 위한 최적화 패턴 및 전략]
  • [PostgreSQL - 트랜잭션과 ACID, Isolation 수준별 차이]
  • [PostgreSQL - Database Lock 교착상태와 읽기/쓰기 성능을 보장하는 MVCC 모델]
  • [PostgreSQL - pgvector와 벡터 저장, 유사도 검색 패턴 개념]
  • [PostgreSQL - 벡터 인덱스 최적화와 벡터 검색과 전문 검색 결합 패턴]
  • [PostgreSQL - PostgreSQL 플러그인]
  • [PostgreSQL - PostGIS - 공간 쿼리와 GIST 인덱스, 지리 타입과 공간 쿼리를 위한 타입과 기본 함수]
  • [PostgreSQL - pg_search - 검색 엔진 없이 텍스트 검색 구현과 주의사항]
  • [PostgreSQL - 단일 인스턴스 한계를 극복하는 분산 패턴과 스케줄링, 분산 환경 구축 방법]
  • [PostgreSQL - Citus - 분산 테이블과 분산 쿼리를 위한 Extension과 데이터 분산 처리]
  • [PostgreSQL - pg_cron - PostgreSQL로 구성하는 CronJob]
  • [PostgreSQL - 스케줄러 + 분산 처리를 동시에 도입하는 주기적 집계 쿼리 패턴]

📖 Workflow-Driven Techniques for Large-Scale Traffic Processing

  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Kafka + Debezium을 활용한 CDC 패턴 설계]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Temporal을 활용한 워크플로우 패턴]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Docker와 경량 이미지를 활용한 환경 구축 방법]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Kafka에서의 메시지 Delivery Guarantee]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - 실시간 동기화의 핵심 CDC]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - MySQL Binary Log 기반의 CDC]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Binary Log 기반의 CDC 구현 플랫폼 Debezium이란?]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Debezium Architecture]
  • [Workflow-Driven Techniques for Large-Scale Traffic Processing - Debezium Architecture Best Practice와 주의사항]

📖 Reactive Programming

📖 ElasticSearch

  • [ElasticSearch - Search Fundamentals & How Search Works]
  • [ElasticSearch - Korean-Optimized Search]
  • [ElasticSearch - Mapping & Data Types]
  • [ElasticSearch - Common Search Features]
  • [ElasticSearch - Building a Product Search Engine with Elasticsearch]
  • [ElasticSearch - Deploying Elasticsearch with Elastic Cloud]
  • [ElasticSearch - Managing Documents]
  • [ElasticSearch - Analysis & Mapping]
  • [ElasticSearch - Search Fundamentals]
  • [ElasticSearch - Query Joins]
  • [ElasticSearch - Processing Search Results]
  • [ElasticSearch - Aggregations]
  • [ElasticSearch - Tips for Improving Search Results]
  • [ElasticSearch - Elasticsearch Clients]
  • [ElasticSearch - Understanding How Elasticsearch Works]
  • [ElasticSearch - Monitoring Elasticsearch]
  • [ElasticSearch - Elasticsearch Troubleshooting]

Clone this wiki locally