목차
RxJava 배압(Backpressure)과 리소스 관리
생산자가 1초에 100건을 만들고 소비자가 1초에 1건만 처리한다면 1초 뒤 대기 데이터가 약 99건 생긴다. 이런 상태가 계속되면 지연이 늘거나 메모리가 부족해질 수 있다. 배압(backpressure)은 소비자가 감당할 수 있는 수량을 request(n)으로 알려 생산·전달 속도를 조율하는 계약이다. RxJava 2에서는 Flowable이 이 계약을 제공하고, Observable은 같은 배압 계약을 제공하지 않는다.

그림은 소비자가 요청한 만큼 생산자가 보내는 흐름을 보여 준다. 그러나 모든 생산자를 멈출 수 있는 것은 아니다. Flowable.range처럼 요청량에 맞춰 값을 만드는 소스는 생산을 늦출 수 있지만, 시계에 맞춰 발생하는 interval이나 외부 센서 이벤트는 시간이 흐르면 새 값이 생긴다. 이런 소스에는 초과 데이터를 어떻게 버퍼링·버림·오류 처리할지 정책이 필요하다.
request(n)을 직접 읽어 보기
import io.reactivex.Flowable;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;
public class DemandExample {
public static void main(String[] args) {
Flowable.range(1, 5).subscribe(new Subscriber<Integer>() {
private Subscription subscription;
private int received;
@Override public void onSubscribe(Subscription s) {
subscription = s;
s.request(2);
}
@Override public void onNext(Integer value) {
System.out.println(value);
received++;
if (received % 2 == 0) subscription.request(2);
}
@Override public void onError(Throwable error) { error.printStackTrace(); }
@Override public void onComplete() { System.out.println("완료"); }
});
// 1, 2, 3, 4, 5, 완료
}
}
구독자는 처음 2건을 요청하고, 2건을 받을 때마다 다시 2건을 요청한다. 다섯 번째 값을 전달한 뒤 원본이 완료한다. request(Long.MAX_VALUE)는 사실상 제한 없이 받겠다는 뜻이라 그 구독자가 건별로 속도를 조절하지 않는다. 원문의 interval 예제에서 오류가 나오는 것은 빠른 시간 기반 생산과 observeOn의 내부 큐 경계 때문이다. 정확히 128에서 오류가 난다는 출력은 환경과 연산자 설정에 따라 달라진다.
멈출 수 없는 생산자에서의 정책
Flowable.create로 외부 콜백을 감쌀 때 BackpressureStrategy를 고른다. 이 전략은 모든 Flowable의 전역 설정이 아니다. 다른 연산자 경계에는 onBackpressureBuffer, onBackpressureDrop, onBackpressureLatest 같은 별도 정책이 있다.
| 정책 | 초과 데이터 | 적합한 예 | 위험 |
|---|---|---|---|
BUFFER | 대기 큐에 보관 | 모든 이벤트가 필요한 짧은 폭주 | 무제한이면 메모리 증가 |
DROP | 새로운 값을 버림 | 일부 표본만 있어도 되는 관측 | 어떤 값이 빠졌는지 기록 필요 |
LATEST | 최신 값으로 교체 | 현재 상태 화면 | 중간 상태를 잃음 |
ERROR | 감당 못 하면 실패 | 누락을 허용하지 않는 경계 | 재시도·오류 처리가 필요 |
MISSING | create에서 별도 처리 없음 | 뒤에서 정책을 명시할 때 | 정책을 빼면 초과 시 실패 가능 |
원문의 표에는 NONE이 있었지만 RxJava 2 BackpressureStrategy의 해당 상수 이름은 MISSING이다. BUFFER가 모든 값을 보존한다고 해도 메모리 한도가 없으면 안정성을 보장하지 않는다. 주문 기록처럼 누락이 불가능한 데이터는 DROP이나 LATEST를 택할 수 없고, 생산자를 제한하거나 내구성 있는 큐에 저장해야 한다. 현재 온도 표시처럼 최신 값이 핵심인 경우에는 LATEST가 지연을 줄일 수 있다.
종료와 모니터링
배압은 구독의 종료를 대신하지 않는다. 무한 소스는 화면·요청의 수명이 끝날 때 Subscription.cancel() 또는 Disposable.dispose()로 연결을 끊는다. 버퍼 크기, 처리 지연, 드롭된 값 수, MissingBackpressureException 발생 수를 함께 관찰한다. 소비자를 무조건 빠르게 만들기 어려운 상황이라면 상류의 생성률을 낮추거나 작업을 나누는 것이 근본적인 해결일 수 있다.
참고: Reactive Streams 사양, RxJava 2 BackpressureStrategy API, Flowable API
