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

4. RxJava Buffer 연산자

목차

buffer: 여러 통지를 목록 하나로 묶기

이벤트를 건별로 처리하는 대신 API 한 번에 여러 건을 보내거나 화면을 일정 간격으로 갱신하려면 buffer를 쓸 수 있다. buffer(3)은 원본의 세 값을 모아 List 하나로 통지한다. 구독자는 더 적은 횟수로 onNext를 받지만 목록에 담긴 총 항목 수는 그대로다. 작업을 병렬로 실행하거나 생산 속도를 자동으로 낮추는 연산자는 아니다.

개수 기준: 마지막 묶음은 어떻게 되는가

import io.reactivex.Flowable;

public class BufferCount {
    public static void main(String[] args) {
        Flowable.range(1, 8)
                .buffer(3)
                .subscribe(System.out::println);
    }
}

출력은 [1, 2, 3], [4, 5, 6], [7, 8]이다. range(1, 8)이 여덟 값을 동기적으로 내보내고 정상 완료하므로, 세 개를 채우지 못한 마지막 목록도 통지된다. 반면 소스가 완료하지 않고 두 값만 보낸 뒤 기다리면 마지막 목록도 기다린다. 이때 목록이 내려오지 않는 이유를 하위 구독의 오류로 착각하기 쉽다. take, timeout, 시간 기준 버퍼처럼 종료 또는 방출 조건을 별도로 정한다.

buffer(3)은 항목 3개마다 한 목록을 만든다. buffer(3, 1)은 세 값을 모으되 한 값씩 창을 이동하므로 [1,2,3], [2,3,4]처럼 목록이 겹친다. 겹친 항목은 여러 목록에 저장되어 메모리와 후속 처리 횟수가 증가한다. 어떤 묶음이 필요한지 결정한 다음 count와 skip을 선택해야 한다.

시간 기준: 도착 속도가 일정하지 않을 때

import io.reactivex.Flowable;
import java.util.concurrent.TimeUnit;

public class BufferTime {
    public static void main(String[] args) {
        Flowable.interval(100, TimeUnit.MILLISECONDS)
                .take(8)
                .buffer(350, TimeUnit.MILLISECONDS)
                .blockingForEach(System.out::println);
    }
}

interval은 기본 계산 스케줄러에서 0부터 값을 내보낸다. 350ms마다 버퍼를 닫으므로 약 100ms 간격의 값들이 묶인다. 시간 경계에 값이 도착하는 경우 두 버퍼 중 어디에 들어갈지는 스케줄링 순서에 좌우되므로 [0,1,2]처럼 정확한 목록 구성을 계약으로 두면 안 된다. .take(8)이 정상 완료하면 마지막 버퍼도 방출되고 blockingForEach가 종료를 기다린다. 무한 interval에 blockingForEach만 붙이면 호출 스레드는 끝나지 않는다.

시간 버퍼는 해당 시간에 값이 없더라도 빈 목록을 낼 수 있다. 빈 배치를 외부 API에 보내면 낭비가 되므로 filter(batch -> !batch.isEmpty())를 뒤에 붙일 수 있다. 개수와 시간 양쪽 상한이 필요하다면 buffer(timespan, unit, count) 오버로드를 검토한다. 드문 이벤트는 오래 묵지 않고, 폭주 때는 한 목록이 무한히 커지지 않도록 하기 위해서다.

메모리와 취소 판단

buffer의 배치는 통지 단위이지 저장 용량 제한이나 배압 정책이 아니다. 상위가 빠르고 하위의 onNext 처리가 느리면 상위 소스의 특성 및 중간 연산자에 따라 대기열이 커지거나 배압 오류가 날 수 있다. Flowable.interval처럼 시간에 맞춰 내보내는 소스는 느린 구독자를 항상 기다리지는 않는다. 배치를 늘리기 전에 처리 속도, 한 목록의 최대 크기, 큐 길이와 오류율을 관찰한다.

하위가 take(2)로 목록 두 개만 받고 취소하면 소스 구독과 버퍼 타이머도 해지된다. 아직 버퍼에 남은 값은 정상 완료 때처럼 추가로 방출되지 않는다. 외부로 전송한 첫 두 묶음이 실제로 저장되었는지는 전송 API의 성공 응답을 별도로 확인해야 한다. 중요한 데이터를 배치로 보낼 때는 실패한 목록의 재시도와 중복 저장 방지 기준도 함께 정한다.

참고: ReactiveX Buffer, RxJava 2 Flowable API.

같은 카테고리의 글