본문으로 건너뛰기
목록으로

Reactive Streams 표준 — Publisher·Subscriber·백프레셔와 JDK Flow

Johny Cho
Software Engineer @ Kurly

Reactor의 Mono·Flux, RxJava의 Observable, Akka Streams… 리액티브 라이브러리는 여럿인데, 이들이 서로 연동되고 같은 개념을 공유하는 건 밑바탕에 공통 표준이 있기 때문입니다. 그 표준이 Reactive Streams — "비동기 스트림을 논블로킹 백프레셔로 주고받는" 최소 규격입니다. WebFlux를 쓰든 다른 라이브러리를 쓰든 이 토대를 알면 "왜 구독해야 실행되나", "백프레셔가 뭔가"가 명확해집니다. (WebFlux에서의 활용은 Spring WebFlux vs MVC 참고.)

왜 표준이 필요했나​

비동기 스트림에는 오래된 문제가 있습니다 — 생산자가 소비자보다 빠를 때입니다.

  • 생산자가 데이터를 밀어내기(push) 만 하면, 느린 소비자는 밀린 데이터를 쌓다가 메모리가 넘치거나(OOM) 데이터를 버리게 됩니다.
  • 그렇다고 소비자가 하나씩 당겨오기(pull) 만 하면, 매번 요청·응답을 주고받아 비동기의 이점이 사라집니다.
Reactive Streams의 해법은 "소비자가 감당할 개수를 미리 알려 주고, 생산자는 그만큼만 비동기로 보내는" 방식입니다 — push와 pull의 절충인 이 흐름 제어가 곧 백프레셔(backpressure) 입니다.

여러 라이브러리가 이걸 제각각 구현하면 서로 연동이 안 되니, 2013~2015년 Netflix·Pivotal(현 VMware)·Lightbend 등이 모여 공통 인터페이스 4개로 표준화한 것이 Reactive Streams입니다.

네 가지 역할​

표준은 인터페이스 4개뿐입니다. 이 관계가 전부입니다.

  • Publisher(발행자) — 데이터를 만들어 발행하는 쪽. subscribe(Subscriber) 하나만 가집니다.
  • Subscriber(구독자) — 데이터를 받는 쪽. onSubscribe·onNext(값)·onComplete(끝)·onError(오류) 네 콜백을 가집니다.
  • Subscription(구독) — Publisher와 Subscriber를 잇는 연결. 구독자가 request(n)으로 받을 개수를 요청하고 cancel()로 취소합니다.
  • Processor — Publisher이자 Subscriber인 중간 변환 단계(직접 구현할 일은 드뭅니다).
PlantUML 코드
@startuml
title Reactive Streams 네 역할
rectangle "Publisher\n데이터를 만들어 발행" as P #EEF3FB
rectangle "Subscriber\nonNext·onComplete·onError 수신" as S #E9F2E9
rectangle "Subscription\nrequest(n)·cancel — 흐름 제어" as Sub #FFF3E0
rectangle "Processor\nPublisher+Subscriber (중간 변환)" as Proc #F3EEFB
P --> S : subscribe / onSubscribe
S --> Sub : request(n) / cancel
Sub --> S : onNext × n
Proc ..> P : (구현: Publisher)
Proc ..> S : (구현: Subscriber)
@enduml

PlantUML Reactive Streams 네 역할

인터페이스로 보면 이렇게 단순합니다.

public interface Publisher<T> {
void subscribe(Subscriber<? super T> s);
}
public interface Subscriber<T> {
void onSubscribe(Subscription s);
void onNext(T item);
void onError(Throwable t);
void onComplete();
}
public interface Subscription {
void request(long n); // 받을 개수 요청 (백프레셔의 핵심)
void cancel();
}

구독과 백프레셔 — request(n)의 흐름​

핵심은 구독자가 request(n)으로 "감당할 만큼만" 당겨 간다는 점입니다. 생산자는 요청받은 개수까지만 onNext를 보냅니다.

PlantUML 코드
@startuml
title Reactive Streams — 구독과 request(n) 기반 백프레셔
participant "Subscriber\n(구독자)" as S
participant "Subscription\n(구독)" as Sub
participant "Publisher\n(발행자)" as P
S -> P : subscribe(subscriber)
P -> S : onSubscribe(subscription)
note right of S : 이제부터 구독자가\n"받을 개수"를 직접 요청
S -> Sub : request(2)
Sub -> S : onNext(item1)
Sub -> S : onNext(item2)
note left of S : 감당한 만큼 처리 후\n다시 요청 (pull = 백프레셔)
S -> Sub : request(2)
Sub -> S : onNext(item3)
Sub -> S : onComplete()
@enduml

PlantUML Reactive Streams 구독·백프레셔 시퀀스

이 pull 기반 요청 덕분에 느린 소비자가 빠른 생산자에게 "천천히"라고 말할 수 있습니다 — 표준이 보장하는 건 결국 이 흐름 제어 하나입니다. (request(Long.MAX_VALUE)로 요청하면 사실상 제한 없는 push가 됩니다.)

