https://dmoritle.tistory.com/282
[Spring WebFlux] 리액티브 프로그래밍이란 무엇인가
1. 들어가며Spring으로 백엔드를 개발하다 보면 대부분 Spring MVC로 시작합니다. 그러다 어느 시점엔가 Spring WebFlux라는 이름을 접하게 됩니다. 계기는 보통 하나로 수렴합니다. MVC 기반 서비스가 트
dmoritle.tistory.com
1. 들어가며
배압(Backpressure)이라는 개념이 없던 시절, 비동기 스트림 처리 라이브러리들은 저마다 다른 방식으로 동작했습니다. 넷플릭스가 만든 초기 RxJava, 라이트벤드의 Akka Streams 등이 각자의 방식으로 비동기 데이터 흐름을 다뤘고, 서로 다른 라이브러리로 만든 컴포넌트끼리는 조합조차 되지 않았습니다.
더 근본적인 문제는 따로 있었습니다. Publisher가 데이터를 계속 밀어내는데 Subscriber가 처리 속도를 못 따라가면, 처리되지 못한 데이터가 메모리에 쌓입니다. 큐가 무한정 자라나고, 결국 OutOfMemoryError로 서비스가 죽습니다. 빠른 발행자와 느린 구독자의 조합이라면 언제든 일어날 수 있는 사고였습니다.
이 두 문제, 즉 상호운용성 부재와 배압 부재로 인한 메모리 폭발을 해결하기 위해 2013년 넷플릭스, 피보탈(Pivotal, 現 Spring), 라이트벤드(Lightbend, Akka)가 시작한 이니셔티브가 Reactive Streams입니다. 1편에서 이름만 스쳐 지나갔던 Publisher, Subscriber, Subscription, Processor 네 인터페이스가 정확히 무엇을 약속하는지, 이번 글에서 하나씩 뜯어보겠습니다.
2. 배압(Backpressure)이란
배압은 한마디로 "받는 쪽이 보내는 쪽의 속도를 제어하는 능력"입니다.
Reactive Streams 이전의 비동기 처리는 대체로 Push 기반이었습니다. 데이터가 준비되는 대로 Publisher가 Subscriber에게 밀어 넣는 방식입니다. 문제는 이 흐름에 브레이크가 없다는 점입니다. Subscriber가 초당 100개를 처리할 수 있는데 Publisher가 초당 1,000개를 밀어 넣으면, 처리되지 못한 900개가 버퍼에 쌓입니다. 이 상태가 지속되면 버퍼는 무한정 자라나고 메모리는 고갈됩니다.
Reactive Streams는 이 흐름을 Pull 기반으로 뒤집었습니다. Subscriber가 request(n)을 호출해서 "지금 n개까지는 처리할 수 있다"고 명시적으로 선언하지 않으면, Publisher는 데이터를 보내지 않습니다. 즉 데이터의 흐름은 Publisher에서 Subscriber로 가지만, 그 흐름의 속도를 결정하는 주도권은 Subscriber가 쥐고 있는 구조입니다.
이 메커니즘 덕분에 Subscriber는 자기 처리 능력만큼만 요청하면 되고, Publisher는 요청받은 양을 절대 넘어서 보낼 수 없습니다. 배압은 이 "받는 쪽이 미리 선언한 한도 안에서만 흐른다"는 규칙 전체를 가리키는 말입니다.
3. Sequence — 정규식 한 줄로 보는 계약의 전체 그림
Reactive Streams 스펙은 Publisher가 Subscriber에게 보내는 일련의 시그널 전체를 Sequence라고 부릅니다. 그리고 이 Sequence가 지켜야 할 규칙을 정규식 한 줄로 정의합니다.
onSubscribe onNext* (onError | onComplete)?
풀어서 읽으면 이렇습니다. onSubscribe는 반드시 정확히 한 번 먼저 호출됩니다. 그 뒤로 onNext가 0번 이상 반복될 수 있습니다. 그리고 마지막에 onError 또는 onComplete 중 하나가, 있다면 딱 한 번만 호출됩니다. 구독이 취소되지 않는 한 이 순서를 벗어날 수 없습니다.
이 한 줄이 사실상 Reactive Streams 계약 전체의 뼈대입니다. 다음 섹션에서 볼 인터페이스별 세부 규칙들은 결국 이 정규식이 허용하는 순서를 더 정교하게 다듬은 것들입니다.
이 개념은 Reactor의 타입 이름과도 직접 연결됩니다. Flux는 이 Sequence가 0개에서 N개의 데이터를 실어 나를 수 있다는 뜻이고, Mono는 0개 또는 1개로 제한된 Sequence라는 뜻입니다. Mono<HttpResponse>가 Flux<HttpResponse>보다 HTTP 응답을 표현하기에 더 적절한 이유도 여기서 나옵니다. HTTP 요청 하나는 응답을 정확히 하나만 만들어내는, 카디널리티가 0~1인 Sequence이기 때문입니다.
4. 4대 인터페이스와 계약
4-1. Publisher.subscribe()는 왜 낯설게 느껴지는가
public interface Publisher<T> {
void subscribe(Subscriber<? super T> s);
}
Publisher.subscribe(Subscriber)라는 형태 자체는 옵저버(Observer) 패턴의 Subject.addObserver(observer)와 구조적으로 거의 같습니다. 그래서 옵저버 패턴에 익숙하다면 이 시그니처 자체는 낯설지 않을 것입니다.
진짜 낯선 지점은 그다음입니다. 일반적인 옵저버 패턴의 구독은, 구독을 신청하면 그걸로 끝이고 이후에는 발행자가 알아서 데이터를 밀어 넣습니다. 반면 Reactive Streams는 subscribe() 호출 한 번으로 끝나지 않습니다.

