목차
리액터 시작
리액티브 프로그래밍은 일련의 작업 단계를 기술하는 것이 아니라 데이터가 전달될 파이프라인을 구성하여 데이터가 전달되는 동안 어떤 형태로든 변경 또는 사용되는 것.
사람의 이름을 가져와 대문자로 변경 후 출력
<명령형 코드>
String name = "devPaik";
String capitalName = name.toUpperCase();
String greeting = "Hello, "+ capitalName + "!";
System.out.println(greeting);
<리액티브 코드>
Mono.just("devPaik")
.map(n -> n.toUpperCase())
.map(cn -> "Hello," + cn + "!")
.subscribe(System.out::println);
위와 같이 리액티브 코드는 데이터가 파이프라인으로 구성하는 것을 볼 수 있다.
파이프라인의 각각의 단계에서 어떻게 하던 데이터는 변경되고, 각 오퍼레이션은 같은 스레드로 실행되거나 다른 스레드로 실행될 수 있다.
리액터에는 Mono, Flux가 있는데 두 개 모두 리액티브 스트림의 Publisher 인터페이스로 구현한 것이다.
-
Mono : 하나의 데이터 항목만을 갖는 데이터셋에 최적회된 타입
-
Flux : 0, 1 or 다수(무한)의 데이터를 갖는 파이프라인
프로젝트에 리액터 추가
리액터 의존성 및 테스트 추가
아래 버전 생략 예시는 Spring Boot의 의존성 관리 또는 Reactor BOM을 사용하는 프로젝트를 전제로 한다. 순수 Java 프로젝트라면 공식 Reactor BOM에서 호환 버전을 지정한다. 옛 compile·testCompile과 Reactor 3.2 버전을 그대로 복사하지 않는다.
gradle
implementation('io.projectreactor:reactor-core')
testImplementation('io.projectreactor:reactor-test')
maven
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-core</artifactId>
</dependency>
<dependency>
<groupId>io.projectreactor</groupId>
<artifactId>reactor-test</artifactId>
<scope>test</scope>
</dependency>
Reactive Type [Flux, Mono]
Reactive Streams의 네 인터페이스는 Publisher<T>, Subscriber<T>, Subscription, Processor<T,R>이다. Reactor는 Publisher를 구현한 Flux<T>와 Mono<T>를 제공한다.
Flux
- Flux는 Publish의 구현체로 0, 1, 또는 여러 요소(0-N개)를 발행하는 일반적인 리액티브 스트림을 정의할 수 있다. (RxJava에서는 Flowable/Observable을 말할 수 있다.)
표현식
onNext x 0..N [ onError | onComplete ]

아래 코드는 1 ~ 5까지 배열로 출력하는 코드이다.
List<Integer> loadStream = Flux.range(1, 5)
// .repeat()
.collectList()
.block();
System.out.println(loadStream);
-
range(1, 5) : 1~5 정수 시퀀스를 생성
-
repeate() : 스트림이 끝나고 다시 스트림을 재구독하는 역할을 한다.
-
collectList() : 생성된 모든 요소를 단일 리스트로 만든다.
-
block() : 실제 구독을 기동하고 최종 결과가 도착할 때까지 실행중인 스레드를 차단한다.
결과
[1, 2, 3, 4, 5]
Mono
Mono는 최대 하나의 요소(0-1개)를 생성할 수 있는 데이터를 스트림을 정의한다.
표현식
onNext x 0..1 [onError | onComplete]

