1. 배경 지식
- 옵저버 패턴
- 명령형 프로그래밍
- 함수형 프로그램밍
1-1. 옵저버 패턴
옵저버 패턴의 기본 개념 - 관찰하는 역할의 옵저버 객체들을 서브젝트라는 객체에 등록한 후, 서브젝트 객체의 상태 변경이 일어나면 여기에 의존성 있는 옵저버들의 메서드를 호출해서 알리는 것

1-2. 자바스크립트 옵저버 패턴 예제
function func1() {
console.log('target click #1');
}
function func2() {
console.log('target click #2');
}
document.querySelector('#target').addEventListener('click', func1);
document.querySelector('#target').addEventListener('click', func2);
실행 결과
target click #1
target click #2
특정 이벤트를 관찰하다가 이벤트가 발생하면 리스너가 호출되는데 이렇게 콜백 함수가 동작하는 상황을 [옵저버가 이벤트를 구독(subscription)한다] 라고 한다.
구독 해제 코드
// 콜백 함수 하나를 구독 해제
document.querySelector('#target').removeEventListener('click', func1);
// 콜백 함수 모두 구독 해제
document.querySelector('#target').removeEventListener('click');
1-3. 함수형 프로그래밍과 순수 함수
RxJS는 비동기로 처리하는 여러 값을 명령형 프로그래밍이 아닌 함수형 프로그래밍 패러다임으로 다룬다.
함수형 프로그래밍에서 다루는 함수는 입력 데이터에 관한 출력 데이터가 항상 같아야 한다. 이러한 특성 때문에 함수형 프로그래밍의 함수는 다른 프로그래밍의 함수와 구분하려고 순수 함수라고도 한다. 순수 함수는 함수의 결과가 항상 보장되므로 잘 사용하면 값을 안전하게 관리하고 디버깅하기 쉽다.
2. 옵저버블
옵저버블이란?
옵저버 패턴을 기반으로 아래 두 가지 개념을 추가한 것
- RxJS의 complete 함수 - 더 이상 데이터가 없음을 알리는 onCompleted 메서드
- RxJS의 error 함수 - 에러가 발생 했음을 알리는 onError 메서드
RxJS는 옵저버 패턴을 적용한 옵저버블(Observable) 이라는 객체를 중심으로 동작
특정 객체를 관찰하는 옵저버에게 여러 이벤트나 값을 보내는 역할

프로미스는 객체를 생성하는 시점에, 옵저버블은 RxJS의 subscribe라는 함수를 호출하여 값이나 이벤트를 소비할 수 있는 시점에 실행되어 데이터를 생산한다.
- 생산자 역할 - 함수 / 이터레이터 / 프로미스 / 옵저버블은 값을 만들어내는 생산자 역할
- 소비자 역할 - function.call / iterator.next / promise.then / 옵저버블에 연결된 옵저버는 소비자 역할
프로미스와 옵저버블은 생산자가 능동적으로 데이터를 생산하면 알림을 받을 수 있는 콜백이 존재한다.

