안동민 개발노트

본문 시작

RxJS와 반응형 프로그래밍

Observable·Observer·구독과 연산자로 시간에 따라 들어오는 여러 값을 처리하고 단일 결과 Promise와의 차이를 비교합니다.

이전 절에서는 Promise와 async/await으로 비동기 작업을 처리하는 방법을 배웠습니다.

이 방식은 단일 작업의 성공/실패를 다루는 데 매우 효과적입니다.

하지만 연속 이벤트 스트림(사용자 입력, 네트워크 이벤트, 타이머 등)을 다루는 복잡한 시나리오에서는 Promise만으로 부족할 수 있습니다.

이 문제를 해결하는 강력한 접근이 반응형 프로그래밍(Reactive Programming)과 라이브러리 RxJS (Reactive Extensions for JavaScript)입니다.

RxJS는 이벤트와 비동기 데이터 흐름을 배열을 다루듯 함수형 방식으로 처리할 수 있게 해줍니다.


반응형 프로그래밍이란?

정의: 반응형 프로그래밍은 데이터 스트림과 변경 사항의 전파에 중점을 둔 비동기 프로그래밍 패러다임입니다.

설명: 기존의 명령형(Imperative) 프로그래밍은 내가 언제 무엇을 할지를 명시적으로 지시하는 방식입니다.

반면 반응형 프로그래밍은 무언가가 발생하면 이렇게 반응하라는 방식으로, 이벤트가 발생했을 때 반응하는 것에 초점을 맞춥니다.

데이터 스트림은 시간이 지남에 따라 발생하는 일련의 이벤트나 값들을 의미합니다.

예를 들어:

  • 사용자 클릭 이벤트
  • 키보드 입력 이벤트
  • HTTP 요청 응답
  • 데이터베이스 변경 알림
  • 타이머 틱

반응형 프로그래밍은 이러한 스트림을 옵저버블(Observable)이라는 개념으로 추상화하고, 다양한 연산자(Operators)를 사용하여 스트림을 변환, 조합, 필터링하는 강력한 도구를 제공합니다.


RxJS의 핵심 개념

RxJS는 반응형 프로그래밍을 자바스크립트/타입스크립트에서 구현하기 위한 라이브러리입니다.

다음 세 가지 핵심 구성 요소를 이해하는 것이 중요합니다.

Observable

정의: Observable은 시간이 지남에 따라 값을 방출(emit)할 수 있는 데이터 스트림을 나타냅니다.

Promise가 단일 값을 비동기적으로 제공하는 반면, Observable은 0개, 1개 또는 여러 개의 값을 동기적 또는 비동기적으로 전달할 수 있습니다.

설명: Observable은 관찰 가능한 데이터 소스입니다.

즉, 누군가 구독(subscribe)하여 그 값을 관찰할 수 있습니다.

Observable은 세 가지 유형의 알림을 방출할 수 있습니다.

  • next: 스트림에서 새로운 값이 방출될 때 호출됩니다.
  • error: 스트림에서 오류가 발생하여 중단될 때 호출됩니다.
  • complete: 스트림이 모든 값을 방출하고 성공적으로 완료될 때 호출됩니다.

Cold Observable은 구독마다 생산 작업을 시작할 수 있습니다. Hot source는 구독과 무관하게 이미 실행 중일 수 있으며, 공유·multicast 정책도 소스와 연산자에 따라 달라집니다.

예시
import { Observable } from 'rxjs';

// Observable 생성
const myObservable = new Observable<number>(subscriber => {
  // next 값을 방출
  subscriber.next(1);
  subscriber.next(2);
  subscriber.next(3);

  // 1초 후 추가 값 방출
  setTimeout(() => {
    subscriber.next(4);
    subscriber.complete(); // 완료 알림
  }, 1000);

  // 구독 해지 시 실행될 클린업 함수 반환
  return () => {
    console.log('Observable: Unsubscribed!');
  };
});

// Observable 구독
console.log('Before subscribe');
const subscription = myObservable.subscribe({
  next: (value) => console.log(`Observable: Got value ${value}`),
  error: (err) => console.error(`Observable: Got error ${err}`),
  complete: () => console.log('Observable: Completed!'),
});
console.log('After subscribe');

// 2초 후 구독 해지 (선택 사항, complete가 호출되면 자동 해지됨)
setTimeout(() => {
  subscription.unsubscribe();
}, 2000);

원문 myObservable은 첫 세 값을 subscribe 호출 중 동기적으로 전달합니다. 반환한 teardown은 로그만 출력하고 setTimeout을 취소하지 않으므로, 조기 해제해도 해당 타이머 작업은 남을 수 있습니다. 아래 수명주기 그림의 자원 정리를 구현하려면 타이머 핸들을 보관해 해제하는 동작까지 필요합니다.

Observer

정의: Observer는 Observable이 방출하는 값에 반응하는 객체입니다.

next, error, complete 메서드를 포함하는 객체입니다.

