1. Reactive Streams란?
- 비동기 Stream Processing에서 Non-Blocking Backpressure를 제공하기 위한 Specification이다.
- ✅ 비동기 환경에서는 Producer와 Consumer가 서로 다른 속도로 동작함
- ✅ Producer가 Consumer의 처리 속도를 고려하지 않고 계속 데이터를 전달하면 처리되지 못한 데이터가 Queue에 계속 쌓임
- ⚠️ Producer → Queue 증가 → Memory 사용 증가 → 시스템 불안정
- ➡️ Consumer가 자신이 처리할 수 있는 양을 Producer에게 알려주는 Backpressure 구조를 정의함
- ➡️ 서로 다른 Reactive 구현체가 호환될 수 있는 공통 규칙을 제공함
2. Reactive Streams가 해결하려는 문제
Producer와 Consumer의 처리 속도 차이
- 비동기 환경에서는 Producer와 Consumer가 동시에 독립적으로 동작할 수 있기 때문에 처리 속도가 항상 같지 않음
- ✅ 아무런 제어가 없다면 대기 데이터가 계속 증가하는 상황이 발생함
- ⚠️ Buffer에 계속 데이터를 저장하는 방식은 일시적으로 속도 차이를 흡수할 수 있지만 지속적인 해결책은 아님
Backpressure
- Reactive Streams에서는 Consumer가 자신이 처리 가능한 양을 Producer에게 전달한다.
- ➡️ Producer가 일방적으로 Push하는 것이 아니라 Consumer의 Demand를 고려하여 데이터를 전달
예시) Backpressure
더보기
- Consumer → request(3) → Producer → Data 1 → Data 2 → Data 3
- 처리가 끝난 뒤 추가 데이터를 원한다면 다시 요청할 수 있다.
3. Reactive Streams 기본 구조
- Publisher ↔ Subscription ↔ Subscriber
흐름
| 순서 | Method / Signal | 방향 | 의미 |
| 1 | subscribe(Subscriber) | Subscriber → Publisher | Publisher 구독 요청 |
| 2 | onSubscribe(Subscription) | Publisher → Subscriber | 구독을 제어할 Subscription 전달 |
| 3 | request(n) | Subscriber → Publisher | 처리 가능한 데이터 n개 요청 |
| 4 | onNext(data) | Publisher → Subscriber | 실제 데이터 전달 |
| 5 | onComplete() | Publisher → Subscriber | 정상적으로 Stream 종료 |
| 5 | onError(error) | Publisher → Subscriber | Error로 Stream 종료 |
| 별도 | cancel() | Subscriber → Publisher | 구독 취소 |
Interface
| Interface | 역할 | 핵심 의미 |
| Publisher<T> | Stream Data 발행 | Subscriber의 Demand에 따라 데이터 발행 |
| Subscriber<T> | Data / Signal 수신 | 데이터를 받고 필요한 양을 요청 |
| Subscription | Publisher ↔ Subscriber 구독 관계 관리 | Demand와 구독 취소 관리 |
| Processor<T,R> | Subscriber + Publisher | 데이터를 받아 처리한 뒤 다시 발행 |
코드) Publisher
더보기
public interface Publisher<T> {
void subscribe(Subscriber<? super T> subscriber);
}
| Method | 방향 | 의미 |
| subscribe(Subscriber<? super T> s) | Subscriber → Publisher | Subscriber가 Publisher를 구독함 |
코드) Subscriber
더보기
public interface Subscriber<T> {
void onSubscribe(Subscription subscription);
void onNext(T item);
void onError(Throwable throwable);
void onComplete();
}
| Method | 방향 | 의미 |
| onSubscribe(Subscription s) | Publisher → Subscriber | 구독이 시작될 때 Subscription 전달 |
| onNext(T item) | Publisher → Subscriber | 실제 Data 전달 |
| onError(Throwable t) | Publisher → Subscriber | Error 발생으로 Stream 종료 |
| onComplete() | Publisher → Subscriber | 정상적으로 Stream 종료 |
코드) Subscription
더보기
public interface Subscription {
void request(long n);
void cancel();
}
| Method | 방향 | 의미 |
| request(long n) | Subscriber → Publisher | 최대 n개의 Data를 요청 |
| cancel() | Subscriber → Publisher | 구독 취소 및 Data 수신 중단 |
코드) Processor
더보기
public interface Processor<T, R> extends Subscriber<T>, Publisher<R> {
}
4. Reactive Streams 생태계와 구현체
- Reactive Streams는 비동기 Stream과 Non-Blocking Backpressure를 위한 Specification
- 이를 기반으로 JDK와 여러 Reactive Library가 구현체를 제공함
JDK Flow
- JDK 9부터 java.util.concurrent.Flow에 Reactive Streams와 의미적으로 동일한 인터페이스가 제공됨
| Reactive Streams | JDK Flow |
| Publisher | Flow.Publisher |
| Subscriber | Flow.Subscriber |
| Subscription | Flow.Subscription |
| Processor | Flow.Processor |
Reactive Streams Specification
- Reactive Streams는 인터페이스 이름만 정의하는 것이 아니라 각 Interface가 따라야 할 동작 규칙까지 Specification으로 정의함
- ✅ Publisher는 Subscriber의 Demand보다 많은 데이터를 전달하면 안 됨
- ✅ onSubscribe → onNext → onComplete / onError 같은 Signal 순서를 지켜야 함
- ✅ Subscription.request(n)과 cancel()의 동작 규칙
- ✅ Error / Cancel 처리 방식
- ✅ 비동기 경계에서의 Thread Safety 관련 규칙
Project Reactor
- Reactive Streams Specification을 기반으로 실제 Reactive Programming을 쉽게 할 수 있도록 만든 Library
- ✅ Publisher, Subscription, Subscriber, request(n) 등의 계약을 실제 코드에서 쉽게 사용할 수 있도록 API를 제공함
- ex) Mono, Flux, map(), flatMap(), filter(), subscribe()
- ✅ Mono와 Flux는 Reactive Streams의 Publisher 구현체로 볼 수 있음
출처
'Java' 카테고리의 다른 글
| [Java][Tutorial] 1-3. Learning the Java Language: Classes and Objects (0) | 2024.01.14 |
|---|---|
| [Java][Tutorial] 1-2. Learning the Java Language: Language Basics (0) | 2024.01.14 |
| [Java] Stream (3) | 2023.10.22 |
| [Java][Tutorial] 3-3. Collections: Aggregate Operations (0) | 2023.10.20 |
| [Java][Tutorial] 3-2. Collections: Queue, Deque, Map (0) | 2023.10.15 |