목차
리액티브 오퍼레이션 적용
Flux와 Mono는 리액터가 제공하는 가장 핵심적인 구성요소이며, 이 오퍼레이션들은 두 타입을 함께 결합하여 데이터가 전달될 수 있는 파이프라인을 생성한다.
Flux와 Mono의 연산자는 사용 목적에 따라 다음처럼 나눠 볼 수 있다. 정확한 API 목록은 사용하는 Reactor 버전의 문서를 확인한다.
-
생성(Creation) 오퍼레이션
-
조합(Combination) 오퍼레이션
-
변환(Transformation) 오퍼레이션
-
로직(Logic) 오퍼레이션
마블 다이어그램은 값, 완료, 오류 신호를 시간축에서 읽는 데 도움이 된다. 아래 예제에서는 각 연산자의 순서 보장과 구독 시점을 코드로 확인한다.
1. 생성(Creation) 오퍼레이션
데이터를 생성하여 방출할 때 사용.
객체로부터 생성
Flux나 Mono로 하나 이상의 객체를 생성하려면 just() 메서드를 사용한다..
//flux
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape");
//mono
Mono<String> orange = Mono.just("Orange");
위 코드처럼 Flux, Mono에서 just() 이용하여 데이터를 생성했지만 Subscribe가 없는데 이 상태는 호스를 수도꼭지에 끼운 것에 비유할 수 있다. 수도꼭지에 끼운 호스에 물을 흐르게 하려면 Subscribe(구독자)를 이용하여 데이터를 흘러나가게 한다.
Mono.subscribe(), or Flux.subscribe()
//flux
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape");
fruit.subscribe(f -> System.out.println("Fruit : " + f));
//mono
Mono<String> orange = Mono.just("Orange");
orange.subscribe(f -> System.out.println("Fruit : "+ f));
리액터에서 StepVerifier를 사용하면 Mono, Flux를 테스트할 수 있다.
StepVerifier가 fruit를 구독한 후 이름과 일치한지 검사한다.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape");
StepVerifier.create(fruit)
.expectNext("Apple")
.expectNext("Orange")
.expectNext("Grape")
.verifyComplete();
컬렉션으로부터 생성하기
Flux는 Array fromArray(), Iterable fromIterable(), Java Stream fromStream()을 생성할 수 있다.
String[] fruits = new String[]{"Apple", "Orange", "Grape"};
List<String> frustList = new ArrayList<>(Arrays.asList(fruits));
//array fromArray()
Flux<String> fruitArray = Flux.fromArray(fruits);
StepVerifier.create(fruitArray)
.expectNext("Apple")
.expectNext("Orange")
.expectNext("Grape")
.verifyComplete();
//iterable fromIterable()
Flux<String> fruitList = Flux.fromIterable(frustList);
StepVerifier.create(fruitList)
.expectNext("Apple")
.expectNext("Orange")
.expectNext("Grape")
.verifyComplete();
//streams fromStream()
Flux<String> fruitStreams = Flux.fromStream(Arrays.stream(fruits));
StepVerifier.create(fruitStreams)
.expectNext("Apple")
.expectNext("Orange")
.expectNext("Grape")
.verifyComplete();
Flux 데이터 생성
데이터없이 매번 새 값을 증가하는 숫자를 보내는 카운터 역할의 Flux만 필요할 때 range()를 사용할 수 있다.
1부터 10까지 증가
Flux<Integer> range = Flux.range(1, 10);
StepVerifier.create(range)
.expectNext(1)
.expectNext(2)
.expectNext(3)
.expectNext(4)
.expectNext(5)
.expectNext(6)
.expectNext(7)
.expectNext(8)
.expectNext(9)
.expectNext(10)
.verifyComplete();
시작 값과 종료 값 대신 값이 방출되는 시간 간격이나 주기를 지정해주는 interval()
Flux<Long> interval = Flux.interval(Duration.ofSeconds(1)).take(5);
StepVerifier.create(interval)
.expectNext(0L)
.expectNext(1L)
.expectNext(2L)
.expectNext(3L)
.expectNext(4L)
.verifyComplete();
조합(Combination) 오퍼레이션
두 개의 리액티브 타입을 결합해야 하거나 하나의 Flux를 두 개 이상의 리액티브 타입으로 분할해야하는 경우 사용
리액티브 타입 결합
mergeWith() : 두 개의 Flux 스트림을 하나의 Flux로 결과를 보여줄 때
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape");
Flux<String> melon = Flux.just("WaterMelon", "Melon", "Kiwi");
Flux<String> merged = fruit.mergeWith(melon);
StepVerifier.create(merged)
.expectNextCount(6) // 두 소스의 상대적 도착 순서에 의존하지 않는다.
.verifyComplete();
Flux<String> sequential = fruit.concatWith(melon);
StepVerifier.create(sequential)
.expectNext("Apple", "Orange", "Grape",
"WaterMelon", "Melon", "Kiwi")
.verifyComplete();
mergeWith()는 두 소스의 결과를 섞는다. 소스마다 내부 순서는 유지될 수 있지만 소스 간 도착 순서를 고정해 테스트하면 스케줄링에 따라 흔들릴 수 있다. 소스 전체의 순서가 중요하다면 concatWith()를 쓴다. 순서를 보장받으려고 나노초 지연을 넣는 방법은 안정적인 테스트가 아니다.
zip() 오퍼레이션을 사용할 수 있다.
zip()은 각 소스의 첫 번째 값끼리, 두 번째 값끼리 짝지어 새 값을 만든다. 더 짧은 소스가 끝나면 결과도 끝난다.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape");
Flux<Integer> count = Flux.just(1,2,3);
Flux<Tuple2<String, Integer>> zipFlux = Flux.zip(fruit, count);
zipFlux.subscribe(s-> System.out.println(s.getT1() + " : " + s.getT2()));
//or
Flux<String> zipFlux2 = Flux.zip(fruit, count, (f,c) -> f + " : "+ c);
zipFlux2.subscribe(System.out::println);
3. 변환(Transformation) 오퍼레이션
데이터가 흐르는동안 일부 값을 필터링하거나 변경할 경우 사용
리액티브 타입으로부터 데이터 필터링
skip() : 맨 앞에서부터 원하는 개수의 항목을 무시하는 것.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape").skip(2);
StepVerifier.create(fruit)
.expectNext("Grape")
.verifyComplete();
2초동안 기다렸다가 값을 방출
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape").delayElements(Duration.ofSeconds(2)).skip(2);
StepVerifier.create(fruit)
.expectNext("Grape")
.verifyComplete();
skip()은 처음부터 여러개의 항목을 건너뛰는 반면, take()는 지정된 항목만 방출한다.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape").take(1);
fruit.subscribe(System.out::println);
StepVerifier.create(fruit)
.expectNext("Apple")
.verifyComplete();
일정 시간이 경과될 동안 방출
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape").delayElements(Duration.ofSeconds(2)).take(1);
fruit.subscribe(System.out::println);
StepVerifier.create(fruit)
.expectNext("Apple")
.verifyComplete();
filter() : Flux를 필터링할 때 사용
Grape가 아닌것만 출력
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape")
.filter(s->!s.equals("Grape"));
fruit.subscribe(System.out::println);
StepVerifier.create(fruit)
.expectNext("Apple")
.expectNext("Orange")
.verifyComplete();
distinct()를 이용하면 중복값을 제거할 수 있다.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape", "Apple")
.distinct();
fruit.subscribe(System.out::println);
StepVerifier.create(fruit)
.expectNext("Apple")
.expectNext("Orange")
.expectNext("Grape")
.verifyComplete();
리액티브 데이터 매핑
발행된 항목을 다른 형태나 타입으로 매핑하는 방법으로 대표적으로 map()과 , flatMap()이 있다.
map() : 반환을 수행하는 Flux를 생성하며, 동기적으로 매핑이 수행된다.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape")
.map(String::toUpperCase);
fruit.subscribe(System.out::println);
StepVerifier.create(fruit)
.expectNext("APPLE")
.expectNext("ORANGE")
.expectNext("GRAPE")
.verifyComplete();
flatMap() : 각 값을 Publisher로 바꾸고 그 결과들을 병합한다. 내부 Publisher가 비동기적으로 동작하면 결과 순서가 바뀔 수 있다.
map()은 한 값을 즉시 다른 값으로 바꾸고, flatMap()은 각 값을 Mono나 Flux로 바꾸어 합친다. 각 내부 Publisher의 실행 스케줄러와 동시성 한도를 별도로 정할 수 있다.
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape")
.flatMap(n -> Mono.fromCallable(n::toUpperCase)
.subscribeOn(Schedulers.parallel()), 2)
.sort(); // 테스트에서 결과 순서를 고정하기 위한 예시
StepVerifier.create(fruit)
.expectNext("APPLE", "GRAPE", "ORANGE")
.verifyComplete();
subscribe()는 소비를 시작하고, subscribeOn()은 해당 Publisher의 구독·생성 작업이 실행될 스케줄러를 지정한다. subscribeOn()을 한 번 붙인 것만으로 모든 요소가 병렬 처리되지는 않는다.
Schedulers는 다음과 같은 메서드를 가지고 있다.
| Scheduler | 주로 쓰는 경우 | 주의점 |
|---|---|---|
immediate() | 호출 스레드에서 실행 | 블로킹 호출은 그대로 호출 스레드를 막음 |
single() | 재사용되는 단일 스레드 | 느린 작업 하나가 뒤 작업을 지연 |
newSingle("name") | 별도 단일 스레드 인스턴스 | 사용 후 폐기 책임 확인 |
parallel() | 짧은 CPU 계산 | JDBC 같은 블로킹 I/O에 쓰지 않음 |
boundedElastic() | 기존 블로킹 I/O를 격리 | 스레드·큐가 무한하지 않음 |
리액티브 스트림의 데이터 버퍼링 하기
buffer() : 데이터를 처리하는 동안 데이터 스트림을 작은 덩어리로 분할
Flux<String> fruit = Flux.just("Apple", "Orange", "Grape","Strawberry", "Banana");
Flux<List<String>> bufferFlux = fruit.buffer(3);
bufferFlux.subscribe(System.out::println);
StepVerifier.create(bufferFlux)
.expectNext(Arrays.asList("Apple", "Orange", "Grape"))
.expectNext(Arrays.asList("Strawberry", "Banana"))
.verifyComplete();
buffer()를 flatMap()과 같이 사용하여 병행으로 처리
Flux.just("Apple", "Orange", "Grape","Strawberry", "Banana")
.buffer(3)
.flatMap(s -> Flux.fromIterable(s)
.map(String::toUpperCase)
.subscribeOn(Schedulers.parallel())
.log())
.subscribe();
5개의 값을 새로운 Flux로 버퍼링 하여 flatMap()에 적용한다. 각 List의 버퍼를 가져와서 해당요소로부터 새로운 Flux를 생성하고 map() 을 사용한다. 버퍼링 된 List는 별도의 스레드에서 병행으로 계속 처리될 수 있다.
log() : 모든 리액티브 시트림 이벤트를 로깅하여 실제 어떻게 처리되는지 파악할 수 있다.
log() : 모든 리액티브 스트림 이벤트를 로깅하여 실제 어떻게 처리되는지 파악할 수 있다.
collectList() : Flux의 모든 항목을 모아 Mono<List<T>>를 반환한다. 무한 스트림에는 끝나지 않고 큰 스트림에는 메모리를 많이 쓴다.
Mono<List<String>> list = Flux.just("Apple", "Orange", "Grape","Strawberry", "Banana").collectList();
StepVerifier.create(list)
.expectNext(Arrays.asList("Apple", "Orange", "Grape","Strawberry", "Banana"))
.verifyComplete();
4. 로직(Logic) 오퍼레이션
Mono나 Flux가 발행한 항목이 어떤 조건과 일치하는지 알아야 할 경우 사용.
-
all(): 모든 메시지가 조건을 충족하는지 확인Flux<String> fruitFlux = Flux.just("apple", "orange", "grape","strawberry", "banana"); Mono<Boolean> hasFruit = fruitFlux.all(a-> a.contains("a")); StepVerifier.create(hasFruit) .expectNext(true) .verifyComplete(); -
any(): 최소 하나의 메시지가 조건을 충족하는지 확인
Flux<String> fruitFlux = Flux.just("apple", "orange", "grape","strawberry", "banana");
Mono<Boolean> hasFruit = fruitFlux.any(a-> a.contains("orange"));
StepVerifier.create(hasFruit)
.expectNext(true)
.verifyComplete();
참조
https://projectreactor.io/docs/core/release/reference/
주문 조회 세 건에 연산자를 적용할 때
예제의 과일 이름을 실제 업무 흐름으로 옮겨 보자. 화면에 주문 ID 세 개가 있고 각 주문의 상세 정보를 비동기 API로 조회한다고 가정한다. map에는 이미 얻은 Order를 화면용 OrderView로 바꾸는 순수 변환이 어울린다. flatMap에는 각 ID에서 Mono<Order>를 만드는 비동기 조회가 어울린다.
Flux<OrderView> views = Flux.fromIterable(orderIds)
.flatMap(id -> orderClient.fetch(id), 4)
.map(order -> new OrderView(order.id(), order.status()));
orderClient.fetch는 비동기 Mono<Order>를 반환한다고 가정한다. 4는 동시에 열어 둘 내부 조회의 상한이다. 무제한으로 호출하면 외부 API 연결 수와 메모리 사용이 급증할 수 있다. 결과가 입력 ID 순서와 달라도 되는 화면이면 flatMap이 적합하다. 순서가 필요하면 flatMapSequential이나 concatMap을 비교한다. concatMap은 내부 조회를 순차적으로 진행하므로 처리량과 순서의 교환이 생긴다.
flowchart LR A["주문 ID Flux"] --> B["flatMap: ID마다 Mono 조회"] B --> C["최대 동시성 4"] C --> D["map: Order → OrderView"] D --> E["구독자 / HTTP 응답"]
flatMap은 순서를 바꿀 수 있고, 어떤 내부 조회가 오류를 내면 전체 스트림이 종료될 수 있다. 한 주문만 실패해도 전체 목록을 실패로 볼지, 실패한 항목을 표시하며 나머지는 내보낼지 업무 규칙을 정해야 한다. onErrorResume를 전체 체인 뒤에 두면 실패 뒤 기본 결과 하나로 대체하는 의미가 될 수 있고, 각 내부 조회 안에 두면 해당 ID만 실패 상태로 바꿀 수 있다. 어디에 오류 처리 연산자를 두는지에 따라 데이터 손실 범위가 달라진다.
시간·메모리·구독 횟수 검증
Flux.interval(Duration.ofSeconds(1)).take(5)를 그대로 검사하면 테스트가 대략 5초 이상 기다린다. 시간 동작만 검증할 때는 StepVerifier.withVirtualTime으로 가상 시간을 사용한다. 다만 생성할 Publisher가 가상 스케줄러를 설치한 뒤 만들어져야 하므로 람다 안에서 생성한다.
StepVerifier.withVirtualTime(
() -> Flux.interval(Duration.ofSeconds(1)).take(3))
.thenAwait(Duration.ofSeconds(3))
.expectNext(0L, 1L, 2L)
.verifyComplete();
collectList는 모든 값을 모은 뒤 한 번에 내보낸다. 무한 스트림에서는 완료되지 않고, 큰 스트림에서는 메모리 사용이 커진다. buffer(3)도 세 값씩 묶는 편의 연산이지 소비자가 느릴 때 전체 입력량을 자동으로 제한하는 안전장치가 아니다. 생산 속도와 소비 속도가 다르면 요청량, 버퍼 상한, 동시성, 드롭 또는 오류 정책을 별도로 정한다.
같은 Flux를 두 번 구독하면 차가운(cold) 소스에서는 외부 API 호출이나 파일 읽기가 두 번 일어날 수 있다. 계산만 있는 Flux.just 예제로는 드러나지 않는다. DB 쓰기나 결제 요청을 Mono에 넣는다면 구독 횟수와 재시도 정책을 명시적으로 검증한다. 값·순서·완료만 확인하는 테스트에 더해 오류, 취소, 느린 구독자 사례가 필요한 이유다.
