목차
concatMap: 내부 스트림을 한 번에 하나씩 구독
flatMap은 입력마다 만든 내부 스트림을 병합한다. concatMap도 각 입력을 스트림으로 바꾸지만, 현재 내부 스트림이 완료된 다음 다음 스트림을 구독한다. 내부 작업의 처리 속도가 달라도 최종 통지 순서는 입력 순서를 따른다. 이는 순서를 지키려는 대가로 느린 작업이 뒤 작업의 시작을 막는다는 뜻이다.

그림의 안쪽 줄은 각 원본 항목에서 만든 스트림이다. 앞줄이 종료된 뒤 뒷줄이 결과에 연결되므로 오른쪽 출력에서는 두 줄의 값이 뒤섞이지 않는다. subscribeOn을 내부 스트림에 붙여 스케줄러를 바꾸더라도 concatMap의 한 번에 하나씩 구독하는 규칙은 유지된다. 단, 결과를 관찰하는 스레드는 실행에 사용한 스케줄러와 observeOn의 위치에 따라 달라질 수 있다.
구독 시점을 출력해 보기
import io.reactivex.Flowable;
public class ConcatMapExample {
public static void main(String[] args) {
Flowable.range(1, 3)
.concatMap(n -> Flowable.just(n, n * 10)
.doOnSubscribe(ignored -> System.out.println("시작: " + n)))
.subscribe(System.out::println);
}
}
동기 예제는 시작: 1, 1, 10, 시작: 2, 2, 20, 시작: 3, 3, 30 순서로 진행한다. 변환 함수가 스트림을 만들었다는 것과 그 스트림이 실제 구독되었다는 것은 구별해야 한다. 소스가 항목을 미리 전달해 대기열에 담더라도 두 번째 내부 스트림의 결과는 첫 번째가 완료하기 전에는 내려오지 않는다. Flowable.interval(...).take(2) 같은 비동기 내부 스트림을 넣으면 각 스트림이 두 번 통지한 후에야 다음 스트림으로 넘어간다.
이 순서는 주문 상태 변경처럼 앞선 작업의 결과가 뒤 작업의 전제일 때 유용하다. 반면 서로 독립인 외부 API 호출을 모두 concatMap으로 직렬화하면 총 지연 시간이 각 호출 지연의 합에 가까워진다. 처리량이 중요하고 순서가 중요하지 않다면 flatMap을, 동시성은 필요하지만 출력 순서도 중요하다면 concatMapEager를 검토한다. 뒤 옵션은 완료된 결과를 대기시키므로 버퍼와 동시 작업 수를 제한해야 한다.
즉시 오류와 지연 오류
import io.reactivex.Flowable;
public class ConcatMapError {
static Flowable<Integer> work(int n) {
return n == 2
? Flowable.error(new IllegalStateException("2 처리 실패"))
: Flowable.just(n);
}
public static void main(String[] args) {
Flowable.range(1, 3)
.concatMap(ConcatMapError::work)
.subscribe(System.out::println,
e -> System.out.println("오류: " + e.getMessage()));
Flowable.range(1, 3)
.concatMapDelayError(ConcatMapError::work)
.subscribe(System.out::println,
e -> System.out.println("지연 오류: " + e.getMessage()));
}
}
첫 구독은 1 다음 오류: 2 처리 실패를 출력한다. 오류가 종료 신호이므로 3을 처리하지 않는다. 두 번째 구독은 1, 3을 출력하고 마지막에 지연 오류: 2 처리 실패를 출력한다. concatMapDelayError는 오류를 성공으로 바꾸지 않는다. 다른 입력의 결과를 최대한 받되 마지막에는 실패 신호가 필요한 경우에 맞는다. 여러 내부 스트림이 실패하면 오류가 합쳐져 전달될 수 있으므로, 실패 항목별 진단이 필요하면 각 내부 스트림에서 해당 입력 ID를 붙여 오류를 기록하거나 결과 타입으로 표현한다.
take(2)처럼 하위에서 취소하면 현재 내부 스트림과 상위 구독이 해지된다. concatMapDelayError도 취소 후에는 남은 입력을 끝까지 처리해 오류를 모으지 않는다. 무한히 완료되지 않는 첫 내부 스트림을 넣으면 뒤 항목은 영원히 시작하지 않으므로 timeout 또는 내부 작업의 종료 조건을 정해야 한다. retry를 전체 결과 뒤에 붙이면 앞서 성공한 입력도 재구독할 수 있다. 재시도 범위를 개별 내부 작업으로 좁혀야 중복 부작용을 통제하기 쉽다.
