1. 서브젝트의 특성
- 콜드 옵저버블 - 멜티캐스팅을 지원하지 않는 옵저버블
- 핫 옵저버블 - 멀티캐스팅을 지원하는 옵저버블
서브젝트란?
멀티캐스팅을 지원하기위해 옵저버이면서 옵저버블이라는 특성이 있다. 따라서, subscribe 함수를 호출했을 때 옵저버를 등록한다. 등록된 옵저버들은 서브젝트가 보내는 값이나 이벤트, 에러, 구독 완료 등의 정보를 받을 수 있다.
/* 옵저버블로 사용하는 서브젝트 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);
[코드 11-1] 옵저버블로 사용하는 서브젝트
서브젝트가 옵저버로서의 특성이 있으므로 각 옵저버로 값, 에러, 완료는 next, error, complete 함수를 호출해서 보낼 수 있다.
/* 옵저버로 사용하는 서브젝트 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);
subject.next(1);
subject.next(2);
subject.next(3);
[코드 11-2] 옵저버로 사용하는 서브젝트
실행 결과
observerA: 1
observerB: 1
observerC: 1
observerA: 2
observerB: 2
observerC: 2
observerA: 3
observerB: 3
observerC: 3
서브젝트의 next 함수를 호출할 때는 subscribe 함수로 등록한 옵저버들에게 값을 전파한다. 옵저버로서 error, complete 함수도 호출할 수 있다.
/* error 함수 호출 후 next 함수 호출 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);
subject.error(new Error('error!'));
subject.next(4);
subject.complete();
[코드 11-3] error 함수 호출 후 next 함수 호출
실행 결과
observerA: Error: error!
observerB: Error: error!
observerC: Error: error!
/* complete 함수 호출 후 next 함수 호출 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);
subject.complete();
subject.next(4);
subject.error(new Error('error!'));
[코드 11-4] complete 함수 호출 후 next 함수 호출
실행 결과
observerA: complete
observerB: complete
observerC: complete
2. 서브젝트와 옵저버블의 연결
/* interval 생성 함수를 이용하는 콜드 옵저버블 동작 */
const { interval } = require('rxjs');
const { take } = require('rxjs/operators');
const intervalSource$ = interval(500).pipe(take(5));
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.log(`observerA: ${e}`),
complete: () => console.log('observerA: complete'),
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.log(`observerB: ${e}`),
complete: () => console.log('observerB: complete'),
};
intervalSource$.subscribe(observerA);
setTimeout(() => {
intervalSource$.subscribe(observerB);
}, 2000);
[코드 11-5] interval 생성 함수를 이용하는 콜드 옵저버블 동작
실행 결과
observerA: 0
observerA: 1
observerB: 0
observerA: 2
observerB: 1
observerA: 3
observerB: 2
observerA: 4
observerA: complete
observerB: 3
observerB: 4
observerB: complete

subscribe 함수를 호출하는 각 옵저버블 구독이 따로 동작하며 매번 새로 구독하는 구조다.
서브벡트를 이용해서 어떻게 기존 옵저버가 멀티캐스팅되는 구조로 연결할 수 있을까? 서브젝트를 선언하고 서브젝트를 구독하는 subject.subscribe 가 필요하다.
/* 서브젝트를 새로 선언해 대체 */
const { interval, Subject } = require('rxjs');
const { take } = require('rxjs/operators');
const subject = new Subject();
const intervalSource$ = interval(500).pipe(take(5));
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.log(`observerA: ${e}`),
complete: () => console.log('observerA: complete'),
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.log(`observerB: ${e}`),
complete: () => console.log('observerB: complete'),
};
subject.subscribe(observerA);
setTimeout(() => {
subject.subscribe(observerB);
}, 2000);
[코드 11-6] 서브젝트를 새로 선언해 대체
[코드 11-5]를 기반으로 서브젝트를 새로 선언하고 intervalSource$ 대신 해당 자리를 subject 로 바꿨다.
다음으로 서브젝트에서 next, error, complete 함수로 값을 전달해야 한다.
/* intervalSource$ 를 구독해 서브젝트로 보내는 작업 추가 */
const { interval, Subject } = require('rxjs');
const { take } = require('rxjs/operators');
const subject = new Subject();
const intervalSource$ = interval(500).pipe(take(5));
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.log(`observerA: ${e}`),
complete: () => console.log('observerA: complete'),
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.log(`observerB: ${e}`),
complete: () => console.log('observerB: complete'),
};
subject.subscribe(observerA);
intervalSource$.subscribe({
next: x => subject.next(x),
error: e => subject.error(e),
complete: () => subject.complete()
});
setTimeout(() => {
subject.subscribe(observerB);
}, 2000);
[코드 11-7] intervalSource$를 구독해 서브젝트로 보내는 작업 추가
실행 결과
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerB: 3
observerA: 4
observerB: 4
observerA: complete
observerB: complete
intervalSource$ 구독을 시작하고 2초 후에 observerA를 구독한다는 것을 보여주어야 하므로 setTimeout 직전에 코드를 위치 시켰다.

