RxJS 실전: Observable 만들기, switchMap·mergeMap·concatMap 선택, 에러 처리, 검색 자동완성
이 글의 핵심
RxJS에서 Observable을 만들고 오퍼레이터로 다듬는 법, 동시 요청을 다루는 세 가지 map 오퍼레이터의 차이, 에러 재시도, 검색 자동완성 구현, Cold와 Hot Observable 공유 문제를 정리합니다.
한참 전에, 검색창에 타이핑할 때마다 API를 때리다가 서버에 혼났던 적이 있습니다. input에 addEventListener를 몇 겹 쌓고, debounce는 손으로 짜고, 이전 요청을 취소할 생각은 못 하고… 그때 “아, 이건 시간에 따라 흐름이 있다”는 걸 머리로는 알았는데 코드로는 계속 if문 덩어리로만 썼습니다. 그게 반응형 프로그래밍이 처음과 만나는 지점입니다. 이벤트·시간·네트워크가 겹칠수록 “지금 무슨 시퀀스가 돌고 있지?”를 선언해 두고 싶어지는 것입니다.
실제로는, 비동기면 Promise로 충분할 때 많아요. 한 번만 resolve되고, 취소가 없으며, 흐름이 직선이면 async/await가 제일 읽기 쉽습니다. RxJS는 그 다음 단계예요. 여러 값이 이어질 수 있고, 중간에 끊을 수 있고, 여러 스트림을 섞을 수 있을 때 이 도구를 꺼내면 이득이 큽니다. “모든 API 호출에 Observable”이면 팀이 미칩니다. 딱 닿는 경계—검색 자동완성, WebSocket, 폼·라우트 이벤트가 뒤엉킨 화면—에 두는 게 맞아요.
Observable 만들기: of, from, interval, fromEvent
RxJS는 그것을 Observable이라는 스트림과, map·filter 같은 오퍼레이터로 이어 붙이는 방식으로 풉니다. of랑 from으로 출발지를 만들며, interval·timer는 시간, fromEvent는 DOM이나 Node 이벤트를 그대로 스트림으로 바꿉니다.
import { of, from } from 'rxjs';
// of: 값으로 생성
const source$ = of(1, 2, 3, 4, 5);
source$.subscribe((value) => console.log(value));
// from: Array, Promise, Iterable로 생성
const array$ = from([1, 2, 3]);
const promise$ = from(fetch('/api/users'));
import { interval, timer } from 'rxjs';
// 1초마다 emit
const interval$ = interval(1000);
interval$.subscribe((value) => console.log(value));
// 3초 후 한 번만 emit
const timer$ = timer(3000);
timer$.subscribe(() => console.log('Timer!'));
import { fromEvent } from 'rxjs';
const button = document.getElementById('btn');
const click$ = fromEvent(button, 'click');
click$.subscribe(() => console.log('Clicked!'));
생성 함수들은 이름만 비슷할 뿐 “끝나는 스트림”과 “끝나지 않는 스트림”으로 나뉜다는 점이 중요합니다. of, from([...]), timer(3000)은 값을 다 내보내면 complete되어 구독이 자동으로 정리됩니다. 반면 interval과 fromEvent는 스스로 끝나지 않으므로, 구독을 해제하지 않으면 컴포넌트가 사라진 뒤에도 콜백이 계속 호출되고 DOM 요소도 메모리에서 해제되지 않습니다. SPA에서 페이지를 오갈 때마다 로그가 두 배, 세 배로 찍히는 증상이 이 누수의 전형적인 모습입니다. subscribe()가 돌려주는 Subscription의 unsubscribe()를 부르거나, take(n)·takeUntil(notifier$)처럼 스트림 자체에 종료 조건을 넣어야 합니다.
from(fetch(...))처럼 Promise를 감쌀 때도 주의할 점이 있습니다. Promise는 만들어지는 순간 이미 실행을 시작하므로, 이 Observable은 구독하기 전에 요청이 나가고 여러 번 구독해도 요청은 한 번뿐입니다. “구독할 때 실행”이라는 Observable의 일반적인 성질과 다르게 동작하므로, 구독 시점에 요청을 보내고 싶다면 defer(() => fetch(...))로 감쌉니다.
오퍼레이터로 스트림 다듬기
중간에 값을 갈고 닦는 건 pipe 안에서 해요. 짝수만 골라서 열 곱한다든지, 검색어를 debounce하고 직전 값과 다를 때만 넘긴다든지.
import { of } from 'rxjs';
import { map, filter } from 'rxjs/operators';
const source$ = of(1, 2, 3, 4, 5);
source$
.pipe(
filter((x) => x % 2 === 0),
map((x) => x * 10)
)
.subscribe((value) => console.log(value));
// 20, 40
import { fromEvent } from 'rxjs';
import { debounceTime, distinctUntilChanged, map } from 'rxjs/operators';
const input = document.getElementById('search');
const search$ = fromEvent(input, 'input');
search$
.pipe(
map((event) => (event.target as HTMLInputElement).value),
debounceTime(500),
distinctUntilChanged()
)
.subscribe((value) => {
console.log('Search:', value);
});
오퍼레이터 순서가 결과를 바꿉니다. debounceTime(500)은 마지막 입력 후 500ms 동안 새 입력이 없을 때만 값을 내보내므로, 빠르게 “react”를 치면 중간의 “r”, “re”, “rea”는 버려지고 “react” 하나만 남습니다. 그 다음 distinctUntilChanged()는 직전에 내보낸 값과 같으면 걸러 냅니다. 사용자가 “react”를 치고 한 글자 지웠다 다시 쳐서 결과적으로 같은 검색어가 되면 요청을 다시 보내지 않는 것입니다. 두 오퍼레이터의 순서를 바꾸면 키 입력마다 비교가 일어나 거의 걸러지지 않습니다. map에서 이벤트를 문자열로 먼저 바꿔 두는 것도 이유가 있는데, 이벤트 객체는 매번 새 객체라 distinctUntilChanged의 참조 비교(===)로는 절대 같다고 판정되지 않기 때문입니다.
여러 요청 처리하기: switchMap·mergeMap·concatMap
요청이 겹칠 때는 switchMap이 가장 먼저 손에 익어요. 새 입력이 오면 안쪽 subscription을 갈아끼우므로 이전 HTTP가 살아 있으면 경쟁 상태(레이스)가 덜해요. 반대로 다 병렬로 때리고 싶으면 mergeMap, 줄 세우고 싶으면 concatMap 쪽 느낌입니다.
import { fromEvent, of } from 'rxjs';
import { switchMap, mergeMap, concatMap } from 'rxjs/operators';
const button = document.getElementById('btn');
const click$ = fromEvent(button, 'click');
// switchMap: 이전 결과 무시 (fetch Promise 자체는 실제로 중단되지 않음)
click$
.pipe(
switchMap(() => fetch('/api/users').then((r) => r.json()))
)
.subscribe((users) => console.log(users));
// mergeMap: 병렬 실행
click$
.pipe(
mergeMap(() => fetch('/api/users').then((r) => r.json()))
)
.subscribe((users) => console.log(users));
// concatMap: 순차 실행
click$
.pipe(
concatMap(() => fetch('/api/users').then((r) => r.json()))
)
.subscribe((users) => console.log(users));
세 오퍼레이터는 “바깥 스트림에 새 값이 왔는데 안쪽 작업이 아직 진행 중일 때” 무엇을 하느냐만 다릅니다. switchMap은 진행 중인 안쪽 구독을 해제하고 새 것으로 갈아탑니다. 가장 최근 요청의 결과만 의미 있는 검색·필터링에 맞습니다. mergeMap은 모두 동시에 실행하고 끝나는 순서대로 결과를 내보내므로, 응답 순서가 요청 순서와 다를 수 있습니다. 두 번째 인자로 동시 실행 수를 제한할 수 있습니다(mergeMap(fn, 3)). concatMap은 앞의 작업이 끝날 때까지 다음 값을 큐에 쌓아 두므로 순서가 보장되고, 저장·결제처럼 순서가 바뀌면 안 되는 쓰기 작업에 맞습니다. 여기에 진행 중이면 새 값을 무시하는 exhaustMap까지 알아 두면 “로그인 버튼 연타 방지”도 한 줄로 해결됩니다.
선택을 잘못했을 때의 증상도 알아 두면 좋습니다. 저장 요청에 switchMap을 쓰면 사용자가 빠르게 두 번 저장할 때 첫 번째 요청의 응답 처리가 취소되어, 서버에는 저장됐는데 화면은 갱신되지 않는 이상한 상태가 됩니다. 반대로 검색에 mergeMap을 쓰면 “reac” 결과가 “react” 결과보다 늦게 도착해 화면에 옛 검색 결과가 남는 레이스가 생깁니다.
또 하나, 위 예제처럼 switchMap 안에서 fetch Promise를 반환하면 이전 결과는 무시되지만 HTTP 요청 자체는 끝까지 진행됩니다. Promise에는 취소 개념이 없기 때문입니다. 서버 부하까지 줄이고 싶다면 rxjs/ajax나 Angular HttpClient처럼 구독 해제 시 요청을 실제로 중단(abort)하는 Observable 기반 HTTP를 쓰거나, new Observable에서 AbortController를 연결해야 합니다.
Subject로 직접 값 밀어 넣기
Subject는 “내가 next로 밀어 넣는 스트림”입니다. 싱크대에서 물 뿌리듯이 여러 구독자한테 뿌릴 수 있습니다. BehaviorSubject는 마지막 값을 들고 있어서 “지금 상태가 뭔지”를 대표로 쓰기 좋으며, ReplaySubject는 늦게 온 애한테도 최근 n개를 되감아 보여줘요.
import { Subject } from 'rxjs';
const subject$ = new Subject<number>();
subject$.subscribe((value) => console.log('A:', value));
subject$.subscribe((value) => console.log('B:', value));
subject$.next(1);
subject$.next(2);
// A: 1, B: 1, A: 2, B: 2
import { BehaviorSubject } from 'rxjs';
const subject$ = new BehaviorSubject<number>(0);
subject$.subscribe((value) => console.log('A:', value));
// A: 0
subject$.next(1);
subject$.next(2);
subject$.subscribe((value) => console.log('B:', value));
// B: 2
import { ReplaySubject } from 'rxjs';
const subject$ = new ReplaySubject<number>(2);
subject$.next(1);
subject$.next(2);
subject$.next(3);
subject$.subscribe((value) => console.log('A:', value));
// A: 2, A: 3
세 Subject의 차이는 “늦게 구독한 쪽이 무엇을 받느냐”입니다. 일반 Subject는 구독 이후의 값만 받아서, 구독 전에 next된 값은 영영 놓칩니다. 이벤트 버스처럼 과거 값이 의미 없는 곳에 맞습니다. BehaviorSubject는 초기값이 필수이고 항상 “현재 값”을 하나 갖고 있어서, 로그인 사용자·테마 같은 상태를 표현할 때 씁니다. .getValue()로 동기적으로 현재 값을 꺼낼 수도 있지만, 이것을 남용하면 반응형 흐름을 벗어난 명령형 코드가 되므로 가급적 구독으로 받는 편이 좋습니다.
Subject를 서비스 밖으로 그대로 노출하면 어느 컴포넌트든 next()를 호출할 수 있어 값이 어디서 바뀌었는지 추적하기 어려워집니다. 서비스 안에서는 private Subject로 두고 밖에는 subject$.asObservable()만 내보내는 것이 흔한 관례입니다.
에러 처리: catchError와 retry
에러는 스트림 안에서 catchError로 흡수하거나, retry로 몇 번 더 굴려보다가 실패 시 fallback 스트림으로 돌릴 수 있습니다. 다만 retry를 POST에 무작정 쓰면 멱등 아닌 요청이 두 번 갈 수 있으니, 재시도 정책은 꼭 같이 봐야 합니다.
import { of, throwError } from 'rxjs';
import { catchError, retry } from 'rxjs/operators';
const source$ = throwError(() => new Error('Error!'));
source$
.pipe(
retry(3),
catchError((error) => {
console.error('Error:', error);
return of('Fallback value');
})
)
.subscribe((value) => console.log(value));
retry(3)은 에러가 나면 소스에 다시 구독하는 방식으로 동작합니다. 그래서 소스가 Cold Observable(구독할 때마다 요청을 새로 보내는 ajax, HttpClient, defer)이어야 의미가 있습니다. from(promise)처럼 이미 실행된 Promise를 감싼 스트림에 retry를 걸면 같은 실패 결과를 세 번 다시 받을 뿐, 요청은 재전송되지 않습니다. 실무에서는 즉시 재시도보다 retry({ count: 3, delay: 1000 })처럼 간격을 두는 편이 일시적 장애에 강합니다.
catchError의 위치도 중요합니다. 위처럼 바깥 파이프라인의 끝에 두면 에러를 잡은 뒤 fallback 값을 내보내고 스트림 전체가 complete됩니다. 검색 자동완성처럼 계속 살아 있어야 하는 스트림에서 이렇게 하면, 요청이 한 번 실패한 뒤로는 입력을 해도 더 이상 아무 반응이 없습니다. 이런 경우에는 아래 예제처럼 switchMap 안쪽 요청 스트림에 catchError를 붙여 실패를 그 요청 하나로 가둬야 합니다.
실전: 검색 자동완성 만들기
검색 자동완성은 위 조합이 그대로 “이야기 한 편”입니다. 타이핑 → 값 선택 → debounce → 이전 쿼리랑 다를 때만 → switchMap으로 요청. 빈 쿼리면 빈 배열. 옛날에 if문으로 엮던 걸 그림으로 펼쳐 놓은 느낌입니다.
import { fromEvent, of } from 'rxjs';
import { debounceTime, distinctUntilChanged, switchMap, map } from 'rxjs/operators';
const input = document.getElementById('search') as HTMLInputElement;
const results = document.getElementById('results');
const search$ = fromEvent(input, 'input').pipe(
map((event) => (event.target as HTMLInputElement).value),
debounceTime(500),
distinctUntilChanged(),
switchMap((query) => {
if (!query) return of([]);
return fetch(`/api/search?q=${encodeURIComponent(query)}`).then((r) => r.json());
})
);
search$.subscribe((items) => {
results.innerHTML = items.map((item) => `<li>${item.name}</li>`).join('');
});
예제를 그대로 쓰기 전에 두 가지를 고쳐야 합니다. 첫째, 검색어는 encodeURIComponent로 인코딩해야 &나 #이 들어간 검색어가 쿼리 스트링을 깨뜨리지 않습니다. 둘째, 서버 응답을 innerHTML에 그대로 넣으면 item.name에 <img onerror=...> 같은 문자열이 섞였을 때 XSS가 됩니다. 실제 코드에서는 textContent로 요소를 만들거나 프레임워크의 템플릿 바인딩을 써야 합니다. 그리고 이 버전은 요청 하나가 실패하면 스트림 전체가 에러로 끝나므로, 아래 “Cold vs Hot” 절의 마지막 예제처럼 안쪽에 catchError를 붙인 형태가 운영용에 가깝습니다.
Angular와 함께 쓰기
Angular와는 원래 잘 맞아요. 컴포넌트가 죽을 때 async 파이프나 takeUntilDestroyed로 구독을 정리해 주는 흐름이 있어서, “스트림 쌓다가 누수”를 줄이기 좋습니다.
import { Component } from '@angular/core';
import { HttpClient } from '@angular/common/http';
import { debounceTime, distinctUntilChanged, switchMap } from 'rxjs/operators';
import { Subject } from 'rxjs';
@Component({
selector: 'app-search',
template: `
<input (input)="search$.next($event.target.value)" />
<ul>
<li *ngFor="let user of users$ | async">{{ user.name }}</li>
</ul>
`,
})
export class SearchComponent {
search$ = new Subject<string>();
users$ = this.search$.pipe(
debounceTime(500),
distinctUntilChanged(),
switchMap((query) => this.http.get(`/api/users?q=${query}`))
);
constructor(private http: HttpClient) {}
}
users$를 템플릿에서 async 파이프로 구독하므로 컴포넌트가 파괴될 때 Angular가 구독을 자동으로 해제합니다. 수동 subscribe()와 unsubscribe()를 쓰지 않는 것이 Angular에서 누수를 막는 가장 확실한 방법입니다. HttpClient는 구독 해제 시 XHR을 abort하므로 여기서의 switchMap은 이전 요청을 실제로 취소합니다.
엄격한 템플릿 타입 검사(strictTemplates)를 켠 프로젝트에서는 $event.target.value에서 Property 'value' does not exist on type 'EventTarget' 에러가 납니다. $event.target의 타입이 EventTarget | null이기 때문입니다. 템플릿 참조 변수(<input #q (input)="search$.next(q.value)">)를 쓰면 타입 문제 없이 값을 넘길 수 있습니다. *ngFor 역시 Angular 17 이후에는 내장 제어 흐름 @for가 권장됩니다.
Cold vs Hot Observable, 그리고 공유
Cold Observable은 구독할 때마다 안이 다시 돌아가요. 따라서 똑같은 걸 컴포넌트 둘이 각각 구독하면 요청이 두 번 나갈 수 있습니다. 그럴 땐 상위에서 한 번 구독해서 내려주거나, share / shareReplay로 “한 소스, 여럿이 share” 쪽을 생각합니다. shareReplay는 최근 n개를 물고 있으므로 설정 한 번 가져온 뒤 여기저기 쓰는 패턴엔 잘 맞는데, refCount를 어떻게 잡느냐에 따라 메모리 이야기가 갈려요.
스케줄러는 값이 “언제” 전달될지를 정합니다. 대부분의 오퍼레이터는 적절한 기본값을 쓰므로 직접 지정할 일은 드물지만, 애니메이션처럼 화면 갱신 주기에 맞춰야 할 때는 animationFrameScheduler를 씁니다. 테스트에서는 debounceTime(500) 같은 시간 기반 코드를 실제로 기다리면 느리고 불안정하므로, TestScheduler의 marble 테스트('a-b-c|' 같은 문자열로 시간 흐름을 표현)로 가상 시간을 돌려 검증합니다. Angular 테스트에서는 fakeAsync와 tick()도 같은 목적입니다.
switchMap이 안쪽을 정리해 준다고 해도, 직접 건 타이머·리스너는 finalize로 정리해 줘야 해요. HTTP는 AbortSignal이랑 엮으면 “이제 이건 버려”가 더 선명해집니다.
import { fromEvent, switchMap, EMPTY } from 'rxjs';
import { ajax } from 'rxjs/ajax';
import { catchError, debounceTime, map } from 'rxjs/operators';
const input = document.getElementById('q') as HTMLInputElement;
fromEvent(input, 'input').pipe(
map((e) => (e.target as HTMLInputElement).value),
debounceTime(300),
switchMap((q) =>
q.trim() === ''
? EMPTY
: ajax.getJSON(`/api/search?q=${encodeURIComponent(q)}`).pipe(
catchError(() => EMPTY)
)
)
).subscribe(render); // render: 결과를 화면에 그리는 함수
이 예제가 앞의 자동완성보다 나은 점은 세 가지입니다. ajax.getJSON은 구독 해제 시 요청을 abort하므로 switchMap이 진짜로 이전 요청을 취소합니다. 빈 검색어에는 EMPTY를 반환해 요청 자체를 보내지 않습니다. 그리고 catchError가 안쪽 요청에 붙어 있어서, 요청 하나가 실패해도 바깥 입력 스트림은 살아 있고 다음 입력에서 다시 검색이 동작합니다.
트러블슈팅 체크리스트
막막할 때는 이렇게 보면 돼요. 아무것도 안 온다—Cold인데 subscribe를 안 했거나 EMPTY/NEVER를 탔는지. 요청이 두 번—Cold를 둘이 따로 먹었는지, share가 필요한지. 메모리 경고—interval이나 DOM이 destroy될 때 안 끊겼는지, Angular면 takeUntilDestroyed 쪽. switch 썼는데도 옛 응답—Promise를 단순히 await만 하고 스트림으로 안 감쌌는지. 테스트만 터진다—가상 시간이 없어서. 표로 정리하던 걸 머릿속에서 이렇게 흘리면 됩니다.
마무리
Promise는 한 번, Observable은 0~무한, 끊을 수 있어서 그림이 달라요. 학습 곡선은 있습니다. 대신 시간과 이벤트가 겹치는 UI에 한 번 익혀 두면, 나중에 Angular 말고 다른 데서도 “아 이건 흐름이네”가 빨리 보이기 시작해요. React 쪽에선 단순히 rxjs만 설치해 써도 되는데, 팀이 hooks 중심이면 “왜 이 프로젝트만 우주기지냐”는 질문은 받을 수 있습니다. 그때는 정말 스트림이 필요한지, 아니면 Promise로 끊을 수 있는지 먼저 합의하는 게 좋고요.