ReactiveX 스트림 조작 연산자의 동작 원리와 비교

스트림 병합 전략: 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}"));

태그: reactive-programming rx-java rx-net observable-pattern async-flows

7월 27일 00:38에 게시됨