서브젝트는 서브젝트 하나에서 스트림 하나가 여러 옵저버로 전파되는 구조이다.
/* 옵저버블에서 값을 바로 서브젝트로 전달 */
const { interval, Subject } = require('rxjs');
const { take } = require('rxjs/operators');
const subject = new Subject();
const intervalSource$ = interval(500).pipe(take(5));
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.log(`observerA: ${e}`),
complete: () => console.log('observerA: complete'),
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.log(`observerB: ${e}`),
complete: () => console.log('observerB: complete'),
};
subject.subscribe(observerA);
intervalSource$.subscribe(subject);
setTimeout(() => {
subject.subscribe(observerB);
}, 2000);
[코드 11-8] 옵저버블에서 값을 바로 서브젝트로 전달
실행 결과
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerB: 3
observerA: 4
observerB: 4
observerA: complete
observerB: complete
서브젝트는 옵저버이기도 하므로 옵저버블에서 값을 바로 서브젝트로 전달해줘서 함수 중복을 피할 수 있다. 즉, 서브젝트의 옵저버블 특성을 잘 활용하면 옵저버블과 연결하거나 직접 옵저버의 함수들을 호출해서 멀티캐스팅하려는 값, 이벤트, 에러, 완료에 관한 정보를 보낼 수 있다.
3. 서브젝트의 에러와 완료 처리
Subject 의 내부 함수
- next 함수 - subscribe 함수 호출 전 전달한 값은 이후 구독하는 옵저버로 전달하지 않는다.
- error 함수 - 호출 결과는 이후에 구독하는 옵저버에게도 전파한다.
- complete 함수 - 호출 결과는 이후에 구독하는 옵저버에게도 전파한다.
즉, 이미 해당 서브젝트에서 에러가 발생했거나 서브젝트 구독을 완료했다는 것을 알려준다.
/* 서브젝트의 에러 발생 상황을 전파 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.error('error');
subject.subscribe(observerA);
subject.subscribe(observerB);
[코드 11-9] 서브젝트의 에러 발생 상황을 전파
실행 결과
observerA: error
observerB: error
/* 서브젝트의 구독 완료 상황을 전파 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.complete();
subject.subscribe(observerA);
subject.subscribe(observerB);
[코드 11-10] 서브젝트의 구독 완료 상황을 전파
실행 결과
observerA: complete
observerB: complete
서브젝트는 unsubscribe 함수도 제공한다. unsubscribe 함수를 호출하면 아무 일도 일어나지 않은 것 같지만 이후 모든 옵저버블 대상으로 멀티케스팅할 수 없다. next, error, complete 함수를 호출할 때도 에러가 발생하며 멀티캐스팅하려고 특정 옵저버가 구독을 시도해도 에러가 발생한다. 이는 등록된 옵저버가 있는 배열을 null 로 만들며, 더 사용할 수 없는 서브젝트로 취급해 closed 플래그를 true 로 인식해 에러가 발생하는 것이다.
/* 서브젝트의 unsubscribe 함수의 동작 */
const { Subject } = require('rxjs');
const subject = new Subject();
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
subject.subscribe(observerA);
subject.subscribe(observerB);
subject.unsubscribe();
// subject 구독 해제 후 다시 구독한다.
subject.subscribe(observerC);
// 하나씩만 주석 처리를 해제한 후 코드를 실행한다.
// subject.next(1);
// subject.error('error');
// subject.complete();
[코드 11-11] 서브젝트의 unsubscribe 함수의 동작
실행 결과
[Error [ObjectUnsubscribedError]: object unsubscribed]
4. 서브젝트의 종류
- BehaviorSubject - 시간과 같은 연속인 값을 다루는 구조에 적합하다. 초기값이 있어 언제 구독해도 항상 값이 있다.
- ReplaySubject - 서브젝트를 생성할 때 인자로 설정한 수만큼 최근 전달받은 아이템을 갖고 있다가 다음 구독할 때 해당 수만큼 이벤트를 전달한다.
- AsyncSubject - 서브젝트 구독 완료 후 가장 마지막에 있는 아이템을 전달한다.
4-1. BehaviorSubject
BehaviorSubject 는 subscribe 함수를 호출하자마자 next 함수에서 전달 받을 수 있는 초기값이 있다. 생성할 때 초기값을 전달하고, 옵저버의 하수가 한 번도 호출되지 않아 아무 값도 전달받지 않는 다면 초기값을 그대로 사용한다. error 나 complete 함수가 호출되지 않았다면 subscribe 함수를 호출할 때마다 최근에 next 함수에서 전달받은 값을 준다. 즉, 초기값이 최근 값이든 subscribe 함수를 호출하자마자 전달받을 수 있는 값이 있는 서브젝트다.
/* BehaviorSubject 의 기본 동작 예 */
const { BehaviorSubject } = require('rxjs');
const behaviorSubject = new BehaviorSubject('초기값');
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const observerC = {
next: x => console.log(`observerC: ${x}`),
error: e => console.error(`observerC: ${e}`),
complete: () => console.log('observerC: complete')
};
behaviorSubject.subscribe(observerA);
behaviorSubject.next('값1');
behaviorSubject.subscribe(observerB);
behaviorSubject.next('값2');
behaviorSubject.subscribe(observerC);
behaviorSubject.next('값3');
behaviorSubject.next('값4');
behaviorSubject.next('값5');
[코드 11-12] BehaviorSubject 의 기본 동작 예
실행 결과
observerA: 초기값
observerA: 값1
observerB: 값1
observerA: 값2
observerB: 값2
observerC: 값2
observerA: 값3
observerB: 값3
observerC: 값3
observerA: 값4
observerB: 값4
observerC: 값4
observerA: 값5
observerB: 값5
observerC: 값5
BehaviorSubject 는 어떤 옵저버든 subcribe 함수로 호출할 때마다 바로 전달받을 수 있는 값이 있다. 이후에는 멀티캐스팅으로 구독하는 모든 옵저버가 next 함수로 값을 전달 받을 수 있는 구조다.
현재 값을 가져올 수 있는 value 또는 getValue 함수
BehaviorSubject 는 서브젝트 중 유일하게 항상 값을 갖고 있기 때문에, 현재 값을 subcribe 함수의 호출 없이 바로 전달 받을 수 있는 게터함수 getValue 를 제공한다. behaviorSubject.value 로 바도 접근할 수 있다.
/* BehaviorSubject 룰 이용하여 구현한 숫자 동작 */
const { BehaviorSubject, interval } = require('rxjs');
const { take, map } = require('rxjs/operators');
const behaviorSubject = new BehaviorSubject(0);
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
const incrementInterval$ = interval(1000).pipe(
take(5),
map(x => behaviorSubject.value + 1), // 최신 값에서 1 증가시킨 값으로 변환
// map(x => behaviorSubject.getValue() + 1)
);
// incrementInterval$ 를 behaviorSubject 에 연결하여 구독 시작
incrementInterval$.subscribe(behaviorSubject);
// observerA 바로 구독
behaviorSubject.subscribe(observerA);
// observerB 는 3.2초 후 구독해 가장 최신 값 3이 바로 나오는지 확인
setTimeout(() => behaviorSubject.subscribe(observerB), 3200);
[코드 11-13] BehaviorSubject 를 이용하여 구현한 숫자 동작
실행 결과
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerB: 3
observerA: 4
observerB: 4
observerA: 5
observerB: 5
observerA: complete
observerB: complete
4-2. ReplaySubject
ReplaySubject 는 next 함수로 지정한 개수만큼 연속해서 전달한 최신 값을 저장했다가 다음 구독 때 해당 개수만큼 옵저버로 발행한다. 그 후 멀티캐스팅되는 값을 발행하는 서브젝트다.
서브젝트를 생성할 때 연속해서 전달해야 하는 값 개수를 지정할 수 있다. 개수를 지정하지 않으면 메모리와 관련한 성능 문제가 발생한다.
/* ReplaySubject 의 기본 사용 예 */
const { ReplaySubject, interval } = require('rxjs');
const { take } = require('rxjs/operators');
const replaySubject = new ReplaySubject(3);
const interval$ = interval(500).pipe(take(8));
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
console.log('try replaySubject.subscribe(observerA)');
replaySubject.subscribe(observerA);
console.log('try interval$.subscribe(replaySubject)');
interval$.subscribe(replaySubject);
setTimeout(() => {
console.log('try replaySubject.subscribe(observerB) setTimeout 2600ms');
replaySubject.subscribe(observerB);
}, 2600)
[코드 11-14] ReplaySubject 의 기본 사용 예
실행 결과
try replaySubject.subscribe(observerA)
try interval$.subscribe(replaySubject)
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerA: 4
try replaySubject.subscribe(observerB) setTimeout 2600ms
observerB: 2
observerB: 3
observerB: 4
observerA: 5
observerB: 5
observerA: 6
observerB: 6
observerA: 7
observerB: 7
observerA: complete
observerB: complete
연속해서 전달할 값 개수는 3개로 지정했고 observerB 를 구독하면 가장 최근 저장한 값 3개인 2, 3, 4를 전달해 발행한다.
제한 없이 값을 전달하는 ReplaySubject
ReplaySubject 생성 부분에서 인자를 사용하지 않으면 제한 없이 연속해서 값을 전달할 수 있다.
/* 제한 없이 값을 전달하는 ReplaySubject 의 사용 예 */
const { ReplaySubject, interval } = require('rxjs');
const { take } = require('rxjs/operators');
const replaySubject = new ReplaySubject();
const interval$ = interval(500).pipe(take(8));
const observerA = {
next: x => console.log(`observerA: ${x}`),
error: e => console.error(`observerA: ${e}`),
complete: () => console.log('observerA: complete')
};
const observerB = {
next: x => console.log(`observerB: ${x}`),
error: e => console.error(`observerB: ${e}`),
complete: () => console.log('observerB: complete')
};
console.log('try replaySubject.subscribe(observerA)');
replaySubject.subscribe(observerA);
console.log('try interval$.subscribe(replaySubject)');
interval$.subscribe(replaySubject);
setTimeout(() => {
console.log('try replaySubject.subscribe(observerB), setTimeout 2600ms');
replaySubject.subscribe(observerB);
}, 2600)
[코드 11-15] 제한 없이 값을 전달하는 ReplaySubject 의 사용 예
실행 결과
try replaySubject.subscribe(observerA)
try interval$.subscribe(replaySubject)
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerA: 4
try replaySubject.subscribe(observerB), setTimeout 2600ms
observerB: 0
observerB: 1
observerB: 2
observerB: 3
observerB: 4
observerA: 5
observerB: 5
observerA: 6
observerB: 6
observerA: 7
observerB: 7
observerA: complete
observerB: complete
4-3. AsyncSubject
AsyncSubject 는 비동기로 실행한 동작이 완료되면 마지막 결과를 받는 역할을 한다. 즉, AsyncSubject 로 비동기 동작을 실행한 후 complete 함수를 호출하기 직전에 비동기 연산의 마지막 결과를 전달해야 한다.
/* 피보나치 수열 옵저버블에 AsyncSubject 연산자 사용 */
const { interval, AsyncSubject } = require('rxjs');
const { take, scan, pluck, tap } = require('rxjs/operators');
const asyncSubject = new AsyncSubject();
const period = 500;
const lastN = 8;
const fibonacci = n => interval(period).pipe(
take(n),
scan((acc, index) => acc ? { a: acc.b, b: acc.a + acc.b } : { a: 0, b: 1 }, null),
pluck('a'),
tap(n => console.log(`tap log: emitting ${n}`))
);
fibonacci(lastN).subscribe(asyncSubject);
asyncSubject.subscribe(result => console.log(`1st subscribe: ${result}`));
setTimeout(() => {
console.log('try 2nd subscribe');
asyncSubject.subscribe(result => console.log(`2nd subscribe: ${result}`));
}, period * lastN + 1000);
[코드 11-16] 피보나치 수열 옵저버블에 AsyncSubject 연산자 사용
실행 결과
tap log: emitting 0
tap log: emitting 1
tap log: emitting 1
tap log: emitting 2
tap log: emitting 3
tap log: emitting 5
tap log: emitting 8
tap log: emitting 13
1st subscribe: 13
try 2nd subscribe
2nd subscribe: 13
finbonacci 함수는 n을 인자로 사용해 해당 개수만큼 period 에서 설정한 시간마다 피보나치 수열을 발행하는 옵저버블을 리턴한다. AsyncSubject 를 테스트하기 위해 첫번째 구독은 실제 구독 완료될 때까지 기다린 후 마지막 결과만 나오는지 확인하는 용도이며, 두번째 구독은 첫번째 구독이 완료된 이후 구독하면 바로 결과가 나오는지 확인하는 용도다.
5. 마치며
서브젝트는 옵저버의 특성과 옵저버블의 특성을 동시에 갖고 멀티캐스팅을 지원한다. 옵저버블과 연결해 멀티캐스팅을 지원하지 않는 옵저버블의 발행 값을 여러 옵저버로 전파할 수 있다.
'RxJS' 카테고리의 다른 글
| 13장. 스케줄러 요약 (0) | 2024.08.08 |
|---|---|
| 12장. 멀티캐스팅 연산자 요약 (0) | 2024.08.07 |
| 10장. 에러 처리 요약 (0) | 2024.08.05 |
| 9장. 조건 연산자 요약 (0) | 2024.08.05 |
| 8장. 유틸리티 연산자 요약 (0) | 2024.08.02 |