목차
리액터 개요
애플리케이션 코드를 개발할 때 명령형(Imperative)과 리엑티브(Reactive) 두 가지 형태로 코드를 작성할 수 있다.
-
명령형 : 순차적으로 연속되는 작업이고, 각 작업은 한 번에 하나씩 그리고 이전 작업 다음에 실행한다. 데이터는 모아서 처리되고 이전 작업 데이터 처리를 끝낸 후에 다음 작업으로 넘어갈 수 있다.
-
리액티브(반응형) : 값을 발행하고 소비하는 관계를 선언하고, 구독·요청량·오류·완료 신호로 흐름을 제어한다. 병렬 실행은 별도로 스케줄러를 지정했을 때 일어날 수 있으며 기본 보장은 아니다.
이런 리액터 프로젝트는 비동기 파이프라인을 구축할 때 콜백 지옥과 깊게 중첩된 코드를 생략하는 목적으로 설계되었다.
현재 사용하는 Project Reactor 3.x는 Reactive Streams 사양의 Publisher를 구현한다. 과거 1.x·2.x의 API와 현재 Flux·Mono API를 혼용하지 않는다.
리엑티브 프로그래밍
보통 대부분의 언어는 동시 프로그래밍을 지원하고 스레드로 동시성을 관리하는 것은 쉽지 않은데 스레드가 많은 수록 더 복잡하기 때문이다.
리액티브 프로그래밍은 명령형 프로그래밍의 대안이 되는 패러다임으로 본질은 함수적이면서 선언적이다.
즉, 순차적으로 처리하는 작업단계가 아니라 데이터가 흘러가는 파이프라인(pipeline)이나 스트림(stream)을 포함한다.
또한 데이터 전체를 사용할 수 있을 때까지 기다리지 않고 사용가능한 데이터가 있을 때마다 처리되므로 입력되는 데이터는 무한할 수 있다.
리엑티브 스트림 정의
리엑티브 스트림은 4개의 인터페이스로 정의할 수 있다.
-
publisher(발행자)
-
Subscriiber(구독자)
-
Subscription(구독)
-
Processor(프로세서)
Publisher는 하나의 Subscription당 Subscriber에 발행(전송)하는 데이터를 생성한다.
Publisher 인터페이스에는 Subscriber가 Publisher를 구독 신청할 수 있는 Subscribe() 메서드 한개를 선언
public interface Publisher<T> {
void subscribe(Subscriber<? super T> subscriber);
}
Subscriber가 구독 신청되면 Publisher로 부터 이벤트르 수신할 수 있고, 이 이벤트들은 Subscriber 인터페이스의 메서드를 통해 전송
public interface Subscriber<T> {
void onSubscribe(Subscription sub);
void onNext(T item);
void onError(Throwable ex);
void onComplete();
}
Subscriber가 수신할 첫 번째 이벤트는 onSubscribe()의 호출을 통해 이루어지며, Publisher가 onSubscribe()를 호출할 때 이 메서드의 인자로 Subscription() 객체를 Subscriber에 전달.
Subscriber는 Subscription객체를 통해 구독을 관리
public interface Subscription {
void request(long n);
void cancel();
}
Subscriber는 request()를 호출하여 전송되는 데이터를 요청하거나, 더 이상 데이터를 수신하지 않고 구독을 취소하는 것을 나타내기 위해 cancel()를 호출할 수 있다.
request()를 호출할 때 Subscriber는 받고자 하는 데이터의 항목 수를 나타내는 long 타입의 값을 인자로 전달하는데 이것을 백 프레셔(Back-Pressure)라고 하며 Subscriber가 처리할 수 있는 것보다 더 많은 데이터를 Publisher가 전송되는 것을 막아준다.
백 프레셔(Back-Pressure) : 데이터가 소비하는(읽는) 컨슈머가 처리하 수 있는 만큼 전달 데이터를 제한하여 빠른 데이터 소스로부터 데이터 전달 폭주를 피할 수 있는 수단.
Subscriber의 요청이 완료되면 데이터가 스트림을 통해 전달되며, onNext() 메서드가 호출되어 Publisher가 전송하는 데이터가 Subscriber에게 전달되고 에러가 발생할 경우 onError()를 호출한다.
그리고 Publisher에서 전송할 데이터가 없고, 데이터를 생성하지 않는다면 Publisher가 onComplete()를 호출하여 작업이 끝난다고 Subscriber에게 알려준다.
Processor 인터페이스는 Subscriber 인터페이스와 Publisher 인터페이스를 결합한 것
public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
}
Processor는 데이터를 수신하고 처리한 후 Publisher는 처리 결과를 자신의 Subscriber에게 발행한다.
신호가 오가는 순서를 따라가기
Publisher와 Subscriber의 계약은 데이터를 담는 컬렉션 인터페이스와 다르다. 구독 후 구독자는 Subscription을 받고, request(n)으로 필요한 양을 알린다. 발행자는 허용된 범위에서 onNext를 보낸다. 끝나면 onComplete, 실패하면 onError 중 하나로 종료한다.
sequenceDiagram participant S as Subscriber participant P as Publisher S->>P: subscribe() P-->>S: onSubscribe(subscription) S->>P: request(2) P-->>S: onNext(1) P-->>S: onNext(2) S->>P: request(2) P-->>S: onNext(3), onNext(4) S->>P: cancel() 또는 추가 요청
그림에서 처음 두 값 뒤 다음 값이 오려면 추가 요청이 필요하다. request(2)는 동시에 두 스레드를 만든다는 뜻이 아니라 수요량(demand)을 뜻한다. cancel()은 더 이상 필요 없는 스트림을 중단하려는 신호다. 취소와 이미 진행 중인 외부 I/O 중단이 언제까지 보장되는지는 해당 발행자의 구현에 달려 있다.
Flux와 Mono로 한정해 보기
Flux<T>는 0개 이상의 값, Mono<T>는 0개 또는 1개의 값을 표현한다. 연산자를 이어 붙여도 실제 처리는 보통 구독이 시작될 때 진행된다. 따라서 다음 코드의 map을 정의하는 것과 값을 얻는 시점을 구별해야 한다.
import reactor.core.publisher.Flux;
Flux<Integer> doubled = Flux.range(1, 3)
.map(number -> number * 2)
.doOnNext(value -> System.out.println("value=" + value));
// 구독이 시작되면 2, 4, 6이 처리된다.
doubled.subscribe();
Flux.range는 이미 메모리에 있는 작은 동기 발행자라 이 예제로 네트워크의 비동기 성능을 증명할 수 없다. map 내부에서 오래 걸리는 JDBC 호출을 하면 호출한 스레드가 기다린다. WebFlux의 이벤트 루프 위에서 같은 작업을 하면 다른 요청 처리까지 지연될 수 있다. 비동기 드라이버를 쓰거나, 블로킹 작업의 실행 경계를 명시적으로 분리해야 한다.
import reactor.core.publisher.BaseSubscriber;
import reactor.core.publisher.Flux;
import org.reactivestreams.Subscription;
Flux.range(1, 5).subscribe(new BaseSubscriber<Integer>() {
private int received;
@Override
protected void hookOnSubscribe(Subscription subscription) {
request(2);
}
@Override
protected void hookOnNext(Integer value) {
System.out.println(value);
received++;
if (received % 2 == 0) request(2);
}
});
이 구독자는 처음 두 개를 요청하고 두 개를 처리할 때마다 다음 두 개를 요청한다. 마지막 요청에서 필요한 값보다 많이 요청해도 발행자는 다섯 값을 보낸 뒤 완료한다. 반면 무한 발행자에 request(Long.MAX_VALUE)를 쓰고 뒤쪽 소비가 느리면 버퍼링 정책에 따라 메모리 압박이 생긴다.
오류 경로도 반드시 정한다. onError로 종료된 스트림은 같은 구독에서 이어서 onNext를 보내지 않는다. onErrorResume 같은 복구 연산자는 기본값을 내거나 대체 발행자로 옮길 수 있지만, 원래 실패를 무작정 숨기면 데이터 누락을 정상 결과로 오해한다. 구독 취소·타임아웃·오류를 각각 테스트한다.
참고: Reactive Streams 사양, Project Reactor 참고 문서, Reactor Core API.