2-1. 옵저버블의 라이프사이클
- 옵저버블 생성(Creating Observables)
- 옵저버블 구독(Subscribing to Observables)
- 옵저버블 실행(Excuting the Observable)
- 옵저버블 구독 해제(Disposing Observables)
옵저버블의 생성
- require('rxjs')에서 불러온 Observable 클래스의 정적 함수 Observable.create 로 직접 생성
- require('rxjs')에서 불러온 range 나 of 로 추상화된 형태로 옵저버블 생성 - 옵저버블 인스턴스에 Observable.prototype 으로 연결된 pipe함수에 다양한 연산자를 인자로 사용하여 세로운 옵저버블 인스턴스를 생성할 수 있다.
구독과 실행
- 데이터를 전달할 콜백을 제공해 함수를 호출(구독)한 후 옵저버블에서 발행하는 값을 사용한다.
- subscribe함수를 이용
- 함수를 여러번 호출해도 해당 함수가 각각 독립적으로 동작한다.
/* 코드 2-2 옴저버블을 구독해 실행하는 자바스크립트 이벤트 처리 */
const {Observable} = require('rxjs');
const observableCreated$ = Observable.create(function (observer) {
for (let i = 1; i <= 10; i++) {
setTimeout(function () {
observer.next(i);
if (i === 10) {
observer.complete();
}
}, 300 * i);
}
});
observableCreated$.subscribe(
function next(item) {
console.log(`observerA: ${item}`);
},
function error(err) {
console.log(`observerA: ${err}`);
},
function complete() {
console.log('observerA: complete');
}
);
setTimeout(function () {
observableCreated$.subscribe(
function next(item) {
console.log(`observerB: ${item}`);
},
function error(err) {
console.log(`observerB: ${err}`);
},
function complete() {
console.log('observerB: complete');
}
);
}, 1350);
[코드 2-2]는 같은 값을 발행 할 수 있도록 멀티캐스팅하지 않는 상황
RxJS의 옵저버블은 멀티캐스팅이 안 될 때와 될 때를 모두 지원한다.
구독 해제
- unsubscribe함수를 사용하여 구독을 해제
주의사항
각 연산자에 넘겨준 함수에 return이나 break를 설정한다고 옵저버블의 동작이 중단되지 않는다. RxJS에서 실행되는 함수에서 return이나 break를 사용하면 해당 함수 안에서만 실행을 중단한다. 해당 옵저버블의 구독을 중단한다는 의미는 아니다.
2-2. 옵저버블 생성하고 실행하기
/* 옵저버블의 생성과 실행 */
const { Observable } = require('rxjs');
const observableCreated$ = Observable.create(function(observer) {
console.log('BEGIN Observable');
observer.next(1);
observer.next(2);
observer.complete();
console.log('END Observable');
});
observableCreated$.subscribe(
function next(item) {
console.log(item);
},
function error(e) {
},
function complete() {
console.log('complete');
}
);
실행 결과
BEGIN Observable
1
2
complete
END Observable
생성한 옵저버블은 옵저버에게 값을 전달하는 함수가 있지만 subscribe 함수가 호출되어야 옵저버블과 옵저버를 연결해 실행한다.

