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

5. RxJava interval 연산자

목차

RxJava interval: 일정 간격으로 번호 보내기

Flowable.interval(1, TimeUnit.SECONDS)는 구독 뒤 약 1초가 지나 0L을 보내고, 이후 약 1초 간격으로 1L, 2L을 이어 보낸다. 기본적으로 완료하지 않는다. 반환값은 시각 자체가 아니라 0부터 증가하는 번호다. 시간은 스케줄러와 시스템 부하의 영향을 받으므로 정밀한 실시간 시계로 사용하지 않는다.

interval 연산자의 시간선

그림에는 반복되는 값 원이 있지만 완료선은 없다. 그래서 concatWith(다음 소스)를 바로 붙이면 다음 소스가 시작되지 않는다. 다음 작업으로 넘어가려면 take(3)처럼 종료 조건을 먼저 만들어야 한다.

세 번만 받고 끝내기

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

public class IntervalExample {
    public static void main(String[] args) throws InterruptedException {
        CountDownLatch done = new CountDownLatch(1);
        Flowable.interval(100, TimeUnit.MILLISECONDS)
            .take(3)
            .subscribe(
                value -> System.out.println("값 " + value),
                error -> { error.printStackTrace(); done.countDown(); },
                () -> { System.out.println("완료"); done.countDown(); }
            );
        done.await(2, TimeUnit.SECONDS);
    }
}

예상되는 값의 순서는 0, 1, 2, 완료다. 정확한 출력 시각과 스레드 이름은 실행 환경에 따라 다르다. interval은 기본적으로 Schedulers.computation()을 사용하므로 구독을 호출한 main이 곧바로 끝나면 출력을 놓칠 수 있다. 위 CountDownLatch는 예제가 끝날 때까지 기다리기 위한 장치이고, take(3)이 실제 스트림의 종료 조건이다.

interval(initialDelay, period, unit, scheduler) 형태로 첫 값까지의 지연과 실행 스케줄러를 명시할 수 있다. 빠르고 결정적인 테스트에는 실제 Thread.sleep 대신 TestScheduler로 가상 시간을 진행한다.

소비가 느리면 무슨 일이 생길까

Flowable.interval(10, TimeUnit.MILLISECONDS)
    .onBackpressureLatest()
    .take(3)
    .subscribe(System.out::println);

시간 기반 소스는 시계가 움직일 때 값을 만들기 때문에 소비자가 request(n)을 늦춘다고 생산 시계를 멈출 수 없다. onBackpressureLatest()는 느린 소비 상황에서 중간 번호를 버리고 최신 값 중심으로 전달하려는 정책이다. 위 작은 코드에서는 소비자가 느리지 않으므로 번호가 건너뛰지 않을 수도 있다. 어떤 값이 반드시 모두 필요하다면 최신값 정책은 맞지 않으며 상류 생성률을 줄이거나 내구성 있는 큐를 고려한다.

화면이 닫히거나 서버 요청이 취소되면 구독도 해제해야 타이머 작업이 남지 않는다. take로 정상 완료시키거나 반환된 Disposable을 dispose()한다. interval은 폴링 신호 같은 주기적 작업에 맞지만, 실제 작업 시간이 주기보다 길다면 작업이 겹칠지·밀릴지·건너뛸지 별도로 설계한다.

참고: ReactiveX Interval, RxJava 2 Flowable API

같은 카테고리의 글