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

7. RxJava defer 연산자

목차

RxJava defer: 구독마다 소스를 새로 만들기

Flowable.defer(factory)는 Flowable을 선언하는 시점에는 factory를 실행하지 않는다. 구독할 때마다 함수를 호출해 그 구독에 사용할 새 Flowable을 얻는다. just는 준비된 값을 통지하고, fromCallable은 구독마다 값 하나를 계산한다. defer는 값뿐 아니라 어떤 소스를 사용할지 구독 시점에 결정한다.

defer 연산자의 시간선

그림에서 구독마다 별도 소스가 만들어지는 모습을 본다. 시간선의 간격은 실제 초 단위가 아니다. defer가 자동으로 스케줄러를 바꾸거나 지연 타이머를 넣는 것도 아니다.

just와 구독 두 번 비교

import io.reactivex.Flowable;
import java.util.concurrent.atomic.AtomicInteger;

public class DeferExample {
    public static void main(String[] args) {
        AtomicInteger counter = new AtomicInteger();
        Flowable<Integer> fixed = Flowable.just(counter.incrementAndGet());
        Flowable<Integer> fresh = Flowable.defer(
            () -> Flowable.just(counter.incrementAndGet())
        );

        fixed.subscribe(System.out::println); // 1
        fixed.subscribe(System.out::println); // 1
        fresh.subscribe(System.out::println); // 2
        fresh.subscribe(System.out::println); // 3
    }
}

fixed를 만드는 줄에서 인자 counter.incrementAndGet()가 먼저 실행된다. fresh를 만드는 줄에서는 카운터가 움직이지 않는다. 각각의 구독 때 팩터리가 새 just 소스를 만들어 2, 3을 통지한다. 별도 스케줄러가 없으므로 위 예제는 main 스레드에서 순서대로 실행된다. 원문의 LocalTime.now() 예제도 같은 원리지만 출력 시각이 항상 다르다는 보장은 없다.

재구독과 소스 선택

AtomicInteger attempts = new AtomicInteger();
Flowable<Integer> task = Flowable.defer(() -> {
    int attempt = attempts.incrementAndGet();
    return attempt == 1
        ? Flowable.<Integer>error(new IllegalStateException("첫 시도 실패"))
        : Flowable.just(attempt);
});
task.retry(1).subscribe(System.out::println); // 2

첫 구독의 팩터리는 오류 소스를 만든다. retry(1)은 그 오류 뒤에 다시 구독하므로 팩터리가 재실행돼 두 번째 소스를 만든다. 이때 결과가 2다. 오류난 작업에서 이미 파일 쓰기나 외부 결제가 일어났다면 재구독은 부작용도 다시 실행할 수 있으므로 멱등성을 확인해야 한다.

구독자별로 다른 설정·사용자 컨텍스트의 소스를 선택해야 할 때 defer가 유용하다. 단순히 파일 내용 한 값을 늦게 읽으려는 목적이라면 fromCallable이 짧다. 팩터리가 예외를 던지거나 null 소스를 반환하면 정상적인 빈 결과가 아니라 오류가 된다. 정상 부재는 Flowable.empty()를 반환한다. 구독을 취소하면 그 소스의 남은 통지는 멈추지만, 이미 시작된 블로킹 I/O의 중단 가능성은 소스 구현에 달려 있다. 비동기로 실행하려면 선택한 소스에 적절한 subscribeOn을 명시한다.

참고: ReactiveX Defer, RxJava 2 Flowable API

같은 카테고리의 글