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 · LIFECYCLE
Observable 파이프라인은 값과 종료·teardown 경로를 함께 가진다
source에서 operator를 거친 값은 observer로 push되고, error·complete는 terminal이며 unsubscribe는 teardown을 producer 자원으로 되돌린다.
- subscribe가 producer 실행이나
subscribe가 producer 실행이나 연결을 시작한다.
- 값은 typed operator
값은 typed operator 체인을 거쳐 observer로 push된다.
- error와 complete는 terminal
error와 complete는 terminal 알림이다.
- unsubscribe는 terminal 알림과
unsubscribe는 terminal 알림과 다르다.
- 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 중 더 간단한 도구를 고르는 것이 좋습니다.
ASYNC STREAM · CHOICE
push·pull과 공유 여부로 비동기 스트림 도구를 고른다
단일 결과인지, 복수 값의 생산 속도를 누가 소유하는지, 여러 소비자가 producer를 공유하는지 질문해 가장 작은 추상화를 고른다.
- 단일 결과는 Promise를
단일 결과는 Promise를 쓴다.
- 소비자가 pull 속도를
소비자가 pull 속도를 정하면 async iterable을 쓴다.
- producer가 여러 값을
producer가 여러 값을 push하면 Observable을 검토한다.
- 공유 producer와 multicast
공유 producer와 multicast 요구를 별도로 확인한다.
- AbortSignal·return·unsubscribe로 정리 계약을
AbortSignal·return·unsubscribe로 정리 계약을 명시한다.
Observable이 모든 비동기 작업의 기본값은 아니며 가장 단순한 모델을 먼저 선택한다.