발행과 구독이 이어지는 구조​

실제 코드는 소스 → 연산자 여러 단계 → 최종 구독자가 사슬처럼 연결된 파이프라인입니다. 방향이 두 갈래인 게 핵심입니다.

  • 데이터는 위→아래(소스 → 구독자) 로 onNext를 타고 내려갑니다.
  • 수요(request(n))는 아래→위(구독자 → 소스) 로 올라갑니다. 각 연산자는 요청을 위로 전파합니다.
PlantUML 코드
@startuml
title 발행 → 연산자 → 구독 — 파이프라인 구조
skinparam defaultTextAlignment center
rectangle "Publisher (소스)\nDB·HTTP·이벤트" as P #EEF3FB
rectangle "연산자 단계\nmap · filter · flatMap" as OP #FFF3E0
rectangle "Subscriber (최종 소비)\nonNext·onComplete·onError" as S #E9F2E9

P -down-> OP : onNext (데이터, 아래로)
OP -down-> S : onNext (변환된 데이터)
S -up-> OP : request(n) (수요, 위로)
OP -up-> P : request(n) (위로 전파)

note right of OP
구독하면 request(n)이
구독자 → 소스로 "위로" 전파되고
데이터는 소스 → 구독자로 "아래로" 전달된다
end note
@enduml

PlantUML 발행-연산자-구독 파이프라인 구조

코드로 보면 이 사슬이 그대로 드러납니다.

Flux.range(1, 100)          // ① 소스: 1~100 발행
.filter(n -> n % 2 == 0) // ② 연산자: 짝수만 통과
.map(n -> n * 10) // ③ 연산자: 변환
.subscribe(System.out::println); // ④ 최종 구독자 — 여기서 실제 실행 시작

연결·실행되는 순서는 이렇습니다.

  1. 조립(assembly) — .filter·.map은 각 단계를 감싼 새 Publisher를 만들 뿐, 아직 아무것도 실행되지 않습니다(위에서 아래로 사슬을 엮기만).
  2. 구독(subscribe) — 최종 subscribe가 호출되면, 구독이 아래에서 위로 거슬러 올라가며 각 단계를 이어 붙입니다(onSubscribe).
  3. 요청(request) — 최종 구독자가 request(n)을 위로 올려 보내고, 각 연산자가 소스까지 전파합니다.
  4. 발행(onNext) — 소스가 요청받은 만큼 데이터를 위에서 아래로 내보내고, 각 연산자를 거쳐 변환된 뒤 구독자에 도착합니다. 끝나면 onComplete, 오류면 onError가 같은 경로로 내려갑니다.
즉 "선언(조립) → 구독 → 위로 요청 → 아래로 발행"의 순서라, 구독 전엔 사슬만 있고 실행은 없다는 점이 절차적 코드와 가장 다릅니다.

두 가지 짚어둘 개념​

구독해야 실행된다 (lazy)​

Publisher는 설계도일 뿐, subscribe가 호출돼야 실제로 데이터가 발행됩니다. 그래서 리액티브 코드는 연산자로 파이프라인을 선언해 두고, 마지막에 구독(프레임워크가 대신 하기도 함)해야 동작합니다.

cold vs hot 발행자​

  • cold(차가운) — 구독하는 그 순간에 데이터 발행이 시작되고, 구독자마다 처음부터 새로 받습니다. 예: DB 조회·HTTP 요청 — 구독할 때마다 쿼리가 다시 실행됩니다.
  • hot(뜨거운) — 구독 여부와 상관없이 소스가 이미 값을 계속 만들어 내보내고 있어서, 구독하면 구독한 시점 이후의 값부터 받습니다(그전 값은 못 봅니다). 예: 마우스 이벤트·실시간 시세 — 늦게 구독하면 이미 지나간 값은 놓칩니다.
즉 cold는 "구독마다 처음부터 다시", hot은 "구독 시점 이후만 수신"입니다 — 예를 들어 같은 시세 스트림이라도 cold면 구독자마다 과거부터 재생되고, hot이면 지금부터의 시세만 받습니다.

JDK 9 Flow API — 표준이 표준 라이브러리로​

Reactive Streams의 네 인터페이스는 JDK 9부터 java.util.concurrent.Flow 에 그대로 편입됐습니다(JEP 266).

java.util.concurrent.Flow.Publisher<T>
java.util.concurrent.Flow.Subscriber<T>
java.util.concurrent.Flow.Subscription
java.util.concurrent.Flow.Processor<T, R>

즉 Reactive Streams와 JDK Flow는 이름공간만 다른 같은 규격입니다. 라이브러리 없이도 표준 타입을 참조할 수 있고, org.reactivestreams와 java.util.concurrent.Flow 사이는 어댑터로 변환됩니다.

라이브러리와의 관계​

