Project Reactor의 Mono와 Flux는 다양한 연산자(operator)를 제공하여 데이터 스트림을 변환, 조합, 필터링 할 수 있도록 한다. 이 연산자들을 통해 복잡한 비동기 로직을 선언적이고
간결하게 표현할 수 있다.
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()연산자는 스트림의 각 항목에 대해 동기적인 변환 함수를 적용하여 새로운 타입 또는 값으로 매핑한다.
입력 스트림의 각 요소가 변환되어 다음 스트림으로 전달된다.
@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()연산자는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()연산자는 주어진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()연산자는 여러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()연산자는 여러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초 대기
}
}
flatMap()과 유사하지만 내부 스트림의 결과를 순서를 보장하여 병합한다.Mono또는 Flux로 대체하여 스트림을 계속 진행한다.Flux이고 단수형인 것은 Mono의 것이다.then(): MonothenMany(): Mono, FluxMono또는 Flux를 실행한다.