리액티브 스트림과 배압(Backpressure) 제어
1. 개요
가. 정의
리액티브 스트림(Reactive Streams)은 비동기·논블로킹 데이터 스트림을 생산자와 소비자 사이에서 처리하되, 소비자가 감당할 수 있는 속도로만 데이터를 흘려보내도록 배압(Backpressure) 신호를 표준화한 명세이자 프로그래밍 모델이다.
리액티브 스트림의 본질은 "밀어내기(push)"와 "끌어오기(pull)"를 하나의 프로토콜로 결합한 데 있다.
전통적인 이벤트 기반 처리는 생산자가 데이터를 소비자에게 일방적으로 밀어내는 push 모델이어서, 소비자가 느리면 버퍼가 넘치거나 메시지가 유실된다.
반대로 소비자가 필요할 때마다 데이터를 요청하는 pull 모델은 안전하지만 왕복 지연이 커져 처리량이 떨어진다.
리액티브 스트림은 소비자가 "지금 n개까지 받을 수 있다"는 수요(demand)를 생산자에게 알리는 request(n) 신호를 도입하여, 평소에는 push의 효율을 유지하면서도 소비자가 포화되면 자동으로 흐름이 제어되는 동적 push-pull 모델을 실현한다.
여기서 배압이란 소비자의 처리 능력이 생산자의 생산 속도보다 낮을 때, 그 부족분을 상류로 되돌려 생산 속도 자체를 늦추는 제어 신호를 뜻한다. 수도관에서 하류가 막히면 상류의 수압이 올라가듯, 소프트웨어 파이프라인에서도 하류의 지연이 상류로 전달되어야 시스템 전체가 붕괴하지 않는다. 배압이 없는 시스템은 부하가 몰릴 때 큐가 무한히 커지다가 메모리 부족으로 프로세스가 종료되거나, 큐를 유한하게 두면 메시지가 조용히 버려지는 두 가지 실패 중 하나에 이른다.
나. 등장 배경과 필요성
리액티브 스트림이 표준으로 정립된 2013~2015년은 스트리밍 데이터와 마이크로서비스가 폭증하던 시기였다.
RxJava, Akka Streams, Project Reactor 등 여러 라이브러리가 각자의 방식으로 비동기 스트림을 다루면서, 서로 다른 라이브러리를 조합할 때 배압 신호가 호환되지 않는 문제가 반복되었다.
이를 해결하기 위해 Netflix, Lightbend, Pivotal 등이 참여하여 Publisher, Subscriber, Subscription, Processor 네 개의 인터페이스로 구성된 최소 명세를 합의했고, 이는 이후 자바 9의 java.util.concurrent.Flow로 표준 라이브러리에 편입되었다.
필요성은 자원 효율과 안정성 두 축에서 설명된다. 첫째, 스레드-당-요청(thread-per-request) 모델은 요청마다 스레드를 점유하므로 동시 접속이 수만 건에 이르면 스레드 스택 메모리와 컨텍스트 스위칭 비용이 폭증한다. 예를 들어 스레드 하나가 1MB 스택을 쓰면 1만 개 스레드는 10GB를 소비하며, I/O 대기 중인 스레드가 대부분이어서 CPU는 놀고 메모리만 고갈된다. 논블로킹 리액티브 모델은 소수의 이벤트 루프 스레드(보통 CPU 코어 수)로 수만 개의 동시 연결을 처리하여 메모리를 크게 절약한다.
둘째, 분산 파이프라인에서 한 구간의 지연이 상류로 전파되지 않으면 장애가 증폭된다. 카프카 컨슈머가 데이터베이스 쓰기 지연으로 느려졌는데도 브로커에서 계속 데이터를 당겨오면 컨슈머의 힙이 넘친다. 배압은 이 상황에서 컨슈머가 처리 가능한 만큼만 폴링하도록 만들어, 지연을 파이프라인 전체가 나누어 흡수하게 한다. 따라서 리액티브 스트림은 단순한 성능 최적화가 아니라 부하 급증 상황에서의 회복탄력성(resilience) 설계의 핵심 수단이다.
다. 특징
리액티브 스트림의 특징은 비동기성, 논블로킹, 배압 기반 흐름 제어, 조합 가능성(composability)으로 요약된다. 연산자(operator)를 함수형으로 연결해 파이프라인을 선언적으로 기술하며, 각 연산자는 상류의 수요와 하류의 수요를 조정하는 역할을 겸한다. 아래 표는 전통적 모델과의 차이를 정리한 것이나, 각 항목의 "왜"는 이어지는 절에서 산문으로 설명한다.
| 구분 | 블로킹 동기 모델 | 순수 push 이벤트 모델 | 리액티브 스트림 |
|---|---|---|---|
| 흐름 제어 | 호출 스택이 자연스러운 배압 | 없음(유실·OOM 위험) | request(n) 기반 명시적 배압 |
| 스레드 사용 | 요청당 스레드 점유 | 이벤트 루프 | 이벤트 루프 + 논블로킹 |
| 지연 전파 | 동기 호출로 전파 | 전파 안 됨 | 수요 신호로 전파 |
| 조합성 | 낮음 | 중간 | 높음(연산자 체인) |
2. 리액티브 스트림 명세와 구성요소
가. 네 개의 핵심 인터페이스
리액티브 스트림 명세는 Publisher, Subscriber, Subscription, Processor라는 네 개의 인터페이스로 상호작용 규약을 정의한다.
Publisher는 데이터를 생산하는 원천이며, subscribe() 호출로 소비자를 연결한다.
Subscriber는 데이터를 소비하며 onSubscribe, onNext, onError, onComplete 네 개의 콜백을 갖는다.
Subscription은 생산자와 소비자 사이의 한 번의 연결을 나타내며, 소비자가 수요를 요청하는 request(n)과 연결을 끊는 cancel()을 제공한다.
Processor는 Publisher이면서 동시에 Subscriber인 중간 단계로, 연산자 체인의 각 노드에 해당한다.
전체 구조를 개념도로 표현하면 다음과 같다.
flowchart LR
P["Publisher(생산자)"] -->|onSubscribe| SUB["Subscription(구독)"]
SUB -->|request n| P
P -->|onNext 데이터| PR["Processor(연산자 체인)"]
PR -->|onNext 변환| S["Subscriber(소비자)"]
S -->|request n 수요| PR
PR -->|request n 재조정| P
S -.->|onError / onComplete| END["종료 신호"]
이 구조의 핵심은 데이터가 왼쪽에서 오른쪽으로 흐르는 동안, 수요 신호(request(n))는 오른쪽에서 왼쪽으로 거슬러 흐른다는 점이다.
소비자가 자신이 감당할 수 있는 개수만큼만 요청하면 그 신호가 연산자를 거쳐 생산자까지 전달되고, 생산자는 요청받은 개수 이하로만 onNext를 방출한다.
이 규약 덕분에 생산자는 소비자의 처리 능력을 몰라도 절대 요청량을 초과해 데이터를 밀어내지 않는다.
나. 신호 규약과 불변식
명세는 신호의 순서와 개수에 관한 엄격한 불변식을 요구한다.
onSubscribe는 정확히 한 번 먼저 호출되어야 하고, onNext는 요청받은 수를 초과할 수 없으며, onError나 onComplete 이후에는 어떤 신호도 발생해서는 안 된다.
이러한 규약은 서로 다른 라이브러리를 안전하게 연결하기 위한 최소 계약이며, TCK(Technology Compatibility Kit)라는 표준 테스트 스위트로 준수 여부를 검증한다.
규약을 지키지 않는 구현은 조합 시 경합 조건이나 자원 누수를 일으키므로, 직접 Publisher를 구현하기보다 검증된 라이브러리의 팩토리 메서드를 사용하는 것이 실무 권장 사항이다.
다. 콜드 스트림과 핫 스트림
스트림은 구독 시점에 따라 콜드(cold)와 핫(hot)으로 나뉜다. 콜드 스트림은 구독자가 붙을 때마다 데이터 생산이 처음부터 새로 시작되며, HTTP 요청·파일 읽기·데이터베이스 조회처럼 요청 단위로 완결되는 작업에 적합하다. 핫 스트림은 구독자 유무와 무관하게 데이터가 계속 흐르며, 센서 이벤트·주식 시세·사용자 클릭처럼 실시간 방송 성격의 원천에 해당한다. 핫 스트림은 배압을 그대로 적용하기 어려운데, 생산 원천이 소비자의 수요를 기다려 주지 않기 때문이다. 이 경우 뒤에서 설명할 완충·표본추출·최신값 유지 같은 유실 허용 전략을 함께 설계해야 한다.
3. 배압 제어 메커니즘과 동작 절차
가. 수요 기반 흐름 제어 절차
배압의 표준 메커니즘은 소비자가 처리 여력만큼 request(n)을 호출하고, 생산자는 그 범위 안에서만 방출하는 왕복 대화다.
소비자는 한 번에 큰 수요를 요청한 뒤 소진될 때마다 보충하거나(예: 256개 요청 후 절반 소진 시 128개 추가), 정확히 하나씩 요청하며 처리 속도에 맞추는 방식을 선택할 수 있다.
아래 시퀀스 다이어그램은 소비자가 느려질 때 수요가 어떻게 조절되는지를 보여 준다.
sequenceDiagram
participant S as Subscriber(소비자)
participant O as Operator(버퍼)
participant P as Publisher(생산자)
S->>O: request(256)
O->>P: request(256)
P-->>O: onNext x256
O-->>S: onNext x256
Note over S: 처리 지연 발생
S->>O: request(64) 만 보충
O->>P: 버퍼 여유만큼만 request
Note over P: 방출 속도 자동 감소
S->>O: cancel() 또는 완료
이 절차에서 주목할 점은 소비자가 지연되면 보충 요청이 늦어지고, 그 결과 생산자의 방출도 자연히 느려진다는 것이다.
어떤 명시적 "속도 제한 명령"도 없이, 오직 수요 신호의 지연만으로 전체 파이프라인의 속도가 소비자에 맞추어진다.
실무에서 Project Reactor의 Flux는 기본적으로 256개 크기의 프리페치 버퍼를 두고 75%가 소진되면 다음 배치를 요청하는 식으로 이 대화를 자동화한다.
나. 배압이 없을 때의 붕괴 시나리오
배압이 제대로 전파되지 않으면 시스템은 두 방향 중 하나로 무너진다.
무한 버퍼를 쓰면 부하 급증 시 큐가 계속 커져 힙이 고갈되고 OutOfMemoryError로 프로세스가 종료된다.
유한 버퍼를 쓰면서 넘칠 때 예외를 던지면 MissingBackpressureException 같은 오류로 요청이 실패한다.
2020년대 초 여러 스트리밍 서비스의 장애 사후분석에서, 소비 지연이 상류로 전파되지 못해 브로커–컨슈머 사이 큐가 폭증한 사례가 반복적으로 보고되었다.
따라서 배압 설계는 "정상 부하에서 잘 도는가"가 아니라 "최대 부하가 소비 능력을 초과할 때 어떻게 우아하게 저하되는가"를 기준으로 검증해야 한다.
다. 유실 허용 전략과 완충
소비자를 무한정 기다릴 수 없는 핫 스트림에서는 데이터 일부를 의도적으로 버리거나 요약하는 오버플로 전략이 필요하다. 대표적으로 버퍼(buffer)는 일정 크기까지 쌓았다가 배치로 넘기고, 드롭(drop)은 소비자가 바쁠 때 새 데이터를 버리며, 최신값(latest)은 가장 최근 값만 유지하고 이전 값을 덮어쓴다. 표본추출(sample)이나 윈도우(window)는 시간·개수 기준으로 데이터를 묶어 하류 부하를 낮춘다. 아래 표는 전략별 특성을 비교한다.
| 전략 | 동작 | 데이터 유실 | 적합 상황 |
|---|---|---|---|
| buffer | 유한 큐에 적재 후 배치 | 넘치면 오류/차단 | 짧은 순간 부하, 유실 불가 |
| drop | 초과분 폐기 | 있음 | 최신성보다 순서 중요할 때 |
| latest | 최신 1건만 유지 | 있음 | 대시보드·시세 표시 |
| error | 즉시 실패 신호 | - | 빠른 실패가 나은 배치 |
전략 선택은 업무 요구사항에서 출발한다. 결제 이벤트처럼 한 건도 유실되면 안 되는 데이터는 buffer + 영속 큐 + 재처리로 무손실을 보장해야 하고, 실시간 온도 게이지처럼 최신값만 의미 있는 데이터는 latest로 오래된 값을 버리는 편이 자원 효율과 사용자 경험 모두에 유리하다.
4. 비교와 실무 적용 사례
가. 명령형 모델과의 트레이드오프
리액티브 모델은 높은 처리량과 자원 효율을 주지만, 디버깅과 학습 곡선이라는 대가를 수반한다. 비동기 파이프라인에서는 스택 트레이스가 실제 실행 경로를 담지 못해 오류 원인 추적이 어렵고, 블로킹 호출 하나가 이벤트 루프를 막으면 전체 처리량이 급락한다. 반면 스레드-당-요청 모델은 코드가 직관적이고 디버깅이 쉬우며, 최근에는 자바의 가상 스레드(Virtual Thread, Project Loom)가 등장하여 블로킹 코드로도 높은 동시성을 얻을 수 있게 되었다. 따라서 "무조건 리액티브"가 아니라, I/O 바운드이면서 스트리밍·배압 제어가 본질적으로 필요한 구간에 선택적으로 적용하는 판단이 중요하다.
나. 웹 서비스 사례 — Spring WebFlux
한 API 게이트웨이가 초당 수만 건의 요청을 소수의 이벤트 루프로 처리해야 하는 상황을 가정한다.
Spring MVC(스레드-당-요청)로는 톰캣 스레드풀 200개가 외부 API 응답을 기다리며 모두 점유되어, 실제 CPU 사용률이 20%인데도 요청이 큐에서 대기하다 타임아웃되는 현상이 발생한다.
Spring WebFlux(Reactor Netty 기반)로 전환하면 CPU 코어 수만큼의 이벤트 루프가 논블로킹으로 수만 연결을 다루고, 하류 API가 느려지면 request(n)이 상류 클라이언트로 전파되어 자연스럽게 유입을 늦춘다.
실측 사례에서 동일 하드웨어로 동시 연결 수용량이 수 배 증가하고 P99 지연이 안정화되는 결과가 보고된다.
다. 데이터 파이프라인 사례 — Kafka 컨슈머
카프카 기반 이벤트 파이프라인에서 컨슈머가 다운스트림 데이터베이스에 초당 5,000건을 쓸 수 있는데 토픽에는 초당 20,000건이 유입된다고 하자.
배압 없이 무제한 폴링하면 컨슈머 힙에 미처리 레코드가 쌓여 GC 폭주와 OOM으로 이어진다.
Reactor Kafka나 Akka Streams의 카프카 커넥터는 다운스트림 싱크의 수요에 맞추어 poll 호출과 오프셋 커밋 속도를 조절하고, 파티션별 프리페치를 제한하여 컨슈머가 처리 가능한 만큼만 당겨온다.
그 결과 순간 유입이 소비 능력을 초과해도 브로커에 데이터가 안전하게 보관된 채 컨슈머는 일정한 속도로 따라잡으며, 파이프라인 전체가 유입 폭주를 완충한다.
5. 심화 — 최신 동향과 유사 기술 연계
리액티브 스트림 표준은 자바 9의 Flow API로 편입되어 언어 표준의 일부가 되었고, R2DBC(리액티브 관계형 데이터베이스 접근)와 리액티브 gRPC처럼 스택 전 계층으로 배압이 확장되고 있다.
데이터베이스 드라이버가 배압을 지원하지 않으면 파이프라인의 한 지점이 블로킹되어 전체 이점이 무너지므로, "종단 간(end-to-end) 논블로킹"이 실무의 핵심 관심사가 되었다.
한편 자바 21에서 정식화된 가상 스레드는 리액티브 모델의 대안으로 부상했다. 가상 스레드는 블로킹처럼 보이는 코드를 논블로킹으로 실행하여 리액티브의 자원 효율을 명령형 코드로 얻게 하지만, 가상 스레드 자체는 배압을 제공하지 않는다는 점이 중요하다. 즉 동시성 문제(스레드 부족)는 가상 스레드가 풀어 주지만, 흐름 제어 문제(생산-소비 속도 불일치)는 여전히 리액티브 스트림이나 세마포어·제한 큐 같은 별도 메커니즘이 필요하다. 따라서 두 기술은 경쟁이 아니라, "가상 스레드로 단순한 I/O를, 리액티브 스트림으로 배압이 필요한 스트리밍을" 다루는 상보 관계로 정착하는 추세다.
배압 개념은 리액티브 스트림에 국한되지 않는다. TCP의 수신 윈도우, gRPC의 HTTP/2 플로우 컨트롤, 카프카의 컨슈머 랙 기반 조절, 서킷 브레이커의 부하 차단은 모두 넓은 의미의 배압·흐름 제어 계열에 속한다. 기술사 관점에서 이들을 하나의 "흐름 제어" 계보로 묶어 설명하면, 개별 기술 암기가 아니라 원리 중심의 답안을 구성할 수 있다.
6. 고려사항 및 시사점
첫째, 적용 범위를 냉정하게 판단해야 한다. 리액티브 스트림은 I/O 바운드·고동시성·스트리밍 구간에서 큰 이점을 주지만, CPU 바운드 배치나 단순 CRUD에는 복잡성만 키운다. 가상 스레드가 성숙한 지금은 "블로킹 코드가 문제인가, 흐름 제어가 문제인가"를 먼저 구분하고, 후자일 때만 리액티브를 선택하는 의사결정 기준이 필요하다.
둘째, 종단 간 논블로킹을 설계 원칙으로 삼아야 한다. 파이프라인 어느 한 지점에서라도 블로킹 호출(JDBC, 동기 파일 I/O 등)이 이벤트 루프에서 실행되면 소수의 이벤트 루프가 막혀 전체 처리량이 붕괴한다. 불가피한 블로킹 작업은 반드시 별도의 경계 스레드풀(bounded elastic scheduler)로 격리하고, 데이터 접근 계층은 R2DBC처럼 논블로킹 드라이버로 통일하는 방향이 바람직하다.
셋째, 오버플로 전략을 업무 위험도에 맞추어 명시적으로 선택해야 한다. 데이터 유실이 허용되지 않는 도메인은 무손실 버퍼와 영속 큐·재처리를 결합하고, 최신성만 중요한 도메인은 유실 허용 전략으로 자원을 아끼는 식으로, 오버플로 정책 자체를 아키텍처 결정 기록(ADR)으로 남겨 운영자가 이해할 수 있게 해야 한다.
넷째, 관측성과 테스트를 함께 설계해야 한다. 배압은 정상 부하에서는 보이지 않다가 최대 부하에서 드러나므로, 프리페치 크기·버퍼 점유율·요청 대비 방출량·컨슈머 랙 같은 지표를 상시 수집하고, 부하 테스트·장애 주입으로 "소비 능력을 초과할 때의 저하 곡선"을 사전에 확인해야 한다. 비동기 스택 트레이스를 보완하기 위한 체크포인트·컨텍스트 전파도 운영 준비의 일부다.
다섯째, 조직 역량과 표준화를 병행해야 한다.
리액티브 코드는 학습 곡선이 가파르고 잘못 쓰면 오히려 성능이 나빠지므로, 검증된 라이브러리의 관용구를 팀 표준으로 정하고, 직접 Publisher를 구현하기보다 팩토리·연산자를 조합하는 규칙을 코드 리뷰에 포함하는 것이 안전하다.
참고자료
- Reactive Streams Specification: https://www.reactive-streams.org/
- Project Reactor Reference — Handling Backpressure: https://projectreactor.io/docs/core/release/reference/
- Spring Framework, Web on Reactive Stack (WebFlux): https://docs.spring.io/spring-framework/reference/web/webflux.html
- OpenJDK, JEP 444: Virtual Threads: https://openjdk.org/jeps/444
한 줄 요약: 리액티브 스트림은
request(n)기반 배압 신호로 소비자의 처리 능력에 맞추어 생산 속도를 되돌려 제어함으로써, 논블로킹의 자원 효율과 부하 급증 시의 회복탄력성을 동시에 달성하는 비동기 스트림 처리 표준이다.