라이브러리표준 타입 대응비고
Project ReactorMono(01)·Flux(0N)스프링 WebFlux의 기본
RxJava 3Flowable(백프레셔 지원)·Observable(미지원) 등안드로이드에서 널리
Akka StreamsSource·Flow·Sink그래프 기반 스트림
JDK FlowFlow.Publisher 등표준 라이브러리 내장

모두 Reactive Streams를 구현하므로, 한 라이브러리의 Publisher를 다른 라이브러리가 구독할 수 있습니다 — 이 상호운용이 표준을 만든 목적입니다. 단, 편의 연산자(map·flatMap 등)는 표준이 아니라 각 라이브러리가 자체 제공합니다.

스프링에서 언제 필요한가 — 구현 예시​

스프링 앱에서 이 표준(과 Reactor)이 실제로 쓸모 있는 상황과 코드입니다.

① 대용량을 메모리에 다 올리지 않고 스트리밍​

전체를 List로 받으면 힙이 위험합니다. Flux로 한 건씩 내보내며 처리하면 메모리를 상수로 유지합니다. R2DBC 리포지토리는 Flux를 반환합니다.

public interface OrderRepository extends ReactiveCrudRepository<Order, Long> {
Flux<Order> findByStatus(String status); // 결과를 한 번에 다 올리지 않고 스트림으로
}

// 대량 데이터를 스트리밍 응답으로 (전체를 메모리에 적재하지 않음)
@GetMapping(value = "/orders/stream", produces = MediaType.APPLICATION_NDJSON_VALUE)
public Flux<Order> stream() {
return orderRepository.findByStatus("PAID");
}

② 실시간 푸시 — SSE​

Flux는 서버가 값을 생기는 대로 계속 내보내는 SSE(Server-Sent Events)에 자연스럽습니다.

@GetMapping(value = "/prices", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<Price> prices() {
return Flux.interval(Duration.ofSeconds(1)) // 1초마다
.map(tick -> priceService.current()); // 최신 시세를 계속 내보냄
}

③ 여러 논블로킹 호출을 병렬로 합치기​

독립적인 외부 호출을 zip으로 동시에 기다렸다 합칩니다(순차 대기 대신).

public Mono<Dashboard> dashboard(Long userId) {
return Mono.zip(userClient.get(userId), // 세 호출을
orderClient.count(userId), // 동시에 논블로킹으로
pointClient.balance(userId)) // 실행하고
.map(t -> new Dashboard(t.getT1(), t.getT2(), t.getT3()));
}

④ 블로킹 라이브러리를 Reactor로 감싸기​

R2DBC가 없거나 블로킹 SDK만 있을 때는, 그 호출을 별도 스케줄러로 분리해 이벤트 루프를 막지 않게 합니다(순수 리액티브의 이점은 줄지만 혼용은 가능).

Mono.fromCallable(() -> legacyBlockingClient.call())  // 블로킹 호출
.subscribeOn(Schedulers.boundedElastic()); // 이벤트 루프 밖 전용 풀에서 실행

⑤ 표준 타입으로 라이브러리 경계 넘기​

메서드 시그니처를 Reactor 전용 Flux 대신 표준 org.reactivestreams.Publisher로 두면, RxJava·JDK Flow 등 어느 구현이 와도 받을 수 있습니다.

// 표준 타입으로 받으면 Reactor Flux·RxJava Flowable 등 무엇이든 구독 가능
public void consume(org.reactivestreams.Publisher<Event> source) {
Flux.from(source).subscribe(this::handle); // Reactor로 감싸 처리
}
정리하면 스프링에서 Reactive Streams가 필요한 건 대용량 스트리밍·실시간 푸시(SSE)·다중 논블로킹 호출 합성처럼 흐름 제어와 논블로킹이 실익을 주는 경우입니다 — 단순 CRUD엔 과합니다(그럴 땐 MVC나 가상 스레드).

표준이 보장하는 것과 아닌 것​

  • 보장: 논블로킹 백프레셔(request(n)), 네 인터페이스의 규약(신호 순서·취소·오류 전파), 라이브러리 간 상호운용.
  • 표준이 아님: 연산자·스케줄러·실행 모델은 각 라이브러리 몫. map이나 subscribeOn은 Reactor·RxJava가 각자 정의합니다.

정리​

  • Reactive Streams는 비동기 스트림을 논블로킹 백프레셔로 주고받는 최소 표준입니다(인터페이스 4개).
  • 핵심은 구독자가 request(n)으로 감당할 만큼만 당겨 오는 흐름 제어입니다 — push/pull의 절충.
  • Publisher는 구독해야 실행되고(lazy), cold/hot에 따라 동작이 다릅니다.
  • JDK 9 Flow가 이 표준을 그대로 담았고, Reactor·RxJava·Akka가 모두 이를 구현해 서로 연동됩니다. 연산자·스케줄러는 표준 밖이라 라이브러리마다 다릅니다.