Java

[Reactive Streams] 1. 비동기 스트림과 Backpressure

noahkim_ 2026. 8. 17. 16:44

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 구현체로 볼 수 있음

 

 

출처