"구독 신청"과 "실제 데이터 흐름 시작" 사이에 또 하나의 핸드셰이크가 끼어 있고, 그 사이에 Subscriber가 능동적으로 흐름을 열어야 하는 지점(request)이 있다는 것이 일반적인 pub/sub 구조와의 결정적인 차이입니다.
4-2. 카프카의 subscribe()와는 무엇이 다른가
카프카는 consumer.subscribe(topics) 형태입니다. subscribe()가 Consumer(구독자) 쪽 객체에 붙어 있고, Consumer가 이를 호출해서 "이 토픽을 구독하겠다"고 선언합니다. 반면 Reactive Streams는 Publisher.subscribe(Subscriber)로, subscribe()가 Publisher(발행자) 쪽 객체에 붙어 있습니다. 실제 코드로 보면 다음과 같습니다.
Flux<Integer> flux = Flux.range(1, 3); // flux 자체가 Publisher
flux.subscribe(mySubscriber); // Publisher(flux)의 메서드를 호출
"구독한다"는 행위의 주체는 의미상 Subscriber인데, 정작 메서드는 Publisher에 붙어 있는 비대칭이 여기서 발생합니다. 이는 브로커 같은 중개자 없이 두 객체가 프로세스 안에서 직접 연결되는 구조이기 때문입니다. 카프카는 브로커가 토픽 이름으로 연결을 중개해주지만, Reactive Streams는 호출하는 쪽 코드가 Publisher 인스턴스와 Subscriber 인스턴스를 둘 다 직접 들고 있어야 합니다.
그 상태에서는 "지금 들고 있는 Publisher 객체"에 메서드를 거는 편이 자연스러운 선택입니다. Iterable.iterator()도 같은 구조입니다. list.iterator()라고 호출하지 iterator.attachTo(list)라고 하지 않는 것처럼, 순회 대상(생산자) 쪽에 메서드를 두고 호출자가 그 생산자 인스턴스를 붙잡고 부르는 방식입니다. 즉 "누가 주도권을 갖고 구독을 요청하는가"(의미상 Subscriber 쪽 의도)와 "메서드가 어느 객체에 붙어 있는가"(구문상 Publisher 쪽)가 서로 다른 축이라는 점이 카프카와의 차이입니다.
4-3. 각 메서드와 계약 규칙
Reactive Streams는 단순히 인터페이스만 정의하지 않습니다. 각 메서드가 지켜야 할 규칙(Rule)을 스펙 문서에 명시적으로 못박아 두고 있습니다. 대표적인 규칙을 추려보면 다음과 같습니다.
| 메서드 | 역할 | 대표 계약 규칙 |
onSubscribe |
구독 성립을 알림 | Subscriber 하나당 정확히 한 번만 호출됨 |
onNext |
데이터를 전달 | Publisher가 보내는 onNext 총 횟수는 그 Subscriber가 요청한 총량을 절대 넘을 수 없음 |
onComplete / onError |
스트림을 종료 | 둘 중 최대 하나만, 그리고 정상적으로 도착하면 그 이후 어떤 시그널도 오지 않음. onComplete는 request 호출 없이도 도착할 수 있음 |
request(n) |
배압 요청 | n이 0 이하면 IllegalArgumentException과 함께 onError가 발생해야 함 |
cancel() |
구독 취소 | best-effort 방식이라 즉시 반영을 보장하지 않으며, 취소 후에도 이미 요청이 걸려 있던 데이터라면 onNext가 더 도착할 수 있음 |
이 표에서 특히 눈여겨볼 점은 cancel()의 즉시성입니다. 취소는 "최대한 빨리 멈춰달라"는 요청이지 즉각적인 차단이 아니므로, cancel() 직후에도 짧게 몇 개의 onNext가 더 들어올 수 있다는 걸 감안하고 후속 처리를 짜야 합니다.
4-4. Subscription이 배압의 핵심 축인 이유
public interface Subscription {
void request(long n);
void cancel();
}
Subscription은 Publisher와 Subscriber 사이에서 이 둘을 연결하는 매개체입니다. Subscriber는 Subscription을 통해서만 request와 cancel을 호출할 수 있고, 이 Subscription은 정확히 그 Publisher-Subscriber 쌍 하나만을 위해 존재합니다.
배압이라는 개념이 실제로 작동하는 지점이 바로 여기입니다. Publisher와 Subscriber는 "무엇을 주고받을지"를 정의하는 역할이고, Subscription은 "얼마나 주고받을지"를 조절하는 밸브 역할을 합니다. 2번 섹션에서 말한 "받는 쪽이 보내는 쪽의 속도를 제어한다"는 문장을 코드 레벨로 옮기면, 결국 Subscription.request(n) 호출 하나로 요약됩니다.
Processor<T, R>는 이 네 인터페이스 중 마지막입니다. Subscriber<T>이면서 동시에 Publisher<R>인 타입으로, 데이터를 받아서(구독자 역할) 변환한 뒤 다시 내보내는(발행자 역할) 중간 처리자입니다. Reactor의 map, filter 같은 연산자들이 내부적으로 담당하는 역할이 바로 이것입니다.
5. Signal, Demand, Emit — 부품에 이름 붙이기
지금까지 쓴 용어들을 정리하고 넘어가겠습니다.
Signal(시그널)은 Publisher와 Subscriber 사이를 오가는 onSubscribe, onNext, onError, onComplete, request(n), cancel 여섯 개 메서드 호출을 통칭하는 말입니다. 방향은 두 갈래입니다. onSubscribe/onNext/onError/onComplete는 Publisher가 Subscriber에게 보내는 알림성 시그널이고, request(n)/cancel은 반대로 Subscriber가 Publisher에게 보내는 요청성 시그널입니다. 3번에서 본 정규식 onSubscribe onNext* (onError|onComplete)?은 이 중 앞쪽, 즉 알림성 시그널이 어떤 순서와 개수로 와야 하는가에 대한 규칙입니다.
Demand(수요)는 request(n)으로 쌓이는, 아직 충족되지 않은 요청량을 말합니다. request(2)를 두 번 호출하면 Demand는 4가 됩니다. 누적되는 값이라는 점이 중요합니다. Publisher는 이 Demand를 초과해서 데이터를 내보낼 수 없습니다.
Emit(방출)은 Publisher가 실제로 onNext(T)를 호출해서 데이터를 내보내는 행위를 가리키는 관용적인 표현입니다. Demand가 있어야만 Emit이 가능하다는 관계가 핵심입니다. 이 한 문장, "Demand 없이는 Emit도 없다"가 배압 메커니즘을 가장 짧게 요약한 문장입니다.
6. Upstream/Downstream과 흐름의 방향성
연산자를 체이닝하다 보면 "업스트림", "다운스트림"이라는 말을 자주 마주칩니다. 이건 스트림 안에서의 상대적인 위치를 가리키는 말입니다.
Source → map() → filter() → Subscriber
map() 연산자를 기준으로 보면, 데이터가 흘러오는 쪽인 Source가 업스트림이고, 데이터가 흘러가는 쪽인 filter()가 다운스트림입니다. 즉 업스트림/다운스트림은 고정된 이름이 아니라, 어느 연산자를 기준점으로 삼느냐에 따라 상대적으로 정해집니다.
여기서 5번의 구분을 다시 가져와 보면, 알림성 시그널과 요청성 시그널은 서로 반대 방향으로 흐릅니다.