즉, 옵저버블 객체 생성 자체는 아무 일도 하지 않고 어떤 일을 해야 할지에 관한 정보만 있고, subscribe 함수를 호출해야 옵저버블이 옵저버에 데이터를 전달하며 동작을 실행한다.
/* next 와 complete 함수 */
const { Observable } = require('rxjs');
Observable.create(function (observer) {
console.log('BEGIN Observable');
observer.next(1);
observer.next(2);
observer.complete();
observer.next(3);
console.log('END Observable');
}).subscribe(
function next(item) {
console.log(item);
},
function error(e) {
},
function complete() {
console.log('complete');
}
);
실행 결과
BEGIN Observable
1
2
complete
END Observable
위의 결과를 보면 옵저버블 객체에서 subscribe 함수를 호출하면 옵저버블이 옵저버의 complete나 error함수를 호출할 때까지 next함수로 값을 발행한다. 하지만 next(3)함수를 호출하기 전에 complete함수가 이미 호출되었으므로 subcribe함수 안에 있는 next함수에 값 3은 발행되지 않는다.
2-3. 구독 객체 관리하기
옵저버는 next, error, complete라는 세 가지 함수로 구성된 객체다.
구독을 멈추게 하는 함수는 unsubscribe다.
/* 옵저버블 안 unsubscribe 함수 이용 */
const { Observable } = require('rxjs');
const observableCreated$ = Observable.create(function subscribe(observer) {
// intervalId 자원 추척
const intervalId = setInterval(function() {
observer.next('hi');
}, 1000);
// intervalId 자원을 해제하고 재배치하는 방법을 제공
return function unsubscribe() {
clearInterval(intervalId);
};
});
/* 옵저버블 구독 해제 */
const { interval } = require('rxjs');
const observable = interval(1000);
// 옵저버와 함께 subscribe 함수를 호출해 옵저버블 실행
const subscription = observable.subscribe(function (x) {
console.log(x);
});
// unsubscribe 함수로 구독 해제(바로 해제됨)
subscription.unsubscribe();
subscription 변수는 Subscription클래스의 인스턴스이다. 이 클래스는 unsubscribe함수 외에 add와 remove함수를 제공한다.
/* 여러개 Subscription 객체의 구독을 모두 해제 */
const { interval } = require('rxjs');
const observable1 = interval(400);
const observable2 = interval(300);
const subscription = observable1.subscribe(function (x) {
console.log(`first: ${x}`);
});
const childSubscription = observable2.subscribe(function (x) {
console.log(`second: ${x}`);
});
subscription.add(childSubscription);
setTimeout(function () {
// Subscription 객체와 하위에 있는 자식 Subscription 객체의 구독을 취소
subscription.unsubscribe();
}, 1000);
위와 같이 Subscription 객체 하나에 여러 Subscription 객체가 추가되었다면, 해당 Subscription객체의 unsubscribe함수를 호출하여 추가된 모든 Subscription객체의 구독을 해제할 수 있다.
3. 서브젝트
서브젝트란?
멀티캐스팅을 지원하는 객체다. 멀티 캐스팅을 지원한다는 것은 여러 옵저버가 이벤트 변경이나 값 전달을 관찰하도록 옵저버블을 구독 한 후, 실제 이벤트 변경이나 값 전달이 발생했을 때 이를 알린다는 뜻이다.
즉, 구독중인 모든 옵저버가 호출되어 같은 값을 전달 받는다는 뜻이다.
서브젝트는 옵저버블이면서 옵저버 역할도 한다.
(next, error, complete 함수를 호출해 같은 결과를 전달 받을 수 있다는 의미이다.)
/* 서브젝트 동작의 예 */
const { Subject } = require('rxjs');
const subject = new Subject();
subject.subscribe({
next: function (v) {
console.log(`observerA: ${v}`);
}
});
subject.subscribe({
next: function (v) {
console.log(`observerB: ${v}`);
}
});
subject.next(1);
subject.next(2);
subject 변수는 값을 보내주고 알려주는 형태의 옵저버블이자 옵저버고, observerA와 B는 이를 구독하는 옵저버다.
RxJS에 속한 Subject를 상속받는 것으로 BehaviorSubject, ReplaySubject, AsyncSubject가 있다.
4. 연산자
RxJS의 연산자는 기본적으로 함수 형태다.
- map연산자 - 값을 어떻게 변환할지 정하는 함수를 인자로 사용
- filter연산자 - 조건을 확인해 true와 false를 리턴하는 함수를 인자로 사용
연산자를 사용하려면 옵저버블을 생성해야 한다.
- Obesrvable.create - 옵저버블을 생성하는 일반적인 생성 함수
- of - 1개의 값만 발행하는 옵저버블을 생성하는 함수
- range - 특정 범위의 값을 순서대로 발행하는 옵저버블을 생성하는 함수
/* 생성 함수와 파이퍼블 연산자를 함께 사용 */
const { interval } = require('rxjs');
const { filter } = require('rxjs/operators');
let divisor = 2;
setInterval(function () {
divisor = (divisor + 1) % 10;
}, 500);
interval(700).pipe(
filter(function (value) {
return value % divisor === 0;
})
).subscribe((value) => console.log(value));
참고 사항 - filter연산자에 있는 값이 true인 것만 선별할 때 divisor값은 외부 참조할 수 있어 값이 바뀔 수 있기 때문에 불변 타입인 const로 선언해 값을 바꿀 수 없는 기본 타입을 참조하거나, 연산자에서 사용하는 함수 안이라는 유효 범위를 갖는 변수를 사용해야 안전하다.
4-1. 파이퍼블 연산자
파이퍼블 연산자는 생성 함수로 만들어진 옵저버블 인스턴스를 pipe함수 안에서 다를 수 있는 연산자이다.
사용법 - 1 : <옵저버블 인스턴스>.pipe(연산자1(), 연산자2(), ...)
사용법 - 2 : <옵저버블 인스턴스>.pipe(연산자1()).pipe(연산자2())...
연결한 파이퍼블 연산자는 각 연산자를 거치며 새로운 옵저버블 인스턴스를 리턴한다.
/* 생성 함수와 파이퍼블 연산자를 연결해 사용하기 */
const { range } = require('rxjs');
const { map, filter } = require('rxjs/operators');
range(1, 10).pipe(
filter(function (value) {
return value % 2 === 0;
}),
map(function (value) {
return value + 1;
})
);
먼저 range 함수가 옵저버블을 생성하고 pipe함수로 filter연산자 뒤에 map연산자를 연결해서 리턴한다.
참고 사항 - 파이퍼블 연산자를 연결해 새로운 옵저버블 인스턴스를 생성할 수 있는 이유는 filter연산자의 구현을 보면 알수 있다. (Observable,js에 있는 lift함수를 사용)
/* lift 함수의 구현 부분 */
lift(poerattor) {
const observable = new Observable();
observable.source = this;
observable.operator = operator;
return observable;
}
lift함수는 기존 옵저버블을 감싸서 새로운 옵저버블 한 단계 끌어 올려주는 역할을 한다.
/* subscribe 함수의 연산자 동작 실행 부분 */
const {toSubscriber} = require("rxjs/internal-compatibility");
const sink = toSubscriber(observerOrNext, error, complete);
if (operator) {
operator.call(sink, this.source);
} else {
sink.add(this._trySubscribe(sink));
}
연산자가 있는 지를 판단한 후 연산자가 있으면 연산자에 해당하는 동작을 실행한다.
연산자가 없을 때는 _trySubscribe에서 subscribe함수에 위치해있는 최종 옵저버(this._trySubscribe(sink))에 결과를 전달한다.
/* _trySubscribe 함수 구현 부분 */
_trySubscribe(sink) {
try {
return this._subscribe(sink);
} catch (err) {
sink.syncErrorThrown = true;
sink.syncErrorValue = err;
sink.error(err);
}
}
4-2. 배열과 비교해 본 옵저버블 연산자 예제
/* 옵저버블을 생성하고 변환 연산자 적용 */
const { Observable } = require('rxjs');
const { map } = require('rxjs/operators');
const observableCreated$ = Observable.create(function (observer) {
observer.next(1);
observer.next(2);
observer.complete();
});
observableCreated$.pipe(
map(function (value) {
return value * 2;
})
).subscribe(function next(item) {
console.log(item);
});
/* array.map 연산자 - 1 */
console.log([1, 2].map(function (value) {
return value * 2;
}));
/* array.map 연산자 - 2 */
console.log([1, 2]
.map(function (value) {
return value * 2;
})
.map(function (value) {
return value + 1;
})
.map(function (value) {
return value * 3;
})
);
실행 결과
2
4
[ 2, 4 ]
[ 9, 15 ]
배열과 옵저버블에서 map연산자의 차이
- 옵저버블 map - 실제 각각의 값 처리는 구독하는 시점에 한다.
- 배열 map - 연산자를 호출할 때마다 새로운 배열을 만든다.
즉, 배열 map은 연산자 수에 비례해 배열을 3번 생성하므로 메모리 공간을 차지해 가비지 컬렉터에서 제거해야하는 문제가 있어 옵저버블이 성능상 장점이 있다.
/* 옵저버블에서 map 연산자를 여러번 호출 */
const { Observable } = require('rxjs');
const { map, toArray } = require('rxjs/operators');
const observableCreated$ = Observable.create(function (observer) {
console.log('Observable BEGIN');
const arr = [1, 2];
for (let i = 0; i < arr.length; i++) {
console.log(`current array: arr[${i}]`);
observer.next(arr[i]);
}
console.log('BEFORE complete');
observer.complete();
console.log('Observable END');
});
function logAndGet(original, value) {
console.log(`original: ${original}, map value: ${value}`);
return value;
}
observableCreated$.pipe(
map(function (value) {
return logAndGet(value, value * 2);
}),
map(function (value) {
return logAndGet(value, value + 1);
}),
map(function (value) {
return logAndGet(value, value * 3);
}),
toArray()
).subscribe(function(arr) {
console.log(arr);
});
실행 결과
Observable BEGIN
current array: arr[0]
original: 1, map value: 2
original: 2, map value: 3
original: 3, map value: 9
current array: arr[1]
original: 2, map value: 4
original: 4, map value: 5
original: 5, map value: 15
BEFORE complete
[ 9, 15 ]
Observable END
map연산자로 값을 감쌀 때마다 새로운 옵저버블 객체만 생성하고 구독할 때 까지 실행되지 않으므로 배열이 생성될 때처럼 실제 연산자가 동작하지 않는다. 구독하는 순간 Observable BEGIN부터 ObservableEND까지의 모든 동작이 한 번에 실행된다.
/* 위 코드의 동작 방식 설명 */
//subscribe 호출 시 toArray에서 필요한 array 생성
const array = [];
//observableCreated 안 for 문에서 observer.next(arr[i]) 를 호출 할 때
const aInput = arr[i];
//observableCreated 안 함수에서 사용하는 시작 옵저저
observerA.next(aInput);
const bInput = aInput * 2;
observerB.next(bInput);
const cInput = bInput + 1;
observerC.next(cInput);
const dInput = cInput * 3;
observerD.next(dInput);
array.put(dInput);
//observableCreated 안에서 observer.complete()을 호출할 때, toArray()에서 실행되는 동작
//subscribe 안에 있는 마지막 옵저버이기도 하다.
observerE.next(array);
중요 포인트 - 구독하는 시점까지 실행을 미룰 수 있으므로 지연 실행이 가능하다는 장점이 있다.
5. 스케줄러
스케줄러란?
옵저버가 옵저버블을 구독할 때 어떤 순서로 어떻게 (동기/비동기 등) 실행 할지 실행 컨텍스트를 관리하는 역할의 자료구조이다.
- 비동기 방식 - setTimeout, setInterval, 마이크로 큐를 이용해 실행하는 asapScheduler, asyncScheduler
- 동기 방식 - 트램폴린 방식으로 큐를 사용하는 queueScheduler
6. 마블 다이어그램
마블 다이어그램은 연산자를 쉽게 이해하는 데 도움을 주고, 개발자가 생각하는 흐름을 그림으로 도식화하는 데 유용하다.