Mono는 결과가 최대 한 개라는 계약을 타입에 드러낸다. 언제나 Flux보다 빠르다고 단정할 수는 없다. 또한 클라이언트에게 작업이 완료됐음을 알리는데 사용할 수 있다.
Mono와 Flux는 서로 변환이 가능하다. 예를들어 Flux<T>.collectList()는 Mono<List<T>>를 반환하고 Mono<T>.flux는 Flux<T>를 반환한다.
from 을이용하여 Flux을 Mono로 변환할 수 있다.
Flux loadStream = Flux.range(1, 5);
Mono.from(loadStream).subscribe(data -> System.out.println(data));
[결과]
1
RxJava 2.x의 리액티브 타입
-
Observable : RxJava 1.x 와 거의 비슷하지만 null값을 허용하지 않고, 배압을 지원하지 않고, Publisher 인터페이스를 구현하지 않는다. 그래서 리액티브 스트림의 스팩과 호환되지 않는다. 반면 Flowable 타입보다 오버헤드가 적다.
-
Flowable : Flux와 동일한 역할을 하고 Reactive Streams의 Publisher를 구현했다. Flowable API는 Publisher 유형의 인수를 사용할 수 있도록 설계되어 있다.
-
Single : 하나의 요소를 생성하는 스트림을 나타내고, Publisher 인터페이스를 상속하지 않고, 배압전략이 필요 없다.
-
Maybe : Mono 타입과 동일한 의도로 구현되었고, Publisher 인터페이스를 통해 구현하지 않아서 리액티브 스트림과 호환할 수 없다.
-
Completable : RxJava 2.x에는 onError, onComplete 신호만 발생시키고 onNext신호는 생성할 수 없는 Complete 유형이 있다. Publisher 인터페이스를 구현하지 않았다.
Flux와 Mono 시퀀스 생성
Flux
//하나씩 생성
Flux<String> flux1 = Flux.just("1","2","3");
//배열형 으로 생성
Flux<Integer> flux2 = Flux.fromArray(new Integer[]{1,2,3});
//Iterable 타입으로 생성
Flux<Integer> flux3 = Flux.fromIterable(Arrays.asList(1,2,3));
//10~15까지 순차적으로 실행하는 데이터 생성
Flux<Integer> flux4 = Flux.range(10,5);
Mono
Mono와 Flux는 서로 변환이 가능하다. 예를들어 Flux.collectList()는 Mono<List>를 반환하고 Mono.flux는 Flux를 반환한다.
//하나씩 생성
Mono<String> mono1 = Mono.just("1");
Mono<String> mono2 = Mono.justOrEmpty(null);
Mono<String> mono3 = Mono.justOrEmpty(Optional.empty());
Mono는 HTTP 요청이나 DB Query 같은 비동기 작업을 Wrapping 하는데 유용하고, 아래의 메서드들도 제공한다.
-
fromCallable
-
fromRunnable
-
fromSupplier
-
fromFuture
-
fromCompletionStage
Reactive Streams 구독
subscribe() 메서드를 통해 구독할 수 있다.
-
Consumer<? super T> consumer : 데이터를 하나하나 가져옴(onNext)
-
Consumer<? super Throwable> errorConsumer : 에러를 통지(onError)
-
Runnable completeConsumer : 완료를 통지(onComplete)
-
Consumer<? super Subscription > subscriptionConsumer : 구독자가 원하는 동작을 추가
Reactor의 subscribe 오버로드는 버전에 따라 세부 형태가 달라질 수 있다. 여기서는 값, 오류, 완료 신호를 각각 처리할 수 있다는 계약만 기억한다. 아래 예제에서 오류 콜백을 생략하지 않는다.
리액티브 스트림을 발행하고 구독
[subscribe 코드]
Flux<String> flux1 = Flux.just("1","2","3");
flux1.subscribe(new Consumer<String>() {
@Override
public void accept(String s) {
System.out.println(s);
}
}, new Consumer<Throwable>() {
@Override
public void accept(Throwable throwable) {
System.out.println("Exception");
}
}, new Runnable() {
@Override
public void run() {
System.out.println("Complete");
}
});
[결과]
1
2
3
Complete
사용자 정의 Subscriber 구현
Subscriber를 통해 인터페이스를 직접 구현 할 수 있다.
public class ReactMain {
public static void main(String[] args) {
Flux.just("Hello", "world")
.subscribe(subscribers());
}
public static Subscriber<String> subscribers() {
Subscriber<String> subscriber = new Subscriber<String>() {
volatile Subscription subscription;
@Override
public void onSubscribe(Subscription s) {
System.out.println("init onsubscribe");
subscription = s;
subscription.request(Integer.MAX_VALUE);
}
@Override
public void onNext(String o) {
System.out.println("onNext");
System.out.println("Data : "+ o);
}
@Override
public void onError(Throwable t) {
System.out.println("Exception!!");
System.out.println(t);
}
@Override
public void onComplete() {
System.out.println("onComplete");
}
};
return subscriber;
}
}
[결과]
init onsubscribe
onNext
Data : Hello
onNext
Data : world
onComplete
구독을 직접 만들어서 구현하게되면 1차원적 코드 흐름이 깨져 오류를 발생하기 쉽다. 그래서 스스로 배압을 관리하고 가입자에 대한 TCK를 잘 구현해야 한다.
그래서 BaseSubscriber 클래스를 상속하여 사용하는 것이 더 좋은 방법으로 TCK에 호환되는 구독자를 훨씬 쉽게 구현할 수 있다. 그리고 hookOnSubscribe, hookOnNext, hookOnError, hookOnCancel, hookOncomplete 등 메서드 들을 재정의 하여 사용할 수 있다.
public class ReactMain {
public static void main(String[] args) {
Flux.just("Hello", "world")
.subscribe(new CustomSubscriber());
}
private static class CustomSubscriber<T> extends BaseSubscriber<T> {
@Override
protected void hookOnSubscribe(Subscription subscription) {
System.out.println("init request");
request(Integer.MAX_VALUE);
}
@Override
protected void hookOnNext(T value) {
System.out.println("onNext");
System.out.println(value);
}
}
}
init request
onNext
Hello
onNext
world
참조
https://projectreactor.io/docs/core/release/reference/
조립, 구독, 취소를 구분하기
위 Mono.just와 Flux.range는 값을 이미 알고 있는 작은 예다. map을 붙이는 것은 변환 규칙을 조립하는 일이고, 구독은 그 규칙을 실행하도록 요청하는 일이다. HTTP 호출이나 파일 읽기처럼 외부 작업을 담을 때는 값이 만들어지는 시점도 확인해야 한다. 예를 들어 Mono.just(repository.find(id))는 find를 먼저 호출한다. 호출을 구독 시점까지 미루려면 Mono.fromCallable(() -> repository.find(id))를 쓸 수 있지만, 이 역시 블로킹 호출이므로 실행 스레드를 분리해야 한다.
flowchart LR A["Flux/Mono 조립"] --> B["subscribe"] B --> C["onSubscribe"] C --> D["request(n)"] D --> E["onNext 0..N"] E --> F["onComplete 또는 onError"] D --> G["cancel"]
onError와 onComplete는 한 구독의 종료 신호다. 오류 뒤 같은 구독으로 값을 계속 보낸다는 설명은 잘못이다. retry는 기존 구독을 되살리는 것이 아니라 소스에 다시 구독할 수 있다. 외부 API나 DB 쓰기를 재시도하면 부작용이 두 번 발생할 수 있으므로 멱등성 조건을 확인한다.
block()은 호출 스레드가 결과를 기다리는 도구다. 작은 콘솔 예제나 테스트의 경계에서는 결과를 확인할 수 있지만 WebFlux 이벤트 루프에서 호출하면 논블로킹 설계를 깨뜨린다. 서버에서는 반환 타입인 Mono·Flux를 그대로 컨트롤러에 전달한다.
결과값 외에 오류와 완료도 검사하기
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
Flux<Integer> values = Flux.range(1, 3)
.map(number -> number * 2);
StepVerifier.create(values)
.expectNext(2, 4, 6)
.verifyComplete();
이 검사는 세 값의 순서와 정상 완료를 확인한다. 오류가 기대되는 시나리오는 expectError(...) 같은 검증을 별도로 둔다. 무한 스트림은 완료가 없으므로 이 테스트를 그대로 쓰면 끝나지 않는다. 제한된 개수만 취하거나 가상 시간·취소 검사를 사용한다.
운영에서 구독자가 느린 경우 buffer를 무제한으로 늘려 해결하지 않는다. 데이터 소스의 속도, 소비 속도, 버퍼 상한과 드롭·오류 정책을 정한다. Mono와 Flux의 차이는 단순한 성능 등급이 아니라 최대 결과 개수의 계약이다. 사용자 조회 한 건은 Mono<User>, 검색 결과 목록은 Flux<User>가 의도를 분명히 표현한다.