스트림 병합 전략: Merge, Concat, Switch, Exhaust
ReactiveX 생태계 (RxJava, RxJS, RxSwift 등) 에서는 여러 개의 소스 스트림을 하나의 결과 스트림으로 통합하는 데 사용되는 연산자들의 행동 패턴이 서로 다릅니다. 핵심적인 네 가지 연산자는 데이터를 병합할 때 시간적 순서와 완료를 어떻게 처리하느냐에 따라 구분됩니다.
- Merge: 각 내부 스트림의 순서를 무시하고, 데이터가 도달한 즉시 출력합니다. 병렬성이 가장 높은 방식입니다.
- Concat: 순차적인 처리를 강제합니다. 첫 번째 스트림이 완전히 종료될 때까지 대기하며, 이후 두 번째 스트림을 시작하는 식입니다.
- Switch / SwitchLatest: 최신 상태만 반영합니다. 새로운 내부 스트림이 생성되면 이전 스트림은 취소되고, 새로 들어온 스트림의 데이터만 전송됩니다.
- Exhaust: 현재 진행 중인 스트림이 완료될 때까지 새로운 요청을 무시합니다. 이전 작업이 끝나야 다음 작업을 고려하는 형태입니다.
변환 및 평탄화: Map 계열 연산자 변형
데이터 스트림의 각 아이템을 다시 스트림으로 변환한 후 이를 편평하게 펼치는 과정에서는 flatMap, concatMap, switchMap, exhaustMap 등의 연산자가 사용됩니다. 이는 내부적으로 앞선 병합 전략과 동일하게 작동하지만, 먼저 매핑 함수를 통해 스트림을 생성한다는 점이 특징입니다.
RxJava 를 활용한 시나리오에서 각 연산자의 차이를 검증하기 위한 코드를 재구현해 보겠습니다. 기존 변수명을 변경하고 지연 로직을 조금 다르게 구성했습니다.
// 시나리오 A: Flatten 을 통한 병렬 처리 확인
val inputValues = listOf("start-01", "start-02", "start-03")
val testEnvironment = TestScheduler()
inputValues.toFlowable()
.flatMap { rawValue ->
val randomDelay = Random().nextDouble() * 5.0
Flowable.just(String.format("%s-ok", rawValue))
.delay(randomDelay, TimeUnit.SECONDS, testEnvironment)
}
.toList()
.subscribeOn(Schedulers.io())
.observeOn(Schedulers.test())
.doOnNext { results -> println("결과 집합: $results") }
.blockingAwait()
// 예상 동작: 입력 순서 없이 도착한 순서대로 결과가 모일 수 있음
testEnvironment.triggerActions()
// 시나리오 B: 순차적 흐름 보장
val sourceData = listOf("task-A", "task-B", "task-C")
val syncScheduler = Single.instance()
sourceData.asFlowable()
.concatMap { identifier ->
val duration = Random().nextInt(3000) + 1000
Flowable.just("${identifier}-complete")
.delay(duration.toLong(), TimeUnit.MILLISECONDS, syncScheduler)
}
.blockingForEach { item -> System.out.println(item) }
// 예상 동작: task-A 완료 후 task-B 가 시작됨
// 시나리오 C: 최신 요청 위주 처리
val requestList = listOf("req-1", "req-2", "req-3")
val virtualClock = TestScheduler()
requestList.toFlowable()
.switchMap { param ->
val latency = Random().nextInt(5)
Flowable.just(param.lowercase())
.delay(latency.toLong(), TimeUnit.SECONDS, virtualClock)
}
.lastOrError() // 마지막 남은 것만 관심 있음
.subscribe { finalResult -> println("최종 값: $finalResult") }
virtualClock.advanceTimeBy(1, TimeUnit.MINUTES)
// 예상 동작: req-3 이 도착하면서 이전 req 들은 취소되고 최종 결과는 'req-3' 관련값 또는 null 일 수 있음
Rx.NET 의 특수 연산자: ManySelect 대 SelectMany
.NET 환경의 Rx 라이브러리에는 다른 프레임워크와 구별되는 몇 가지 연산자가 존재합니다. 특히 ManySelect는 현재 시점에서부터 미래까지 나올 데이터들의 윈도우를 생성하여 계산하는 반면, SelectMany는 일반적인 평탄화 연산자와 유사하게 개별 아이템을 스트림으로 전환합니다.
다음 예제는 범위를 생성하여 두 연산자의 결과를 시각화합니다.
using System;
using System.Collections.Generic;
using System.Linq;
using System.Reactive.Linq;
public class DemoProgram
{
public static void Main(string[] args)
{
var rangeSource = Observable.Range(1, 10);
// ManySelect: 누적 합산 윈도우 활용
Console.WriteLine("--- ManySelect Sequence ---");
rangeSource
.ManySelect(x => x.Sum(), Scheduler.CurrentThread)
.Concat()
.Subscribe(result => Console.WriteLine($"Current Sum: {result}"));
// SelectMany: 일반적 분산 매핑
Console.WriteLine("--- SelectMany Sequence ---");
rangeSource
.SelectMany(num =>
Observable.Range(num, 10 - num + 1)
.Sum()
)
.Subscribe(val => Console.WriteLine($"Mapped Result: {val}"));
Console.WriteLine("Process Complete.");
}
}
다중 스트림 동기화: Zip, CombineLatest 및 기타
두 개 이상의 독립된 스트림으로부터 데이터를 추출하여 결합해야 할 때는 Zip, CombineLatest, WithLatestFrom, ForkJoin 중 적절한 연산자를 선택해야 합니다. 각각 데이터 수집의 타이밍 조건이 엄격하게 정의되어 있습니다.
- Zip: 두 스트림 모두 동일한 인덱스 위치 (N 번째 항목) 에서 데이터를 가져왔을 때만 결합합니다. 느린 쪽이 전체 속도를 좌우합니다.
- CombineLatest: 한쪽에서 새로운 데이터가 발생하면 상대방의 최근 데이터를 가져와서 결합합니다. 최소한 양쪽 모두 한번 이상 배출되었어야 합니다.
- WithLatestFrom: 메인 스트림에서만 트리거되며, 서브 스트림의 마지막 값을 참조합니다.
- ForkJoin: 모든 스트림이 성공적으로 종결될 때 마지막 요소들끼리만 결합합니다.
C# 기반의 시뮬레이션 코드로 각 연산자의 타임라인 차이를 확인해 봅니다.
public IObservable<int> CreateSlowStream(int seed, List<int> steps)
{
return Observable.Generate(seed,
s => s < seed + steps.Count - 1,
s => s + 1,
s => s,
s => TimeSpan.FromMilliseconds(steps[s - seed] * 100));
}
var streamOne = CreateSlowStream(10, new List<int> { 5, 5, 5 });
var streamTwo = CreateSlowStream(20, new List<int> { 4, 4, 4, 4 });
// Zip 실행: 쌍으로 된 데이터 생성
streamOne.Zip(streamTwo, (a, b) => new { Left = a, Right = b })
.Subscribe(p => Console.WriteLine($"[{p.Left}, {p.Right}]"));
// CombineLatest 실행: 상태 변화마다 즉시 결합
streamOne.CombineLatest(streamTwo, (a, b) => $"L:{a}-R:{b}")
.Subscribe(msg => Console.WriteLine($"State: {msg}"));
// ForkJoin 실행: 최종 종료 시점에만 결합
streamOne.ForkJoin(streamTwo, (a, b) => $"{a}+{b}")
.Subscribe(finalVal => Console.WriteLine($"Final Join: {finalVal}"));