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

2. RxJava flatMap 연산자

목차

flatMap: 한 입력에서 여러 값을 만들기

map은 입력 하나에서 값 하나를 반환한다. flatMap의 변환 함수는 값 대신 Publisher를 반환한다. 따라서 입력 하나를 여러 출력으로 펼치거나 Flowable.empty()로 출력하지 않을 수 있다. 반환한 내부 스트림은 상위 스트림의 구독을 통해 구독되고, 결과가 하나의 하위 구독자로 병합된다. 여기서는 RxJava 2의 io.reactivex.Flowable을 사용한다.

입력마다 내부 스트림을 만드는 flatMap 다이어그램

그림에서 위쪽 점 하나가 아래쪽 작은 스트림 하나를 만든다. 작은 스트림들이 내보내는 점의 개수는 같을 필요가 없다. 내부 스트림이 서로 다른 시점에 값을 내보내면 오른쪽 결과는 입력 순서 대신 실제 도착 순서가 된다.

0개, 1개, 여러 개로 변환

import io.reactivex.Flowable;

public class FlatMapExample {
    public static void main(String[] args) {
        Flowable.just("A", "", "B")
                .flatMap(text -> text.isEmpty()
                        ? Flowable.<String>empty()
                        : Flowable.just(text.toLowerCase(), text))
                .subscribe(System.out::println, Throwable::printStackTrace);
    }
}

이 동기 예제의 출력은 a, A, b, B다. 첫 입력 A는 두 값을 만들고 빈 문자열은 empty()로 값을 만들지 않는다. empty()도 정상 완료 신호는 보내므로 다음 입력 B의 처리를 막지 않는다. 변환 함수에서 null을 반환하는 것은 건너뛰기가 아니다. RxJava 스트림은 null 항목을 허용하지 않으며 오류 경로로 간다.

입력과 결과를 함께 출력하려면 내부 스트림의 map에서 입력을 캡처한다.

Flowable.just("A", "B")
        .flatMap(original -> Flowable.just(original.toLowerCase())
                .map(converted -> original + " -> " + converted))
        .subscribe(System.out::println);
// A -> a
// B -> b

이 코드는 위의 Flowable import가 있는 위치에서 이어 쓰는 코드 조각이다. flatMap(mapper, resultSelector) 오버로드도 같은 목적에 쓸 수 있지만, 단순 조합은 내부 map이 데이터 흐름을 읽기 쉽다.

출력 순서는 왜 달라지는가

flatMap은 내부 스트림 하나가 완료할 때까지 다음 내부 스트림의 구독을 기다리지 않는다. 예를 들어 주문 ID마다 비동기 조회를 시작하면 나중에 시작한 조회가 먼저 완료할 수 있다. 아래 delay는 타이머 스케줄러를 사용하므로 출력은 입력 순서 보장의 예가 아니다.

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

public class FlatMapOrder {
    public static void main(String[] args) {
        Flowable.just(1, 2, 3)
                .flatMap(id -> Flowable.just(id)
                        .delay(id == 1 ? 300 : 50, TimeUnit.MILLISECONDS))
                .blockingForEach(System.out::println);
    }
}

첫 입력 1의 완료가 늦으므로 2나 3이 먼저 출력될 수 있다. 동시에 구독하는 내부 스트림 수를 제한하려면 flatMap(mapper, maxConcurrency) 오버로드를 살펴본다. 이는 동시성의 상한이며 결과 순서를 복구하지는 않는다. 순서가 업무 규칙이면 concatMap을 선택한다. 동시 실행과 입력 순서를 함께 원하면 concatMapEager를 검토하되, 먼저 끝난 결과를 순서가 올 때까지 보관하므로 메모리 사용량이 늘 수 있다.

여러 내부 스트림의 결과가 병합되는 다이어그램

그림의 핵심은 내부 스트림의 시간축이 겹칠 수 있다는 점이다. subscribeOn으로 각 작업을 다른 스케줄러에 배치해도 flatMap이 순서를 맞춰주지는 않는다. 반대로 flatMap만 붙여서는 작업이 자동으로 다른 스레드에서 실행되지 않는다.

오류와 취소를 추적하기

Flowable.just(1, 2, 0, 4)
        .flatMap(n -> n == 0
                ? Flowable.<Integer>error(new ArithmeticException("0으로 나눔"))
                : Flowable.just(12 / n))
        .subscribe(System.out::println,
                error -> System.out.println("실패: " + error.getMessage()));

동기 실행에서는 12, 6 다음 오류가 오고 4에 대한 결과는 없다. 기본 flatMap은 내부 오류를 전체 결과의 오류로 전파하며 활성 내부 구독도 정리한다. 비동기 입력에서는 이미 진행 중인 결과와 오류가 경합할 수 있으므로, 오류 전까지 관찰한 값의 목록을 계약으로 삼으면 안 된다. flatMapDelayError는 오류를 미뤄 다른 내부 스트림을 처리하지만 최종 결과는 여전히 오류다.

하위 구독자가 take(2)로 두 결과만 받은 뒤 취소하면 상위 및 활성 내부 구독에도 취소가 전달된다. 이미 외부 서비스에 보낸 요청을 네트워크 수준에서 취소할 수 있는지는 그 요청 API와 연결한 Flowable 구현에 달려 있다. 구독 해제가 데이터베이스 거래를 자동으로 되돌린다고 가정해서는 안 된다.

참고: ReactiveX FlatMap, RxJava 2 Flowable API.

같은 카테고리의 글