Subscriber(가장 다운스트림에 위치)가 request(n)으로 Demand를 요청하면, 그 요청성 시그널이 체인을 거슬러 Publisher 방향(업스트림)으로 전달됩니다. 반대로 실제 데이터를 담은 onNext 같은 알림성 시그널은 Publisher에서 시작해 체인을 따라 Subscriber 방향(다운스트림)으로 흘러갑니다. 이 반대 방향 흐름이야말로 배압 메커니즘의 본질입니다. 데이터는 밀려 내려오지만, 그 데이터의 양은 밑에서 끌어올리는 만큼만 허용됩니다.
7. 동작 흐름 종합 다이어그램 및 다음 편 예고
지금까지 짚은 개념들을 하나의 흐름으로 이어보면 다음과 같습니다.

- Subscriber가
Publisher.subscribe()를 호출해 구독 신청 - Publisher가
Subscriber.onSubscribe(Subscription)를 호출해 연결 수립 — 이 시점까지는 아직 데이터가 오가지 않음 - Subscriber가
Subscription.request(n)을 호출해 Demand를 쌓음 - Publisher가 Demand 한도 안에서
onNext(T)를 호출해 데이터를 Emit - 더 보낼 데이터가 없으면
onComplete, 실패하면onError로 Sequence를 종료
3편에서는 이 스펙이 실제로 어떻게 구현되어 있는지 살펴봅니다.
https://dmoritle.tistory.com/284
[Spring WebFlux] Reactive Streams, 실무에서는 - 구현체와 오해
1. 들어가며Spring 프로젝트에서 결제 SDK나 메시징 클라이언트 같은 외부 라이브러리를 붙이다 보면, 가끔 낯선 타입을 마주칩니다. 내 서비스 레이어는 전부 Mono와 Flux로 짜여 있는데, SDK가 반환
dmoritle.tistory.com
"Data Source"와 "Operator"라는, 스펙에는 없지만 실무에서 자주 쓰는 용어들을 정리하고, Project Reactor와 RxJava, Akka Streams, java.util.concurrent.Flow가 이 스펙을 각각 어떻게 구현했는지 비교해보겠습니다.
'Spring Framework > Spring WebFlux' 카테고리의 다른 글
| [Spring WebFlux] 5편: Reactor의 배압 제어 - request(n)에서 전략까지 (0) | 2026.09.18 |
|---|---|
| [Spring WebFlux] 4편: Cold/Warm/Hot Sequence (1) | 2026.09.15 |
| [Spring WebFlux] 3편: Reactive Streams, 실무에서는 - 구현체와 오해 (0) | 2026.09.14 |
| [Spring WebFlux] 1편: 리액티브 프로그래밍이란 무엇인가 (0) | 2026.09.14 |