개발자가 리액티브한 코드를 작성하기 위해서는 이러한 코드 구성을 용이하게 해주는 리액티브 라이브러리가 있어야 된다. 이 리액티브 라이브러리를 어떻게 구현할지 정의해 놓은 별도의 표준 사양이 있는데, 이것을 바로 리액티브 스트림즈라고 부른다.
리액티브 스트림즈는 한마디로 '데이터 스트림을 Non-Blocking이면서 비동기적인 방식으로 처리하기 위한 리액티브 라이브러리의 표준 사양'이라고 표현할 수 있다.
리액티브 스트림즈를 구현한 구현체로 RxJava, Reactor, Akka Streams, Java 9 Flow API 등이 있는데, 그중에서 Spring Framework와 가장 궁합이 잘 맞는 것은 Reactor이다.
리액티브 스트림즈를 통해 구현해야 되는 API 컴포넌트에는 Publisher, Subscriber, Subscription, Processor가 있다.
| 컴포넌트 | 설명 |
|---|---|
| Publisher | 데이터를 생성하고 통지(발행, 게시, 방출)하는 역할을 한다. |
| Subscriber | 구독한 Publisher로부터 통지(발행, 게시, 방출)된 데이터를 전달받아서 처리하는 역할을 한다. |
| Subscription | Publisher에 요청할 데이터의 개수를 지정하고, 데이터의 구독을 취소하는 역할을 한다. |
| Processor | Publisher와 Subscriber의 기능을 모두 가지고 있다. 즉, Subscriber로서 다른 Publisher를 구독할 수 있고, Publisher로서 다른 Subscriber가 구독할 수 있다. |

