본문으로 건너뛰기
홈
기술
기술 전체
프로그래밍68
컴퓨터 과학63
AI48
웹 개발36
인프라33
데이터31
소프트웨어 공학18
소개
← 목록으로프로그래밍 › 비동기 › 리액티브 프로그래밍

[RxJava] 배압(Backpressure)과 리소스 관리

목차

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감당 못 하면 실패누락을 허용하지 않는 경계재시도·오류 처리가 필요
MISSINGcreate에서 별도 처리 없음뒤에서 정책을 명시할 때정책을 빼면 초과 시 실패 가능

원문의 표에는 NONE이 있었지만 RxJava 2 BackpressureStrategy의 해당 상수 이름은 MISSING이다. BUFFER가 모든 값을 보존한다고 해도 메모리 한도가 없으면 안정성을 보장하지 않는다. 주문 기록처럼 누락이 불가능한 데이터는 DROP이나 LATEST를 택할 수 없고, 생산자를 제한하거나 내구성 있는 큐에 저장해야 한다. 현재 온도 표시처럼 최신 값이 핵심인 경우에는 LATEST가 지연을 줄일 수 있다.

종료와 모니터링

배압은 구독의 종료를 대신하지 않는다. 무한 소스는 화면·요청의 수명이 끝날 때 Subscription.cancel() 또는 Disposable.dispose()로 연결을 끊는다. 버퍼 크기, 처리 지연, 드롭된 값 수, MissingBackpressureException 발생 수를 함께 관찰한다. 소비자를 무조건 빠르게 만들기 어려운 상황이라면 상류의 생성률을 낮추거나 작업을 나누는 것이 근본적인 해결일 수 있다.

참고: Reactive Streams 사양, RxJava 2 BackpressureStrategy API, Flowable API

같은 카테고리의 글