- 왼쪽에서 오른쪽으로 향하는 가로 줄은 시간에 따라 next 함수에서 발행하는 값들을 표시한다. 입력(input) 옵저버블이라고도 한다.
- 원 모양이고 그 안에 실제 값이 표시되어 있다. 이렇게 가로줄이 하나의 옵저버블이다.
- 제일 마지막에 있는 수직선( | )은 구독 완료(complete 함수 호출)를 의미한다.
- 가운데 이름이 있는 박스는 연산자를 가리킨다.
- 가운데 연산자 아래 있는 옵저버블은 연산자를 이용해서 생성한 옵저버블이다. 출력(output) 옵저버블이라고도 한다. 각 값들은 색깔이나 숫자로 구분해서 어떤 입력값으로 어떤 발행 값을 만들었는지 구분할 수 있게 했다. 원, 세모, 네모를 사용해서 값을 구분하기도 한다.
- X 표시는 에러가 발생해 옵저버블 실행을 종료했음을 의미한다.
7. 프로미스와 함께 본 옵저버 콜백 및 에러 처리
/* 정상적인 프로미스의 값 처리 방식 */
const promise = new Promise((resolve, reject) => {
resolve(1);
});
promise.then(function (value) {
console.log(value);
});
실행 결과
1
프로미스 생성자에서 사용하는 함수에 resolve의 1이라는 값을 인자로 설정하는 방식으로 값 1을 전달할 수 있다.
/* 프로미스의 에러 처리 방식 */
const promise = new Promise((resolve, reject) => {
reject(new Error('에러 발생'));
});
promise.then(
function (value) {
console.log(value);
},
function (error) {
console.error(error);
}
);
실행 결과
Error: 에러 발생
at C:\RxJS\chapter2\2-18.js:3:10
at new Promise (<anonymous>)
at Object.<anonymous> (C:\RxJS\chapter2\2-18.js:2:17)
at Module._compile (node:internal/modules/cjs/loader:1376:14)
at Module._extensions..js (node:internal/modules/cjs/loader:1435:10)
at Module.load (node:internal/modules/cjs/loader:1207:32)
at Module._load (node:internal/modules/cjs/loader:1023:12)
at Function.executeUserEntryPoint [as runMain] (node:internal/modules/run_main:135:12)
at node:internal/main/run_main_module:28:49
두번째 인자인 reject는 에러를 처리하므로 값을 전달하면 에러가 발생한다.
프로미스는 1개의 값을 취급하므로 resolve나 reject 둘 중 하나를 먼저 사용하면 then함수로 둘 중 하나의 결과를 받을 수 있다.
RxJS는 옵저버블에서 여러 개의 값을 취급한다. 옵저버에 있는 next함수로 값을 전달할 수 있고, error함수로 에러를 전달 할 수 있다. 에러를 전달하면 해당 옵저버블은 구독을 종료한다.
에러 없이 실행되면 complete함수를 호출해 구독을 종료한다.
/* 옵저버블을 이용한 값 전달과 에러 처리 */
const { Observable } = require('rxjs');
const observableCreated$ = Observable.create(function(observer) {
try {
observer.next(1);
observer.next(2);
throw("throw err test");
} catch (err) {
observer.error(err);
} finally {
observer.complete();
}
});
observableCreated$.subscribe(
function next(item) {
console.log(item);
},
function error(err) {
console.error('error: ' + err);
},
function complete() {
console.log('complete');
}
);
실행 결과
1
2
error: throw err test
- next : 다음에 전달할 값 또는 이벤트를 발행한다.
- error : 에러나 예외가 발생하면 이를 전달 받는다. 구독 종료
- complete : 정상적으로 옵저버블 구독을 완료하면 호출한다. 구독 종료
'RxJS' 카테고리의 다른 글
| 6장. 조합 연산자 요약 (0) | 2024.08.01 |
|---|---|
| 5장. 변환 연산자 요약 (0) | 2024.07.31 |
| 4장. 필터링 연산자 요약 (0) | 2024.07.29 |
| 3장. 생성 함수 요약 (0) | 2024.07.25 |
| 1장. RxJS 소개와 개발 환경 구축 요약 (1) | 2024.07.10 |