- 먼저 Subscriber는 전달받을 데이터를 구독한다(subscribe).
- 다음으로 Publisher는 데이터를 통지(발행, 게시, 방출)할 준비가 되었음을 Subscriber에 알린다(onSubscribe).
- Publisher가 데이터를 통지할 준비가 되었다는 알림을 받은 Subscriber는 전달받기를 원하는 데이터의 개수를 Publisher에게 요청한다(Subscription.request).
- 다음으로 Publisher는 Subscriber로부터 요청받은 만큼의 데이터를 통지한다(onNext).
- 이렇게 Publisher와 Subscriber 간에 데이터 통지, 데이터 수신, 데이터 요청의 과정을 반복하다가 Publisher가 모든 데이터를 통지하게 되면 마지막으로 데이터 전송이 완료되었음을 Subscriber에게 알린다(onComplete). 만약에 Publisher가 데이터를 처리하는 과정에서 에러가 발생하면 에러가 발생했음을 Subscriber에게 알린다(onError).
Subscriber가 Subscription.request를 통해 데이터의 요청 개수를 지정하는 이유는, Publisher와 Subscriber는 각각 다른 스레드에서 비동기적으로 상호작용하는 경우가 대부분이기 때문에 Publisher가 통지하는 속도가 Publisher로부터 통지된 데이터를 Subscriber가 처리하는 속도보다 더 빠르면 처리를 기다리는 데이터는 쌓이게 되고, 이는 결과적으로 시스템 부하가 커지는 결과를 낳게 되기 때문이다.
이러한 문제를 방지하기 위해 Subscription.request를 통해 데이터 개수를 제어하는 것이다.
리액티브 스트림즈의 컴포넌트는 실제 코드에서 인터페이스 형태로 정의되며, 이 인터페이스들을 구현해서 해당 컴포넌트를 사용하게 된다.
public interface Publisher<T> {
public void subscribe(Subscriber<? super T> s;
}
subscriber 메서드는 파라미터로 전달받은 Subscriber를 등록하는 역할을 한다.
Kafka의 경우 Publisher와 Subscriber 중간에 메시지 브로커가 있고 이 브로커 내에 여러 개의 토픽이 존재하는데, Publisher와 Subscriber는 브로커에 있는 특정 토픽을 바라보는 구조로 이루어져 있다. 그래서 Kafka에서의 Publisher와 Subscriber는 각각 브로커 내의 특정 토픽만 바라보면 되기 때문에 Publisher는 특정 토픽으로 메시지 데이터를 전송하기만 하면 되고, Subscriber는 특정 토픽을 구독하고 해당 토픽에 전달되는 메시지 데이터를 전달받기만 하면 된다. 이는 Publisher와 Subscriber의 느슨한 결합 구조라고 볼 수 있다.
반면에 리액티브 스트림즈에서의 Publisher와 Subscriber는 개념상으로는 Subscriber가 구독하는 것이 맞는데 실제 코드상에서는 Publisher가 subscribe 메서드의 파라미터인 Subscriber를 등록하는 형태로 구독이 이루어진다.
public interface Subscriber<T> {
public void onSubscribe(Subscription s);
public void onNext(T t);
public void onError(Throwable t);
public void onComplete();
}
- onSubscribe: 구독 시작 시점에 어떤 처리를 하는 역할은 한다. 여기서의 처리는 Publisher에게 요청할 데이터의 개수를 지정하거나 구독을 해지하는 것을 의마하는데, 이것은 onSubscribe 메서드의 파라미터로 전달되는 Subscription 객체를 통해서 이루어진다.
- onNext: Publisher가 통지한 데이터를 처리하는 역할을 한다.
- onError: Publisher가 데이터 통지를 위한 처리 과정에서 에러가 발생했을 때 해당 에러를 처리하는 역할을 한다.
- onComplete: Publisher가 데이터 통지를 완료했음을 알릴 때 호출되는 메서드이다. 데이터 통지가 정상적으로 완료될 경우에 어떤 후처리를 해야 한다면 onComplete 메서드에서 처리 코드를 작성하면 된다.
public interface Subscription {
public void request(long n);
public void cancel();
}
Subscription 인터페이스는 Subscriber가 구독한 데이터의 개수를 요청하거나 또는 데이터 요청의 취소, 즉 구독을 해지하는 역할을 한다.
request 메서드를 통해서 Publisher에게 데이터의 개수를 요청할 수 있고, cancel 메서드를 통해서 구독을 해지할 수 있다.
public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {}
별도로 구현해야 하는 메서드가 없다. 다른 인터페이스들과 다른 점은 Subscriber 인터페이스와 Publisher 인터페이스를 상속한다는 것이다.
- Publisher가 Subscriber 인터페이스 구현 객체를 subscribe 메서드의 파라미터로 전달한다.
- Publisher 내부에서는 전달받은 Subscriber 인터페이스 구현 객체의 onSubscribe 메서드를 호출하면서 Subscriber의 구독을 의미하는 Subscription 인터페이스 구현 객체를 Subscriber에게 전달한다.
- 호출된 Subscriber 인터페이스 구현 객체의 onSubscribe 메서드에서 전달 받은 Subscription 객체를 통해 전달받을 데이터의 개수를 Publisher에게 요청한다.
- Publisher는 Subscriber로부터 전달받은 요청 개수만큼의 데이터를 onNext 메서드를 호출해서 Subscriber에게 전달한다.
- Publisher는 통지할 데이터가 더 이상 없을 경우 onComplete 메서드를 호출해서 Subscriber에게 데이터 처리 종료를 알린다.
Publisher와 Subscriber 간에 주고받는 상호작용을 signal이라고 표현한다.
리액티브 스트림즈의 인터페이스 코드에서 볼 수 있는 onSubscribe, onNext, onComplete, onError, request 또는 cancel 메서드를 리액티브 스트림즈에서는 signal이라고 표현한다.
onSubscribe, onNext, onComplete, onError 메서드는 Subscriber 인터페이스에 정의되지만 이 메서드들을 실제 호출해서 사용하는 주체는 Publisher이기 때문에 Publisher가 Subscriber에게 보내는 Signal이라고 볼 수 있다.
그리고 request와 cancel 메서드는 Subscription 인터페이스 코드에 정의되지만 이 메서드들을 실제로 사용하는 주체는 Subscriber이므로 Subscriber가 Publisher에게 보내는 Signal이라고 볼 수 있다.
Demand는 Subscriber가 Publisher에게 요청하는 데이터를 의미한다. 더 구체적으로 얘기하면 Publisher가 아직 Subscriber에게 전달하지 않은 Subscriber가 요청한 데이터를 말한다.
Publisher가 Subscriber에게 데이터를 전달하는 것을 일컬을 때 데이터를 '통지(발행, 게시, 방출)한다'라는 용어를 사용했다. 리액티브 프로그래밍과 관련된 영문 문서에서 가장 많이 볼 수 있는 용어는 바로 Emit이다.
public class Example {
public static void main(String[]) {
Flux
.just(1, 2, 3, 4, 5, 6)
.filter(n -> n % 2 == 0)
.map(n -> n * 2)
.subscribe(System.out::println);
}
}
Reactor 코드로 구성되어 있다.
just 메서드를 사용해서 데이터를 생성한 후 emit하게 되는데, 여기서 just 메서드는 리액티브 스트림즈의 컴포넌트 중에서 Publisher의 역할을 한다. 메서드 체인 방식으로 호출할 수 있는 이유는 호출하는 각각의 메서드들이 모두 같은 타입의 객체를 반환하기 때문이다. 각각의 메서드들이 반환하는(subscribe 메서드 제외) 반환 값이 모두 Flux 타입의 객체이기 때문에 이런 식의 메서드 호출이 가능하다.
데이터 스트림의 관점에서 볼 때, just 메서드 호출을 통해 반환된 Flux 입장에서는 filter 메서드 호출을 통해 반환된 Flux가 자신보다 더 하위에 있기 때문에 Downstream이 된다.
반면에 filter 메서드 호출을 통해 반환된 Flux 입장에서는 just 메서드 호출을 통해 반환된 Flux가 자신보다 더 상위에 있기 때문에 Upstream이 된다.
Sequence는 Publisher가 emit하는 데이터의 연속적인 흐름을 정의해 놓은 것 자체를 의미하는데, 이 Sequence는 Operator 체인 형태로 정의된다. Flux를 통해서 데이터를 생성, emit하고 filter 메서드를 통해서 필터링한 후, map 메서드를 통해 변환하는 과정 자체를 바로 Sequence라고 부른다.
just, filter, map 같은 메서드들을 리액티브 프로그래밍에서는 Operator라고 부른다.
Data Source, Source Publisher, Source Flux 등을 들 수 있는데, 대부분 '최초의'라는 의미로 사용된다. Source와 비슷한 의미로 종종 Original이라는 용어도 사용된다.
리액티브 스트림즈 표준 사양에는 리액티브 스트림즈 컴포넌트를 어떻게 구현해야 되는지에 대한 규칙들이 컴포넌트별로 정의되어 있다.
| 번호 | 규칙 |
|---|---|
| 1 | Publisher가 Subscriber에게 보내는 onNext signal의 총 개수는 항상 해당 Subscriber의 구독을 통해 요청된 데이터의 총 개수보다 더 작거나 같아야 한다. |
| 2 | Publisher는 요청된 것보다 적은 수의 onNext signal을 보내고 onComplete 또는 onError를 호출하여 구독을 종료할 수 있다. |
| 3 | Publisher의 데이터 처리가 실패하면 onError signal을 보내야 한다. |
| 4 | Publisher의 데이터 처리가 성공적으로 종료되면 onComplete signal을 보내야 한다. |
| 5 | Publisher가 Subscriber에게 onError 또는 onComplete signal을 보내는 경우 해당 Subscriber의 구독은 취소된 것으로 간주되어야 한다. |
| 6 | 일단 종료 상태 signal을 받으면 (onError, onComplete) 더 이상 signal이 발생되지 않아야 한다. |
| 7 | 구독이 취소되면 Subscriber는 결국 signal을 받는 것을 중지해야 한다. |
- 2번 규칙의 경우 예외가 존재한다. Publisher가 처리할 데이터가 끊임없이 발생하는 무한 스트림의 경우, 처리 중 에러가 발생하기 전까지는 종료 자체가 없기 때문에 2번 규칙에 대한 예외적인 경우에 해당한다. IoT 디바이스 센서에서 발생하는 데이터처럼 끊임없이 발생하는 데이터를 떠올려 보면 쉽게 이해할 수 있다.
- 3번 규칙은 Publisher가 처리를 진행할 수 없는 상황이 되면 Subscriber에게 알려서 Publisher에서 발생한 실패를 Subscriber가 처리할 수 있는 기회를 가지도록 하는 데 주목적이 있다.
- 4번 규칙은 publisherrk 종료 상태가 되었음을 Subscriber에게 알림으로써 Subscriber가 리소르를 정리하는 등의 후처리를 할 수 있도록 하는 데 그 목적이 있다.
| 번호 | 규칙 |
|---|---|
| 1 | Subscriber는 Publisher로부터 onNext signal을 수신하기 위해 Subscription.request(n)를 통해 Demand signal을 Publisher에게 보내야 한다. |
| 2 | Subscriber.onComplete() 및 Subscriber.onError(Throwable t)는 Subscription 또는 Publisher의 메서드를 호출해서는 안된다. |
| 3 | Subscriber.onComplete() 및 Subscriber.onError(Throwable t)는 signal을 수신한 후 구독이 취소된 것으로 간주해야 한다. |
| 4 | 구독이 더 이상 필요하지 않은 경우 Subscriber는 Subscription.cancel()을 호출해야 한다. |
| 5 | Subscriber.onSubscriber()는 지정된 Subscriber에 대해 최대 한 번만 호출되어야 한다. |
- 1번 규칙의 목적은 데이터를 언제, 얼마나 수신할 수 있는지를 결정하는 책임이 Subscriber에게 있다는 것을 확립하는 것이다. 리액티브 스트림즈에서는 한 번에 하나의 데이터를 요청하기보다는 Subscriber가 처리할 수 있는 적절한 상한선만큼의 데이터 개수 요청을 권장한다.
- 2번 규칙은 Subscriber가 완료 signal 또는 에러 signal을 처리하는 동안 Publisher/Subscription과 Subscriber 간의 순환 및 경쟁 조건(Race Condition)을 방지하기 위함이다.
@Override public void onComplete() { // 잘못된 예: onComplete()에서 다시 Subscription의 메서드를 호출함. subscription.request(1); // 순환 호출 발생 가능성 // 정답: onComplete()는 아무런 추가 호출 없이 끝내야 함. } @Override public void onError(Throwable t) { // 잘못된 예: onError()에서 다시 Publisher의 메서드를 호출함. publisher.subscribe(this); // 새로운 구독을 시도할 경우 경쟁 조건 발생 가능성 // 정답: onError()는 문제가 발생했음을 알리고 끝내야 함. }
- 4번 규칙의 목적은 Subscriber가 구독이 더 이상 필요하지 않을 때 명시적으로 구독을 취소함으로써 해당 구독이 유지하고 있는 리소스를 적절한 시기에 안전하게 해제할 수 있도록 하는 것이다.
- 5번 규칙의 의미는 동일한 구독자가 최대 한 번만 구독할 수 있다는 의미와 같다.
| 번호 | 규칙 |
|---|---|
| 1 | 구독은 Subscriber가 onNext 또는 onSubscribe 내에서 동기적으로 Subscription.request를 호출하도록 허용해야 한다. |
| 2 | 구독이 취소된 후 추가적으로 호출되는 Subscription.request(long n)는 효력이 없어야 한다. |
| 3 | 구독이 취소된 후 추가적으로 호출되는 Subscription.cancel()은 효력이 없어야 한다. |
| 4 | 구독이 취소되지 않은 동안 Subscription.request(long n)의 매개변수가 0보다 작거나 같으면 java.lang.IllegalArgumentException과 함께 onError signal을 보내야 한다. |
| 5 | 구독이 취소되지 않은 동안 Subscription.cancel()은 Publisher가 Subscriber에게 보내는 signal을 결국 중지하도록 요청해야 한다. |
| 6 | 구독이 취소되지 않은 동안 Subscription.cancel()은 Publisher에게 해당 구독자에 대한 참조를 결국 삭제하도록 요청해야 한다. |
| 7 | Subscription.cancel(), Subscription.request() 호출에 대한 응답으로 예외를 던지는 것을 허용하지 않는다. |
| 8 | 구독은 무제한 수의 request 호출을 지원해야 하고 최대 2^63 - 1개의 Demand를 지원해야 한다. |
- 7번 규칙은 Subscription의 메서드를 호출했을 때 메서드 내부로 예외가 던져지지 않도록 규정한다. Java에서 일반적으로 어떤 메서드를 호출하여 예외가 발생하면 메서드를 호출한 쪽으로 예외를 던지는데, 리액티브 스트림즈에서는 예외가 발생하면 해당 예외를 onError signal과 함께 보내도록 규정한다.
RxJava에서 Rx는 Reactive Extensions라는 의미이다. RxJava는 .NET 환경의 리액티브 확장 라이브러리를 넷플릭스에서 Java 언어로 포팅하여 만든 JVM 기반의 대표적인 리액티브 확장 라이브러리이다.
RxJava는 2.0부터 리액티브 스트림즈 사양을 지원하기 시작했는데, 이 때문에 리액티브 스트림즈 사양을 지원하지 않았던 RxJava 1.x와 2.0의 기능들이 함께 사용되고 있다.
1.x 버전과 2.0 이후 버전의 근본적인 동작 구조는 크게 차이 나지 않지만 두 버전에는 Backpressure를 지원하느냐 그렇지 않느냐 하는 차이점이 있다.
Project Reactor는 Spring Framework 팀에 의해 주도적으로 개발된 리액티브 스트림즈의 구현체이다.
Akka는 JVM상에서의 동시성과 분산 애플리케이션을 단순화 해주는 오픈소스 툴킷이다.
Akka는 Actor Model을 적극적으로 사용하는 대표적인 기술이다. 이러한 Actor들 간의 통신은 메시지를 통해서만 이루어지고 Actor들은 서로 독립적이기 때문에 느슨한 결합과 높은 응집력이 보장된다.
이 Akka라는 Actor 기반의 동시성 모델을 사용하는 툴킷 위에 리액티브 스트림즈를 구현한 것이 바로 Akka Streams이다.
Flow API는 Reactor, RxJava, Akka Streams처럼 리액티브 스트림즈를 구현한 구현체가 아니라 리액티브 스트림즈의 표준 사양이 SPI(Service Provider Interface)로써 Java API에 정의되어 있다.
Java에서 이 Flow API를 정의해 높은 이유를 유추해 보자면, JDBC처럼 사용자들이 리액티브 스트림즈와 관련된 API를 사용하기 위해서 하나의 인터페이스만 바라볼 수 있는 단일 창구로서의 역할을 기대하기 때문이라고 생각할 수 있다.
또 리액티브 스트림즈를 구현한 여러 구현체들을 Flow API로 변환하는 상호운용성 측면에서 생각해 볼 수 있다.