설명: subscribe() 메서드를 호출할 때 Observer 객체를 인자로 전달합니다.

Observer는 Observable의 라이프사이클 이벤트를 처리하는 콜백 함수들의 묶음입니다.

예시 (위와 동일한 subscribe 부분)
const subscription = myObservable.subscribe({
  next: (value) => console.log(`Observer: Got value ${value}`),
  error: (err) => console.error(`Observer: Got error ${err}`),
  complete: () => console.log('Observer: Completed!'),
});

Operators

정의: Operator는 Observable을 입력으로 받아 새로운 Observable을 출력으로 반환하는 함수입니다.

스트림을 변환, 필터링, 결합하는 데 사용됩니다.

설명: Operators는 RxJS의 가장 강력한 부분 중 하나입니다.

이들은 원본 Observable을 변경하지 않고, 새로운 Observable을 생성하여 스트림 처리의 함수형 및 불변성을 유지합니다.

Operators는 크게 두 가지 범주로 나눌 수 있습니다.

  • 생성 연산자 (Creation Operators): 새로운 Observable을 생성합니다. (of, from, interval, timer, fromEvent 등)
  • 파이프 가능한 연산자 (Pipeable Operators): Observable 스트림을 변환하고 조합합니다. pipe() 메서드 내에서 체인 형태로 사용됩니다. (map, filter, debounceTime, switchMap, mergeMap, takeUntil 등)
예시
import { of, fromEvent, interval } from 'rxjs';
import { map, filter, debounceTime, take } from 'rxjs/operators'; // pipeable operators는 'rxjs/operators'에서 임포트

// 1. 생성 연산자: of
of(10, 20, 30)
  .subscribe(value => console.log(`of operator: ${value}`)); // 10, 20, 30

// 2. 생성 연산자: fromEvent (DOM 이벤트 스트림)
const clicks = fromEvent(document, 'click');
clicks
  .pipe(
    map((event: MouseEvent) => `Click at ${event.clientX}, ${event.clientY}`), // 이벤트를 문자열로 변환
    debounceTime(300) // 300ms 안에 다른 클릭이 없으면 값을 방출 (과도한 이벤트 방지)
  )
  .subscribe(message => console.log(`fromEvent + pipe: ${message}`));

// 3. 생성 연산자: interval + 파이프 가능한 연산자: take, filter
interval(1000) // 1초마다 0부터 증가하는 숫자 방출
  .pipe(
    take(5),    // 처음 5개의 값만 취함
    filter(num => num % 2 === 0), // 짝수만 필터링
    map(num => `Even number: ${num}`) // 문자열로 변환
  )
  .subscribe({
    next: (val) => console.log(`Interval Stream: ${val}`),
    complete: () => console.log('Interval Stream: Completed!'),
  });

pipe()로 값의 변환뿐 아니라 구독과 종료 경로를 구성합니다.

Observable 파이프라인은 값과 종료·teardown 경로를 함께 가진다

OBSERVABLE · LIFECYCLE

Observable 파이프라인은 값과 종료·teardown 경로를 함께 가진다

source에서 operator를 거친 값은 observer로 push되고, error·complete는 terminal이며 unsubscribe는 teardown을 producer 자원으로 되돌린다.

Observable 파이프라인은 값과 종료·teardown 경로를 함께 가진다 source에서 operator를 거친 값은 observer로 push되고, error·complete는 terminal이며 unsubscribe는 teardown을 producer 자원으로 되돌린다. Producer cold · hot source Typed operators map · filter Observer next handler Terminal error complete Teardown unsubscribe cleanup
  1. subscribe가 producer 실행이나

    subscribe가 producer 실행이나 연결을 시작한다.

  2. 값은 typed operator

    값은 typed operator 체인을 거쳐 observer로 push된다.

  3. error와 complete는 terminal

    error와 complete는 terminal 알림이다.

  4. unsubscribe는 terminal 알림과

    unsubscribe는 terminal 알림과 다르다.

  5. teardown이 timer·listener·socket 자원을

    teardown이 timer·listener·socket 자원을 정리한다.

Observable의 cold/hot 여부와 multicast 정책은 라이브러리 연산자와 source 설계에 따라 달라진다.


RxJS와 타입스크립트

RxJS는 타입스크립트로 작성되었으며, 타입스크립트와 함께 사용할 때 그 진가가 발휘됩니다.

  • 강력한 타입 추론: Observable의 제네릭 타입(Observable<T>), Observer의 콜백 함수 인자, 그리고 Operator 체인 전반에 걸쳐 타입이 정확하게 추론됩니다. 이는 개발자가 런타임 오류 대신 컴파일 타임에 타입 관련 문제를 발견하고 해결할 수 있게 합니다.
  • IDE 지원: 타입 정보 덕분에 IDE는 자동 완성, 리팩터링, 오류 검출 등 강력한 개발자 경험을 제공합니다.
  • 명확한 인터페이스: Observable, Observer, Subscription 등의 인터페이스가 명확하게 정의되어 있어 코드의 의도를 쉽게 파악할 수 있습니다.
