목차
Reactive Streams
데이터 생산자가 소비자보다 빠른 비동기 경계에서는 대기열이 계속 커질 수 있다. Reactive Streams는 생산자와 소비자 사이에 수요(demand)를 전달하는 JVM 표준 계약이다. 변환 연산자나 스레드풀을 구현한 라이브러리가 아니라, 서로 다른 구현이 연결될 때 지켜야 할 신호와 배압(backpressure)의 규칙을 정의한다. 명세는 비동기 경계에서의 비차단 처리를 목표로 하지만, 개별 Publisher는 같은 스레드에서 동기적으로 신호를 보낼 수도 있다.

그림의 화살표를 따라가면 Subscriber가 subscribe로 연결하고, Publisher가 onSubscribe로 Subscription을 건네준다. 그다음 소비자가 request(n)으로 받을 수 있는 개수를 알리고, 생산자는 그 범위 안에서 onNext를 보낸다. 데이터를 다 보내면 생산자가 onComplete, 실패하면 onError를 보낸다. 소비자가 완료를 요청하는 프로토콜은 아니다.
네 인터페이스와 신호 순서
| 인터페이스 | 역할 |
|---|---|
Publisher<T> | 구독을 받아 수요에 따라 T를 발행한다. |
Subscriber<T> | onSubscribe, onNext, onError, onComplete를 받는다. |
Subscription | 특정 구독의 request(long n)과 cancel()을 제공한다. |
Processor<T,R> | Subscriber<T>이면서 Publisher<R>인 처리 단계다. |
이 네 인터페이스는 org.reactivestreams 패키지의 API다. 옛 글에서 언급하던 rxjava-reactive-streams는 RxJava 1과 명세 사이의 어댑터 라이브러리였으며 인터페이스의 소유 패키지가 아니다. 아래 예제는 RxJava 2의 Flowable을 사용한다. Flowable은 Reactive Streams Publisher를 구현하지만 RxJava Observable과 Observer 조합은 같은 Subscription.request(n) 계약을 제공하지 않는다. 두 타입을 혼동하면 배압을 적용했다고 착각할 수 있다.
한 구독에서 관찰할 수 있는 신호 순서는 다음과 같다.
onSubscribe → onNext × 0회 이상 → (onComplete | onError)
onSubscribe는 구독당 한 번이다. 두 종료 신호 중 하나가 온 뒤에는 추가 신호를 보내지 않는다. onNext 수는 요청한 총량을 넘지 않지만 생산자가 요청량을 모두 채우기 전에 정상 완료하거나 실패할 수 있다. onComplete와 onError에는 별도 데이터 요청이 필요하지 않다. null 항목은 허용하지 않으며 한 Subscriber에게 보내는 신호는 순차적으로 전달해야 한다.

두 번째 그림에서 request가 다시 나타나는 이유는 소비자가 처리 여력을 회복한 뒤 수요를 보충하기 때문이다. 예를 들어 처음 두 건을 요청했다면, 생산자는 누적 두 건을 초과해 보낼 수 없다. 소비자가 한 건을 처리하고 다시 request(1)을 부르면 총 요청량은 세 건이 된다. 이 숫자는 처리 완료 확인이나 거래 커밋 신호가 아니다.
요청량이 실제 출력을 제한하는 예
import io.reactivex.Flowable;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;
public class DemandExample {
public static void main(String[] args) {
Flowable.range(1, 4).subscribe(new Subscriber<Integer>() {
private Subscription subscription;
@Override
public void onSubscribe(Subscription s) {
subscription = s;
s.request(2);
}
@Override
public void onNext(Integer value) {
System.out.println(value);
subscription.request(1);
}
@Override
public void onError(Throwable error) {
error.printStackTrace();
}
@Override
public void onComplete() {
System.out.println("완료");
}
});
}
}
이 동기 range 예제는 1, 2, 3, 4, 완료를 출력한다. 처음 두 건의 수요를 만든 뒤 onNext에서 한 건씩 보충하기 때문이다. onNext에서 request(1)을 제거하면 처음 두 값까지만 받는다. 그 상태에서 기다리는 것이 오류는 아니다. 생산자가 더 보내려면 추가 수요가 필요하다. request(Long.MAX_VALUE)는 사실상 제한 없는 수요를 나타내므로 처리량과 메모리 한계를 모르는 채 기본값처럼 쓰면 안 된다.
request(n)의 n은 양수여야 한다. request(0)으로 일시 중지하는 방식은 명세 위반이며 오류로 처리된다. 일시 정지는 더 요청하지 않는 방식으로 표현한다. 여러 스레드에서 직접 request와 cancel을 호출하는 소비자를 만들 때에는 호출 사이의 순차 관계를 보장해야 한다. 위 예제처럼 같은 신호 처리 흐름에서 호출하는 편이 단순하다.
취소와 운영 시 판단
소비자가 더 이상 필요하지 않으면 subscription.cancel()을 호출한다. 취소는 종료 통지 요청과 다르다. 취소 직후에도 이미 전달 중인 신호가 잠시 보일 수 있고, onComplete를 반드시 받는 것도 아니다. 재구독은 새 Subscriber와 새 Subscription을 기준으로 생각해야 하며, 기존 소비자 인스턴스를 재사용해도 상태가 자동 초기화되는 것은 아니다.
배압 계약이 있다고 모든 메모리 문제가 사라지지는 않는다. 중간 연산자가 데이터를 버퍼링하거나 외부 이벤트 소스가 속도를 늦출 수 없다면 별도의 큐 한계, 드롭/최신값 정책, 과부하 오류를 선택해야 한다. 소비자에게 request(1)만 반복하면 왕복 비용이 커질 수 있으므로 처리 가능한 크기의 묶음으로 요청하는 편이 낫다. 실제 시스템에서는 구독별 미처리 건수, 큐 길이, 처리 지연, 취소율을 함께 본다.
참고: Reactive Streams JVM 명세, Reactive Streams Subscriber API, RxJava 2 Flowable API.