목차
RxJava fromCallable: 구독할 때 값 계산하기
Flowable.fromCallable(callable)은 소스를 만들 때 함수를 실행하지 않는다. 구독할 때 Callable.call()을 호출해 반환값 하나를 onNext로 보내고 onComplete로 끝낸다. 함수가 예외를 던지면 그 구독에는 onError가 오고 완료 신호는 오지 않는다. “값 하나”라는 결과는 just와 같아 보여도 값을 만드는 시점이 다르다.

그림에서 함수를 통해 값 하나가 만들어져 내려온다. 그림은 함수를 언제 구독하는지나 어느 스레드에서 실행하는지를 보장하지 않으므로 코드에서 확인해야 한다.
두 번 구독하면 두 번 계산
import io.reactivex.Flowable;
import java.util.concurrent.atomic.AtomicInteger;
public class CallableExample {
public static void main(String[] args) {
AtomicInteger calls = new AtomicInteger();
Flowable<Integer> source = Flowable.fromCallable(calls::incrementAndGet);
System.out.println(calls.get()); // 0: 아직 실행하지 않음
source.subscribe(System.out::println); // 1
source.subscribe(System.out::println); // 2
}
}
fromCallable을 호출한 시점의 카운터는 0이다. 첫 구독에서 1, 다음 구독에서 2가 된다. 별도 스케줄러가 없으므로 위 출력은 main 스레드에서 동기적으로 진행된다. 원문의 System.currentTimeMillis() 예제도 구독 시 시계를 읽지만 두 구독이 아주 가까우면 같은 밀리초 값을 반환할 수 있다. 값이 다르다는 것보다 함수가 다시 실행된다는 것이 계약이다.
Flowable.just(calls.incrementAndGet())를 사용하면 인자식이 소스를 만들 때 먼저 계산되고 두 구독자 모두 같은 1을 받는다. 요청 때마다 새 설정·파일 내용·DB 조회 결과를 얻어야 한다면 fromCallable이 맞고, 이미 확정된 값이면 just가 단순하다. 함수 자체가 다른 Flowable을 만들어 반환해야 한다면 defer를 사용한다.
I/O를 넣을 때의 스케줄러와 취소
Flowable<String> read = Flowable.fromCallable(() ->
java.nio.file.Files.readString(java.nio.file.Path.of("sample.txt")));
read.subscribeOn(io.reactivex.schedulers.Schedulers.io())
.subscribe(System.out::println, Throwable::printStackTrace);
fromCallable 자체가 비동기는 아니다. 위처럼 subscribeOn을 지정하면 구독과 계산이 해당 스케줄러에서 실행된다. 따라서 main이 바로 종료하는 짧은 프로그램에서는 값이 출력되기 전에 프로세스가 끝날 수 있다. 운영 코드에서는 구독 수명, 파일 읽기의 타임아웃·예외, 결과를 어느 스레드에서 사용할지를 함께 설계한다.
구독을 취소하면 아직 시작되지 않은 작업의 실행은 막을 수 있지만, 이미 시작된 블로킹 Callable이 반드시 즉시 멈춘다고 보장할 수는 없다. 사용한 I/O API의 취소·인터럽트 계약을 확인한다. 함수가 null을 반환하면 RxJava 2의 값 계약을 위반해 오류가 된다. 정상적으로 값이 없을 수 있는 함수라면 Maybe.fromCallable 등의 계약을 선택한다.