예시
import { of, Observable } from 'rxjs';
import { map, filter } from 'rxjs/operators';

interface Product {
  id: number;
  name: string;
  price: number;
  available: boolean;
}

const products$: Observable<Product[]> = of([
  { id: 1, name: 'Laptop', price: 1200, available: true },
  { id: 2, name: 'Mouse', price: 25, available: false },
  { id: 3, name: 'Keyboard', price: 75, available: true },
]);

products$.pipe(
  // map을 통해 각 제품 배열을 순회하며 가용성 확인
  map(products => products.filter(p => p.available)),
  // map을 통해 필터링된 제품들의 이름만 추출
  map(availableProducts => availableProducts.map(p => p.name.toUpperCase()))
).subscribe(
  (productNames: string[]) => { // productNames가 string[] 타입임을 정확히 추론
    console.log('Available Products (Uppercase):', productNames.join(', '));
  },
  (error) => console.error('Error:', error)
);

// 컴파일러는 pipe 체인의 각 단계에서 타입이 어떻게 변형되는지 추론할 수 있습니다.
// 예를 들어, 첫 번째 map 후에는 Observable<Product[]> -> Observable<Product[]>,
// 두 번째 map 후에는 Observable<string[]> 으로 변환됨을 압니다.

RxJS의 활용 사례

  • UI 이벤트 처리: 사용자 입력(클릭, 키 입력, 마우스 이동)을 스트림으로 처리하여 디바운싱, 스로틀링, 드래그 앤 드롭 구현 등에 활용됩니다.
  • HTTP 요청: Angular와 같은 프레임워크에서는 HTTP 클라이언트가 RxJS Observable을 반환하여 비동기 데이터 흐름을 처리합니다.
  • 실시간 데이터: 웹소켓을 통한 실시간 채팅, 주식 시세 등 스트리밍 데이터를 다룰 때 유용합니다.
  • 복잡한 비동기 흐름 제어: 여러 비동기 작업을 결합(combineLatest, forkJoin, merge), 경쟁(race), 순차 처리(concatMap), 진행 중 새 입력 무시(exhaustMap), 최신 내부 구독으로 전환(switchMap)하는 데 사용됩니다.
  • 상태 관리: Redux-Observable이나 Ngrx/Effects와 같은 라이브러리에서 비동기 사이드 이펙트(side effect)를 관리하는 데 사용됩니다.

Promise vs Observable

비교 기준Promise와 Observable의 차이
결과Promise는 한 번 정착한 성공 값 또는 거부 이유, Observable은 여러 next 뒤 선택적으로 terminal 알림
시작Promise executor는 동기 호출. Observable은 cold/hot source와 공유 정책에 따라 다름
정리Promise 자체 취소는 없지만 작업이 AbortSignal 등을 지원할 수 있음. unsubscribe는 등록된 teardown을 실행하며 실제 작업 취소는 그 구현에 달림
오류Observable의 error도 해당 구독을 종료. retry는 같은 구독이 계속되는 것이 아니라 재구독
조합Promise는 결과 조합, Observable은 시간에 따른 값과 구독 전환을 연산자로 구성

RxJS를 무조건 쓰는 것이 답은 아닙니다.

이벤트가 얼마나 자주 오고, 취소가 필요한지, 여러 스트림을 합쳐야 하는지에 따라 Promise, async iterator, Observable 중 더 간단한 도구를 고르는 것이 좋습니다.

push·pull과 공유 여부로 비동기 스트림 도구를 고른다

ASYNC STREAM · CHOICE

push·pull과 공유 여부로 비동기 스트림 도구를 고른다

단일 결과인지, 복수 값의 생산 속도를 누가 소유하는지, 여러 소비자가 producer를 공유하는지 질문해 가장 작은 추상화를 고른다.

push·pull과 공유 여부로 비동기 스트림 도구를 고른다 단일 결과인지, 복수 값의 생산 속도를 누가 소유하는지, 여러 소비자가 producer를 공유하는지 질문해 가장 작은 추상화를 고른다. 아니오 예 비동기 데이터 개수 · 제어 · 공유 복수 push가 필요한가? producer가 속도 소유 Promise · async iterable 단일 또는 pull Observable push · composition · multicast
  1. 단일 결과는 Promise를

    단일 결과는 Promise를 쓴다.

  2. 소비자가 pull 속도를

    소비자가 pull 속도를 정하면 async iterable을 쓴다.

  3. producer가 여러 값을

    producer가 여러 값을 push하면 Observable을 검토한다.

  4. 공유 producer와 multicast

    공유 producer와 multicast 요구를 별도로 확인한다.

  5. AbortSignal·return·unsubscribe로 정리 계약을

    AbortSignal·return·unsubscribe로 정리 계약을 명시한다.

Observable이 모든 비동기 작업의 기본값은 아니며 가장 단순한 모델을 먼저 선택한다.