#include <stdio.h>
int add(int a, int b) {
return a + b;
}
int main() {
printf("%d students.\n",
add(printf("Good "), printf("morning ")));
return 0;
}
옵저버가 옵저버블을 구독할 때 값을 전달받는 순서와 실행 컨텍스트를 관리하는 역할을 하는 자료구조다.
1. 이벤트 루프와 RxJS의 스케줄러 개념
RxJS는 이벤트 루프 구조에서 로직 처리를 미루는 스케줄려를 구현하려고 플랫폼 환경(브라우저 또는 Node.js 환경)에 따라 제공하는 API 를 적절하게 활용한다.
주요 API
setTimeout - 특정 시간 뒤로 로직 실행을 미룸
setInterval - 특정 시간마다 반복해서 로직을 실행
setImmediate - 마이크로소프트 계열 브라우저와 Node.js 에서 현재 이벤트 루프 주기 끝에 로직을 실행
process.nextTick - Node.js 에서 이벤트 루프와 관계없이 무조건 현재 작업이 완료된 직후 로직을 실행
window.requestAnimationFrame - 브라우저에서 프레임이 끊기지 않도록 각 프레임마다 로직을 실행
RxJS 에서 제공하는 스케줄러
시간 기반의 스케줄러
비동기 처리를 위한 스케줄러
브라우저의 애니메이션 프레임 손실을 막는 스케줄러
RxJS 공식문서 상 스케줄러의 구성 요소
자료 구조 - 작업물을 우선 순위나 다른 기준에 따라서 저장하고 큐잉한다.
실행 컨텍스트 - 작업(task)을 실행하는 때와 위치를 가리킨다.
(가상) 클락 - now 함수라는 스케줄러의 시간을 가리키는 게터 함수를 제공한다. 특정 스케줄러에 스케줄한 작업들은 클락으로 설정한 시간에 맞춰 동작한다.
스케줄러 사용 방법
subscribeOn, observeOn 연산자나 스케줄러를 인자로 사용하는 연산자를 사용하는 방법
직접 연산자를 구현할 때 스케줄러에 있는 schedule 함수를 호출하는 방법
대표적인 스케줄러
AsyncScheduler - 일정 시간 이후에 실행되도록 만든다.
AsapScheduler - 비동기 동작을 최대한 빨리 실행하게 만든다.
QueueScheduler - 내부에 큐(queue)를 두고 작업을 넣어 동기로 실행하게 만든다.
기타 - 애니메이션과 테스트 코드에 사용하는 스케줄러
2. 스케줄러 구조
스케줄러는 schedule 함수를 호출해서 동작한다. schedule 함수는 해당 스케줄러와 매칭하는 액션 객체를 생성해 해당 액션을 실행한다.
/* Scheduler 클래스의 구현 코드 일부 */
constructor(SchedulerAction, now = Scheduler.now) {
// ...생략
}
schedule(work, delay = 0, state) {
return new this.SchedulerAction(this, work).schedule(state, delay);
}
[코드 13-1] Scheduler 클래스의 구현 코드 일부
액션은 Action 클래스를 상속받는데 내부적으로 동작해야 하는 장업(work) 함수를 인자로 사용한다. schedule 함수를 호출할 때 이 작업이 실행해야 할 상태 값인 state 를 전달 받는다.
스케줄러는 액션에 상태값을 전달하는 역할을 하고, 액션은 스케줄러의 작업 단위이다.
/* AsyncScheduler 의 액션 생성 */
import {AsyncScheduler} from "rxjs/internal/scheduler/AsyncScheduler";
import {AsyncAction} from "rxjs/internal/scheduler/AsyncAction";
export const async = new AsyncScheduler(AsyncAction);
[코드 13-2] AsyncScheduler 의 액션 생성
3. 대표 스케줄러
3-1. AsyncScheduler
AsyncScheduler 는 대표 스케줄러들이 상속받는 부모 스케줄러다. 각 작업 단위로 보면 setTimeout 함수처럼 일회성으로 일정 시간 후 정의한 작업을 동작시키는 스케줄러다.
AsyncScheduler 는 내부에 setInterval 함수를 두고 일정 간격마다 요청이 오는 작업을 해당 스케줄러를 사용 완료할 때까지 처리한다.
AsyncScheduler 는 delay 를 인자로 사용해 일정 시간 후 작업을 처리한다. 이 스케줄러를 상속받는 다은 스케줄러는 schedule 함수에 delay 를 사용했을 때 부모인 AsyncScheduler 를 이용해 작업을 처리한다.
asyncScheduler 에 schdule 함수로 work 함수를 사용하면 AsyncAction 인스턴스에서 실행되며 this 가 가리키는 인스턴스는 AsyncAction 이 된다. 따라서 work 함수 안에 있는 selfAction 은 AsyncAction 인스턴스다.
[코드 13-3]에서 처음 호출하는 schedule 함수는 스케줄러에서 호출하는 함수이고, 그 안에서 호출하는 schedule 함수는 액션에서 호출하는 함수이다.
/* 스케줄러에서 schedule 함수를 호출할 때마다 새로 생성하는 액션 */
schedule(work, delay = 0, state) {
return new this.SchedulerAction(this, work).schedule(state, delay);
}
[코드 13-4] 스케줄러에서 schedule 함수를 호출할 때마다 새로 생성하는 액션
스케줄러에서 schedule 함수를 호출할 때마다 해당 스케줄러의 액션 객체를 새로 생성해준 후 work 함수를 실행하고, 액션에서 호출하는 schedule 함수는 액션 객체 안에서 별개의 작업을 실행하는 역할이다.
3-2. AsapScheduler
AsapScheduler 는 각 플랫폼에 맞게 동기로 작업을 처리한 후 가능하면 빠르게 비동기로 작업을 처리하는 스케줄러다.
현재 이벤트 처리의 끝이나 현재 실행 로직 다음에 실행해야 할 이벤트 처리보다 더 빠르게 처리해야 하는 작업이 있을 때 사용
AsapScheduler 의 구현 원리 AsapScheduler 는 setImmeediate 함수를 호출한 후 actions 배열에 있는 액션을 매번 꺼내 비동기 동작을 한다. 그리고 work 함수를 호출할 때 해당 액션의 상태 값 (state) 을 전달해 동작을 실행한다.
0보다 큰 delay 값이 있으면 super 를 이용해 부모인 AsyncAction 의 동작을 호출한다. 그렇지 않으면 액션(this) 자체를 actions 배열에 푸시한다.
/* 스케줄러의 actions 배열 사용 방식 (AsapScheduler 의 flush 메서드) */
import {AsyncScheduler} from "rxjs/internal/scheduler/AsyncScheduler";
export class AsapScheduler extends AsyncScheduler {
flush(action) {
this.active = true;
this.scheduled = undefined;
const {actions} = this;
// 생략...
action = action || actions.shift();
do {
if (error = action.execute(action.state, action.delay)) {
break;
}
} while (++index < count && (action = actions.shift()));
// 생략...
}
}
[코드 13-6] 스케줄러의 actions 배열 사용 방식 (AsapScheduler 의 flush 메서드)
actions 에서 하나하나 값을 꺼내 동작할 때는 execute 함수를 호출한다. 이 때 state와 함께 work 함수를 호출한다.
/* AsapScheduler 를 이용한 동기 및 비동기 처리 예 */
const { of, asapScheduler } = require("rxjs");
console.log('start');
of(1, 2, 3, asapScheduler).subscribe(x => console.log(x));
console.log(`actions length: ${asapScheduler.actions.length}`);
console.log('end');
[코드 13-7] AsapScheduler 를 이용한 동기 및 비동기 처리 예
실행 결과
start
actions length: 1
end
1
2
3
start와 end가 동기로 먼저 실행되고, 1부터 3까지는 비동기로 한 번에 실행된다. 또한, 마이크로 큐에서 1개의 액션을 처리하려고 actions 배열에 1개의 액션을 추가했다. 이 1개의 액션은 인자로 나열된 1부터 3까지의 값을 하나하나 꺼내 전달하는 역할을 한다.
/* ArrayObservable 의 구현 코드 일부 */
static of(...array) {
// ...생략
if (len > 1) {
return new ArrayObservable(array, scheduler);
}
// ...생략
}
static dispatch(state) {
const {array, index, count, subscriber} = state;
if (index >= count) {
subscriber.complete();
return;
}
subscriber.next(array[index]);
if (subscriber.closed) {
return;
}
state.index = index + 1;
this.schedule(state);
}
_subscribe(subscriber) {
// ...생략
if (scheduler) {
return scheduler.schedule(ArrayObservable.dispatch, 0, {array, index, count, subscriber});
}
// ...생략
}
[코드 13-8] ArrayObservable 의 구현 코드 일부
of 함수 - 내부 array 에 나열된 값을 담아 실행
dispatch - work 함수이며, 상태 값으로 전달되는 객체에는 array, index, count, subscriber 가 있음
/* AsapScheduler 의 재귀 호출 */
const { asapScheduler } = require("rxjs");
console.log('start');
asapScheduler.schedule(function work(value) {
value = value || 1;
console.log(value);
var selfAction = this;
if (value < 3) {
selfAction.schedule(value + 1);
}
});
console.log(`actions length: ${asapScheduler.actions.length}`);
console.log('end');
[코드 13-9] AsapScheduler 의 재귀 호출
실행 결과
start
actions length: 1
end
1
2
3
3-3. QueueScheduler
QueueScheduler 는 동기 방식의 스케줄러다. actions 배열을 반복 실행하며 먼저 들어온 값을 먼저 사용하는 큐 자료구조를 사용한다.
/* QueueScheduler 의 구현 코드 일부 */
// AsyncScheduler 를 상속받을 뿐 구현은 없음
import {AsyncScheduler} from "rxjs/internal/scheduler/AsyncScheduler";
export class QueueScheduler extends AsyncScheduler { }
[코드 13-10] QueueScheduler 의 구현 코드 일부
QueueScheduler 는 동기 방식이므로 반복 실행 시작 전 플래그를 표시하고, 반복 실행 중이면 actions 배열에 푸시만 한다. 즉, actions 배열에 아직 실행해야 할 동작이 남아 있으면 배열 요소를 모두 실행할 때까지 반복해서 동작한다.
/* QueueAction 의 schedule 함수 구현 코드 */
// request, recycle 은 호출하지 않고 실행만 한다.
schedule(state, delay = 0) {
if (delay > 0) { // delay 값이 0보다 크면 AsapScheduler 를 상속받아 실행
return super.schedule(state, delay);
}
this.delay = delay;
this.state = state;
this.scheduler.flush(this); // AsyncScheduler 의 flush 메서드 호출
return this;
}
execute(state, delay) {
return (delay > 0 || this.closed) ?
super.execute(state, delay) : // delay == 0 이고 !this.closed 이므로 실행함.
this._execute(state, delay);
}
[코드 13-11] QueueAction 의 schedule 함수 구현 코드
스케줄러의 flush 메소드를 동기 방식으로 실행한다.
/* AsyncScheduler 의 flush 메소드 구현 코드 */
flush(action) {
const { actions } = this;
if (this.active) {
actions.push(action);
return;
}
let error;
this.active = true;
do {
if (error = action.execute(action.state, action.delay)) {
break;
}
} while (action = actions.shift()); // 스케줄러 큐 모두 사용
this.active = false;
if (error) {
while (action = actions.shift()) {
action.unsubscribe();
}
throw error;
}
}
[코드 13-12] AsyncScheduler 의 flush 메소드 구현 코드
flush 함수 안에서 스케줄러의 active 플래그를 반복 실행 시작 전후에 설정한다.
/* AsyncAction 의 _execute 함수 구현 코드 */
_execute(state, delay) {
let errored = false;
let errorValue = undefined;
try {
this.work(state);
} catch (e) {
errored = true;
errorValue = !!e && e || new Error(e);
}
if (errored) {
this.unsubscribe();
return errorValue;
}
}
[코드 13-13] AsyncAction 의 _execute 함수 구현 코드
QueueScheduler 사용 - 연산자 안에서 동기 방식 및 콜스택이 아닌 반복문으로 꼬리 재귀를 호출해야 하거나 큐에 넣어 순서를 맞춰야 할 때 사용하면 좋다.
두번째 인자 null 은 delay 값을 지정하지 않겠다는 의미이다. 이유는 QueueScheduler 도 delay 값을 지정하면 AsyncScheduler 를 상속받아 동작하므로 null 로 지정했다. 세번째 인자는 초기값을 넣었다.index 값을 1씩 증가시켜 재귀로 a, b 의 값을 누적 시킨다.
큐에서 하나식 꺼내서 순차적으로 동작함을 확인할 수 있다.
4. 스케줄러에서 사용하는 연산자
4-1. subscribeOn 연산자
subscribeOn 연산자는 구독하는 옵저버블 자체를 인자로 사용할 스케줄러로 바꿔준다.
before subscribe
BEGIN source
END source
after subscribe
1
2
3
'after subscribe' 출력까지는 동기 방식으로 실행되고 observeOn 연산자 다음으로 생성되는 옵저버블은 1초 후 스케줄러를 이용해서 실행된다. 즉, subscribe 함수 안에 있는 next 함수의 동작이 스케줄러의 영향을 받아 1부터 3까지 출력만 1초 후 비동기로 처리한다.
observeOn 연산자 다음에 바로 subscribe 함수를 호출하지 않고 다른 연산자를 추가했어도 그 다음에 추가하는 연산자부터는 스케줄러를 이용해 옵저버블을 실행한다.
subjectOrSubjectFactory - 소스 옵저버블 요소 순서로 서브젝트 팩토리 함수를 실행한다.
selector? - 선택자 함수로 소스 옵저버블을 여러 번 구독하지 않고 서브젝트를 이용해 소스 옵저버블을 필요할 때마다 사용할 수 있다.
2-1. multicast 연산자의 connect 함수로 서브젝트와 연결
multicast 연산자는 서브젝트를 생성하는 팩토리 함수나 서브젝트 자체를 인자로 사용한다.
connect 함수를 호출하면 소스 옵저버블에서 값을 발행하여 해당 서브젝트로 전달하고 해당 옵저버블을 구독하도록 등록된 옵저버들은 서브젝트로 같은 값을 전달받을 수 있다.
즉, 서브젝트를 직접 제공해서 커넥터블 옵저버블을 만든다.
/* connect 함수로 서브젝트와 연결 */
const { interval, Subject } = require('rxjs');
const { take, multicast } = require('rxjs/operators');
const sourceObservable$ = interval(500).pipe(take(5));
const multi = sourceObservable$.pipe(multicast(() => new Subject()));
// 첫번째 인자로 사용하는 팩토리 함수에서 리턴한 서브젝트를 구독하는 부분
const subscriberOne = multi.subscribe(val => console.log(val));
const subscriberTwo = multi.subscribe(val => console.log(val));
// 소스 옵저버블이 서브젝트를 구독하는 부분
multi.connect();
[코드 12-3] connect 함수로 서브젝트와 연결
실행 결과
0
0
1
1
2
2
3
3
4
4
multi.connect() 이 부분을 주석 처리하면 아무 일도 일어나지 않는다. 주석 처리하지 않으면 서브젝트가 해당 소스 옵저버블을 구독하는 것을 확인할 수 있다.
/* multicast 연산자의 서브젝트 구독 확인 */
const { interval, Subject } = require('rxjs');
const { take, multicast } = require('rxjs/operators');
const subject = new Subject();
const sourceObservable$ = interval(500).pipe(take(5));
const multi = sourceObservable$.pipe(multicast(() => subject));
// 다음 주적 처리한 코드를 사용해도 된다.
// const multi = sourceObservable$.pipe(multicast(subject));
// 첫번째 인자로 사용하는 팩토리 함수에서 리턴한 서브젝트를 구독하는 부분
const subscriberOne = multi.subscribe(val => console.log(val));
const subscriberTwo = multi.subscribe(val => console.log(val));
// 소스 옵저버블이 서브젝트를 구독하는 부분
subject.next(1);
[코드 12-4] multicast 연산자의 서브젝트 구독 확인
실행 결과
1
1
multi 를 두본 구독했지만 서브젝트를 구독한 것과 같은 효과가 있다. connect 함수를 호출하지 않았으므로 sourceObservable$ 은 동작하지 않았다. 하지만, multicast 연산자에서 사용하는 서브젝트에 next 함수로 1을 전달하면 1을 구독한 두 옵저버블의 발행 값을 출력한다. 이는 multicast 연산자로 만든 옵저버블 구독이 서브젝트를 구독하는 것과 같다는 뜻이다.
multicast 연산자는 연산자를 호출하는 쪽에서 서브젝트 팩토리 함수까지 제공하므로 서브젝트와의 의존성이 생긴다.
2-2. multicast 연산자의 선택자 함수
multicast 연산자는 두번째 인자로 선택자 함수를 사용한다. 선택자 함수 사용 시 multicast 연산자가 리턴하는 옵버버블이 커넥터블 옵저버블로 변환되지 않고 다른 방식으로 멀티캐스팅한다. 멀티캐스팅을 하지만 connect 함수를 제공하지 않고 동작하는 것이다.
주의할 점 - 구독할 때마다 팩토리 함수를 호출한다.
/* 같은 옵저버블을 두 번 구독할 때 multicast 연산자를 사용 안함 */
const { interval, zip, timer, Subject } = require('rxjs');
const { take, mergeMap, tap } = require('rxjs/operators');
interval(1500).pipe(
take(6)
).subscribe(x => console.log(`${(x + 1) * 1500}ms elapsed`));
const sourceObservable$ = interval(1500).pipe(
take(5),
tap(x => console.log(`tap ${x}`))
);
zip(sourceObservable$, sourceObservable$, (a, b) => a + ',' + b)
.subscribe(val => console.log("value : " + val));
[코드 12-5] 같은 옵저버블을 두 번 구독 할 때 multicast 연산자를 사용 안함
실행 결과
1500ms elapsed
tap 0
tap 0
value : 0,0
3000ms elapsed
tap 1
tap 1
value : 1,1
4500ms elapsed
tap 2
tap 2
value : 2,2
6000ms elapsed
tap 3
tap 3
value : 3,3
7500ms elapsed
tap 4
tap 4
value : 4,4
9000ms elapsed
1.5초마다 값을 발행하는 소스 옵저버블을 zip 연산자로 두 번 합해서 구독했다. 같은 소스 옵저버블(콜드 옵저버블)에서 발행한 값을 zip 연산자에 전달했더라도 각각 따로 동작한다.
'tap 숫자' 형식의 메세지가 두 번 출력된다. 1.5초마다 값 각각을 새로 발행한다는 것을 알 수 있다.
/* 멀티캐스팅할 때 서브젝트 안에 있는 선택자 함수 이용 */
const { interval, timer, zip, Subject } = require('rxjs');
const { take, tap, multicast, mergeMap } = require('rxjs/operators');
interval(1500).pipe(
take(6)
).subscribe(x => console.log(`${(x + 1) * 1500}ms elapsed`));
const sourceObservable$ = interval(1500).pipe(
take(5),
tap(x => console.log(`tap ${x}`))
);
const multi = sourceObservable$.pipe(
multicast(
() => new Subject(),
subject => zip(subject, subject, (a, b) => a + ',' + b)
)
);
multi.subscribe(val => console.log("value : " + val));
[코드 12-6] 멀티캐스팅할 때 서브젝트 안에 있는 선택자 함수 이용
실행 결과
1500ms elapsed
tap 0
value : 0,0
3000ms elapsed
tap 1
value : 1,1
4500ms elapsed
tap 2
value : 2,2
6000ms elapsed
tap 3
value : 3,3
7500ms elapsed
tap 4
value : 4,4
9000ms elapsed
소스 옵저버블 대신 소스 옵저버블과 연결된 서브젝트를 사용하므로 소스 옵저버블을 한 번만 구독한다. 때문에, 소스 옵저버블을 한 번만 구독한 후 발행한 값을 선택자 함수에서 제공하는 스트림을 거쳐서 출력한다.
tap 0
a: 0
b: 0
tap 1
a: 1
b: 1
tap 2
a: 2
b: 2
tap 3
a: 3
b: 3
tap 4
a: 4
b: 4
timeout
tap 0
tap 1
tap 2
tap 3
tap 4
tap 연산자를 호출할 때 a, b 로 멀티캐스팅 되지만 다임아웃 후에는 tap 연산자만 호출하고 멀티캐스팅 되지 않는다.
실행 순서
a, b 를 출력하는 두 옵저버블에서 publish 연산자를 호출해 만든 testSource$ 를 구독하고 connect 함수를 호출하여 500ms 마다 0~4까지 5개의 숫자를 멀티캐스팅한다.
5개 숫자를 모두 발행한 후 complete 함수를 호출한 3초 후에 재구독하도록 setTimeout 함수를 호출한다.
3초가 지난 후에는 setTimeout 함수의 콜백 함수가 실행되어 기존 구독을 모두 해제한다. 그리고 새로운 옵저버블인 c를 구독한 후 connect 함수를 호출한다. 이 때 publish 연산자 실행 전의 소스 옵저버블은 동작하므로 'tap 0'부터 'tap 4'까지 출력하지만 이를 구독하는 c로 시작하는 부분은 서브젝트 구독이 완료되었으므로 출력되지 않는다.
이는 서브젝트 구독이 완료되어 더 이상 next 함수로 값을 전달해도 이를 수용하지 않아 발생하는 현상이다.
주석처리된 부분 (서브젝트를 새로 생성하는 팩토리 함수를 multicast 연산자로 바꾼 testSource$)으로 교체하면 connect 함수를 호출할 때 서브젝트의 팩토리 함수를 호출하여 새로운 서브젝트를 만들어 주기 때문에 값을 잘 전달 받아 발행한다. 아래 결과를 확인해 보자.
실행 결과 - 주석 처리된 부분으로 대체
tap 0
a: 0
b: 0
tap 1
a: 1
b: 1
tap 2
a: 2
b: 2
tap 3
a: 3
b: 3
tap 4
a: 4
b: 4
timeout
tap 0
c: 0
tap 1
c: 1
tap 2
c: 2
tap 3
c: 3
tap 4
c: 4
3-1. publishXXX 연산자
publishBehavior, publishReplay, publishLast 연산자는 특정 서브젝트 자체를 멀티캐스팅하는 연산자다.
publishBehavior, publishReplay, publishLast 의 공통점
선택자 함수를 사용하지 않고 해당 서브젝트를 만드는데 필요한 것만 사용한다.
무조건 커넥터블 옵저버블을 리턴한다.
publishBehavior, publishReplay, publishLast 의 차이점
publishBehavior - BehaviorSubject 를 multicast 연산자를 사용해 커넥터블 옵저버블을 만들어 준다.
publishReplay - ReplaySubject 를 multicast 연산자를 사용해 커넥터블 옵저버블을 만들어 준다.
publishLast - AsyncSubject 를 multicast 연산자를 사용해 커넥터블 옵저버블을 만들어 준다.
연산자 원형 - publishBehavior publishBehavior<T>(value: T): UnaryFunction<Observable<T>, ConnetableObservable<T>>
연산자 원형 - publishLast publishLast<T>(): UnaryFunction<Observable<T>, ConnectObservable<R>>
/* publishBehavior 연산자의 구현 코드 */
import {BehaviorSubject} from "rxjs";
import {multicast} from "rxjs/operators";
export function publishBehavior(value) {
return (source) => multicast(new BehaviorSubject(value))(source);
}
[코드 12-10] publishBehavior 연산자의 구현 코드
/* publishReplay 연산자의 구현 코드 */
import {ReplaySubject} from "rxjs";
import {multicast} from "rxjs/operators";
export function publishReplay(bufferSize,
windowTime,
selectorOrScheduler,
scheduler) {
if (selectorOrScheduler && typeof selectorOrScheduler !== 'function') {
scheduler = selectorOrScheduler;
}
const selector = typeof selectorOrScheduler === 'function' ?
selectorOrScheduler : undefined;
const subject = new ReplaySubject(bufferSize, windowTime, scheduler);
return (source) => multicast(() => subject, selector)(source);
}
[코드 12-11] publishReplay 연산자의 구현 코드
/* publishLast 연산자의 구현 코드 */
import {multicast} from "rxjs/operators";
import {AsyncSubject} from "rxjs";
export function publishLast() {
return (source) => multicast(new AsyncSubject())(source);
}
[코드 12-12] publishLast 연산자의 구현 코드
4. refCount 연산자
refCOunt 연산자는 커넥터블 옵저버블을 구독하는 옵저버의 구를 카운트 한 후 최초로 1이 되면 connect 함수를 자동으로 호출한다. 또한, 옵버버블 구독을 1개 해제할 때마다 count를 1씩 줄이다가 0이 되면 unsubscribe 함수까지 자동으로 호출해준다.
tap 0
a: 0
b: 0
tap 1
a: 1
b: 1
tap 2
a: 2
b: 2
tap 3
a: 3
b: 3
tap 4
a: 4
b: 4
timeout
[코드 12-9] 와 다른점은 타임아웃 후 connect 함수를 호출하지 않는다는 것이다. 아래 구현 코드를 확인하자.
/* ConnectableObservable 의 구현 코드 일부분 */
class RefCountOperator {
constructor(connectable) {
this.connectable = connectable;
}
call(subscriber, source) {
const { connectable } = this;
connectable._refCount++;
const refCounter = new RefCountSubscriber(subscriber, connectable);
const subscription = source.subscribe(refCounter);
if (!refCounter.closed) {
refCounter.connection = connectable.connect();
}
return subscription;
}
}
[코드 12-14] ConnectableObservable 의 구현 코드 일부분
5. share 연산자
share 연산자는 pipe(publish(), refCount()) 를 추상화한 연산자다. 기존 옵저버블을 커넥터블 옵저버블로 바꾼후 refCount 연산자를 사용해 connect 함수를 호출할 필요 없는 핫 옵저버블을 만든다.
연산자 원형 share<T>(): MonoTypeOperatorFunction<T>
/* share 연산자의 구현 코드 일부 */
import {multicast, refCount} from "rxjs/operators";
import {Subject} from "rxjs";
function shareSubjectFactory() {
return new Subject();
}
export function share() {
return (source) => refCount()(multicast(shareSubjectFactory)(source));
}
[코드 12-15] share 연산자의 구현 코드 일부
새로운 서브젝트를 리턴하는 팩토리 함수를 multicast 연산자에서 사용한다. 즉, publish 연산자에서 서브젝트 자체를 사용했을때 발생하는 재구독 문제를 피할수 있다.
multicast 연산자에서 사용하는 팩토리 함수 덕분에 소스 옵저버블에서 값을 다 발행하고 구독 완료했다면, refCount 연산자가 발행한 값이 0이 된 이후 재구독을 하여도 새로운 값을 전달 받을 수 있다.
서브젝트는 옵저버이기도 하므로 옵저버블에서 값을 바로 서브젝트로 전달해줘서 함수 중복을 피할 수 있다. 즉, 서브젝트의 옵저버블 특성을 잘 활용하면 옵저버블과 연결하거나 직접 옵저버의 함수들을 호출해서 멀티캐스팅하려는 값, 이벤트, 에러, 완료에 관한 정보를 보낼 수 있다.
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();
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');
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. 마치며
서브젝트는 옵저버의 특성과 옵저버블의 특성을 동시에 갖고 멀티캐스팅을 지원한다. 옵저버블과 연결해 멀티캐스팅을 지원하지 않는 옵저버블의 발행 값을 여러 옵저버로 전파할 수 있다.
tap 연산자에서 해당 값을 parseInt 로 바꿔 정수인지 검사하는데 r은 정수값이 아니므로 TypeError 를 전달한다. 이 때 에러각 발생한 지점의 index 값을 error에 넣고 catchError 연산자의 선택자 함수는 err 객체를 전달 받는다. 해당 레러가 정수 확인 중 전달된 에러가 맞다면 에러 메세지를 출력한 후 해당 index 부터 나머지 값을 자례로 발행한다. 다른 에러면 에러 메세지만 발행하도록 옵저버블을 리턴한다.
subscribe 함수에 error 함수도 설정했지만, catchError 연산자가 에러를 적절히 처리해줘서 해당 함수를 호출하지 않는다. 하지만, catchError 연산자의 선택자 함수에서 리턴받아 구독한 옵저버블에서 에러가 발생하면 error 함수를 호출한다.
1-1. mergeMap 연산자를 사용한 catchError 연산자 응용
[코드 10-1]은 에러가 발생한 원래 스트림에서 에러가 발생한 값만 처리하고 나머지 값들을 그대로 처리할 수 없다. 따라서 mergeMap 연산자를 추가한후 소스 옵저버블에서 발행하는 값 각각을 옵저버블로 감싸서 여기에 catchError 연산자를 적용해야 한다. 이렇게 하면 에러가 발생해도 나머지 값들을 계속 발행할 수 있다.
0
1
2
3
4
RANDOM ERROR 5
6
7
8
RANDOM ERROR 9
10
11
12
13
14
RANDOM ERROR 15
RANDOM ERROR 16
RANDOM ERROR 17
18
19
RANDOM ERROR 20
21
22
23
24
25
26
27
28
RANDOM ERROR 29
참고로 소스 옵저버블 구독을 처음부터 다시 재시도하려면 mergeMap 연산자를 사용하지 않아도 된다. 하지만, 여러 값을 발행한는 스트림 각각에 재시도를 해야 한다면 mergeMap 연산자를 사용해야 한다.
3. retryWhen 연산자
retrtWhen 연산자는 소스 옵저버블 구독을 재시도한다는 점에서는 retry 연산자와 비슷하지만, 구독을 재시도하기 전 에러를 전달받아 특정한 옵저버블에서 발행한 후 구독을 재시도하는 연산자다.
retryWhen 연산자의 마블 다이어그램
연산자 원형 retryWhen<T>( notifier: (errors: Observable<any>) => Observable<any> ): MonoTypeOperatorFunction<T>
notifier 함수 - 소스 옵저버블에서 에러가 발생했을 때 에러를 errors 라는 옵저버블로 다룬다. 이 옵저버블의 스트림을 전달받아 notifier에 리턴하는 옵저버블을 구독한다.
notifier 에서 리턴하는 옵저버블의 값을 발행하면 이어소 소스 옵저버블 구독을 재시도한다. 에러가 발행하면 전체 스트림의 구독을 종료한다. 따라서, 구독 재시도가 필요하면 재시도하기 전 해줘야 할 일을 처리한 후 아무 값이든 발행해 해당 스트림을 종료시키지 않도록 해야한다.
find 연산자는 인자로 사용하는 predicate 함수로 소스 옵저버블에서 발행하는 값 중 처음으로 함수 조건을 만족했을 때 true 를 리턴하는 값을 발행하고 구독을 완료하는 연산자다. 구독을 완료 할 때까지 조건을 만족하는 값이 없었다면 undefined 라는 값을 발행한다.
find 연산자의 마블 다이어그램
연산자 원형 find<T>( predicate: (value: T, index: number, source: Observable<T>) => boolean, thisArgs?: any ): MonoTypeOperatorFunction<T>
predicate 함수를 호출해 동작하는 중 에러가 발행하면 error 함수를 호출해 에러를 전달 받는다. 소스 옵저버블에서 에러가 발생해도 error 함수로 에러를 전파한다.
nextOrObserver? - 함수일 때 소스 옵저버블에서 발행하는 다음 값을 전달받는 next 콜백 함수로 동작한다. 객체이면 옵저버 객체로 다룬다.
error, complete - 에러가 발생하거나 소스 옵저버블 구독을 완료했을 때 발생하는 콜백이다.
1-1. next 콜백 사용
/* next 함수에서 발행하는 값만 부수 효과로 처리 */
const { range } = require('rxjs');
const { tap, filter, map } = require('rxjs/operators');
range(1, 10).pipe(
tap(x => console.log(`stream 1 (range 1, 10) ${x}`)),
filter(x => x % 2 === 0),
tap(x => console.log(` stream 2 (filter x % 2 === 0) ${x}`)),
map(x => x + 1),
tap(x => console.log(` stream 3 (map x + 1) ${x}`))
).subscribe(x => console.log(` result ${x}`));
[코드 8-1] next 함수에서 발행하는 값만 부수 효과로 처리
실행 결과
stream 1 (range 1, 10) 1
stream 1 (range 1, 10) 2
stream 2 (filter x % 2 === 0) 2
stream 3 (map x + 1) 3
result 3
stream 1 (range 1, 10) 3
stream 1 (range 1, 10) 4
stream 2 (filter x % 2 === 0) 4
stream 3 (map x + 1) 5
result 5
stream 1 (range 1, 10) 5
stream 1 (range 1, 10) 6
stream 2 (filter x % 2 === 0) 6
stream 3 (map x + 1) 7
result 7
stream 1 (range 1, 10) 7
stream 1 (range 1, 10) 8
stream 2 (filter x % 2 === 0) 8
stream 3 (map x + 1) 9
result 9
stream 1 (range 1, 10) 9
stream 1 (range 1, 10) 10
stream 2 (filter x % 2 === 0) 10
stream 3 (map x + 1) 11
result 11
단계별로 filter 연산자의 조건을 만족하지 못하면 해당 부분만 출력되고, 조건을 만족하면 다음 스트립을 출력한다.
1-2. error 콜백 함수 사용
/* error 콜백 함수를 사용하는 예 */
const { range } = require('rxjs');
const { map, tap } = require('rxjs/operators');
range(1, 8).pipe(
map(x => x === 8 ? x.test() : x + 1),
tap(
x => console.log(`tap next: ${x}`),
err => console.error(`tap ERROR: ${err}`)
)
).subscribe(
x => console.log(`result: ${x}`),
err => console.error(`subscribe ERROR: ${err}`)
);
[코드 8-2] error 콜백 함수를 사용하는 예
실행 결과
tap next: 2
result: 2
tap next: 3
result: 3
tap next: 4
result: 4
tap next: 5
result: 5
tap next: 6
result: 6
tap next: 7
result: 7
tap next: 8
result: 8
tap ERROR: TypeError: x.test is not a function
subscribe ERROR: TypeError: x.test is not a function
8에 해당 값은 에러 때문에 map 연산자로 값을 변환하지 못하므로 출력 결과가 없고 tap 연산자의 두번째 인자인 error 콜백 함수를 호출한 후 그 다음 subscribe 에 있는 옵저버의 error 콜백 함수를 호출한다.
tap next: 1 STREAM 1
result: 1
tap next: 2 STREAM 1
result: 2
tap next: 3 STREAM 1
result: 3
tap next: 4 STREAM 1
result: 4
complete STREAM 1
tap next: 5 STREAM 2
result: 5
tap next: 6 STREAM 2
result: 6
tap next: 7 STREAM 2
result: 7
complete STREAM 2
subscribe complete
실행 결과는 [코드 8-3] 과 같다. 즉, 옵저버 객체 자체를 전달할 수 있다.
2. finalize 연산자
finalize 연산자는 옵저버블 스트림 실행을 완료하거나 에러가 발생했을 때 인자로 사용하는 콜백 함수를 호출하는 연산자다.
연산자 원형 finalize<T>(callback: () => void): MonoTypeOperatorFunction<T>
finalize 연산자는 기존 구독하는 소스 옵저버블에 영향을 주지 않고 옵저버블 라이프사이클이 끝날 때 호출되는 콜백 함수를 인자로 사용한다.
/* finalize 연산자의 구현 코드 일부 */
import {Subscription} from "rxjs";
export function finalize(callback) {
return (source) => source.lift(new FinallyOperator(callback));
}
class FinallyOperator {
constructor(callback) {
this.callback = callback;
}
call(subscriber, source) {
return source.subscribe(new FinallySubscriber(subscriber, this.callback));
}
}
class FinallySubscriber {
constructor(destination, callback) {
super(destination);
this.add(new Subscription(callback));
}
}
finalize 연산자의 인자로 사용하는 콜백 함수는 FinallySubscriber 클래스에 unsubscribe 함수의 콜백 함수 Subscription 객체를 add 함수로 추가해 구독한다. 그러므로, 소스 옵저버블에서 complete 함수나 error 함수를 호출했을 때 해당 콜백 함수도 같이 호출한다.
tap 연산자와는 달리 subscribe 함수에서 complete 함수를 호출한 후 finalize 연산자가 인자로 사용하는 콜백 함수를 호출한다.
2-1. 에러가 발생했을 때의 finalize 연산자 사용
/* 에러가 발생했을 때 finalize 연산자 사용 예 */
const { range } = require("rxjs");
const { finalize, tap } = require("rxjs/operators");
range(1, 3).pipe(
tap(x => x === 3 && x.test()),
finalize(() => console.log('FINALLY CALLBACK'))
).subscribe(
x => console.log(`result ${x}`),
err => console.error(`ERROR: ${err}`),
);
[코드 8-7] 에러가 발생했을 때 finalize 연산자 사용
실행 결과
result 1
result 2
ERROR: TypeError: x.test is not a function
FINALLY CALLBACK
3은 tap 연산자의 x.test() 를 호출해서 일부러 에러를 발생 시켰다. 에러가 출력된 후 finalize 연산자의 콜백 함수를 호출하였다.
3. toPromise 함수
toPromise 함수는 호출 후 구독해서 동작하는 것이 아니다. 호출하자마자 새로 생성한 프로미스를 리턴해준다. 또한, 프로미스 안 함수에서 소스 옵저버블인 this를 사용해 구독하므로 프로미스 생성과 동시에 소스 옵저버블을 구독한다. 그리고 toPromise 함수에서 리턴하는 프로미스는 소스 옵저버블 구독이 완료되었을 때 가장 최근 값을 resolve 로 갖는다. 중간에 에러가 발생하면 해당 프로미스의 reject 로 에러를 전달하도록 동작한다.
/* roPromise 함수의 사용 예 */
const { interval } = require('rxjs');
const { take, tap } = require('rxjs/operators');
interval(100).pipe(
take(10),
tap(x => console.log(`interval tap ${x}`))
).toPromise().then(
value => console.log(`프로미스 결과 ${value}`),
reason => console.error(`프로미스 에러 ${reason}`)
);
[코드 8-8] toPromise 함수의 사용 예
실행 결과
interval tap 0
interval tap 1
interval tap 2
interval tap 3
interval tap 4
interval tap 5
interval tap 6
interval tap 7
interval tap 8
interval tap 9
프로미스 결과 9
toPromise 함수를 호출할 때 프로미스를 생성하며 소스 옵저버블을 구독하므로 100ms 마다 0부터 1씩 증가하는 10개 숫자가 순서대로 tap 연산자 안에서 호출된다. 그리고 이 소스 옵저버블의 구독이 완료되면 발행한 값을 프로미스 결과로 리턴하는 것이다. 또한, 해당 프로미스에 then 함수를 호출해 결과 값을 전달받으면 결과값을 출력한다.
3-1. toPromise 함수의 reject 에러 처리
/* toPromise 함수의 소스 옵저버블에서 에러가 발생한 예 */
const { interval } = require('rxjs');
const { take, tap } = require('rxjs/operators');
interval(100).pipe(
take(10),
tap(x => console.log(`interval tap ${x < 3 ? x : x.test()}`))
).toPromise().then(
value => console.log(`프로미스 결과 ${value}`),
reason => console.error(`프로미스 에러 ${reason}`)
);
[코드 8-9] toPromise 함수의 소스 옵저버블에서 에러가 발생한 예
실행 결과
interval tap 0
interval tap 1
interval tap 2
프로미스 에러 TypeError: x.test is not a function
tap 연산자 안에서 값이 3 미만일 때 해당 값을 출력하다가 그 이후에는 x.test()를 호출해서 일부러 에러를 발생시켰다. 숫자 값만 발행하므로 test 를 호출하면 타입 에러가 발생하는 것이다.
에러가 발생하면 해당 프로미스의 then 함수 안에 있는 에러 처리 함수를 호출해 에러 메세지를 출력한다.
4. toArray 연산자
toArray 연산자는 소스 옵저버블에서 발행한 값을 내부에 생성한 배열에 저장하다가 소스 옵저버블 구독이 완료되면 해당 배열을 next 함수로 발행하도록 동작하는 연산자다. 새로운 옵저버블을 리턴하고 이를 구독해야만 옵저버블 안에서 결과를 전달받을 수 있으며 subscribe 함수 호출만으로 구독하는 일은 없다.
toArray 연산자의 마블 다이어그램
연산자 원형 toArray<T>(): OperatorFunction<T, T[]>
소스 옵저버블에서 에러가 발생하면 error 함수를 호출한다.
사용시 주의 사항 구독을 완료할 때까지 배열에 값을 저장하므로 무한 스트림에서 사용하지 않도록 해야한다. 너무 많은 값을 저장하면 배열이 커지므로 메모리 이슈에 주의 해야 한다.
/* toArray 연산자의 사용 예 */
const { range } = require('rxjs');
const { filter, toArray } = require('rxjs/operators');
range(1, 30).pipe(
filter(x => x % 2 === 0),
toArray()
).subscribe(
value => console.log(`배열여부: ${Array.isArray(value)}, 값: ${value}`)
);
소스 옵저버블은 0이란 값 1개만 발행하지만 누적자 함수로 초기값 1을 함께 누적해 최종 1이라는 값을 발행한다.
[코드7-3]실행 과정
소스 옵저버블 0 -> 누적자 함수 호출 (초기값 acc = 1, curr = 0), 1 + 0 = 1 계산 후, 누적값 = 1
complete -> 누적값 1 발행 후 구독 완료 (next(1), complete 함수를 차례로 호출)
/* 초기값 1을 설정해 여러 개 값을 발행하는 예 */
const { range } = require('rxjs');
const { reduce } = require('rxjs/operators');
range(1, 4).pipe(reduce((acc, cur) => acc + cur, 1))
.subscribe(result => console.log(`result: ${result}`));
[코드 7-4] 초기값 1을 설정해 여러 개 값을 발행하는 예
실행 결과
result: 11
[코드7-4]에서 누적자 함수를 호출해 누적하는 과정
소스 옵저버블 1 -> 누적자 함수 호출 (초기값 acc = 1, curr = 1), 1 + 1 = 2 계산 후, 누적값 = 2
소스 옵저버블 2 -> 누적자 함수 호출 (acc = 2, curr = 2), 2 + 2 = 4 계산 후, 누적값 = 4
소스 옵저버블 3 -> 누적자 함수 호출 (acc = 4, curr = 3), 4 + 3 = 7 계산 후, 누적값 = 7
소스 옵저버블 4 -> 누적자 함수 호출 (acc = 7, curr = 4), 7 + 4 = 11 계산 후, 누적값 = 11
complete -> 누적값 11 발행 후 구독 완료 (next(10), complete 함수를 차례로 호출)
1-3. 누적자 함수의 index 파라미터
0부터 시작해 누적자 함수를 호출할 때 소스 옵저버블의 몇 번째 값을 전달하는지를 나타낸다.
초기값이 없으면 index 는 1이 된다. 초기값이 있으면 index 는 0이 된다.
2. max 연산자
max 연산자는 발행되는 값 중 가장 큰 값을 출력하는 함수다. reduce 연산자에 제일 큰 값을 누적한 함수의 값을 전달하는 방법으로 동작한다.
max 연산자의 마블 다이어그램
연산자 원형 amx<T>(comparer?: (x: T, y: T) => number): MonoTypeOperatorFunction<T>
comparer? - 두 값을 비교하려고 기본값 대신 사용할 비교 함수
/* max 연산자의 구현 코드 일부분 */
export function max(comparer) {
const max = (typeof comparer === 'function')
? (x, y) => comparer(x, y) > 0 ? x : y
: (x, y) => x > y ? x : y;
return this.lift(new ReduceOperator(max));
}
[코드 7-5] max 연산자의 구현 코드 일부분
reduce 연산자를 이용해 소스 옵저버블에서 complete 함수를 호출해야 지금까지 누적한 가장 큰 값을 발행한다. 인자가 없으면 부등호로 값의 크기를 비교한다. 따라서, 부등호로 비교할 수 있는 값만 발행해야 정상 동작한다. 그렇지 않는 값이나 객체는 comparer 함수를 인자로 사용해 비교해야 한다.
/* comparer 함수를 사용하지 않는 예 */
const { range } = require('rxjs');
const { max } = require('rxjs/operators');
range(1, 10).pipe(max())
.subscribe(result => console.log(`result: ${result}`));
comparer 함수는 x.avg에서 y.avg를 빼서 평점을 비교한다. x의 평균 평점이 y의 평균 평점보다 클 때 0보다 큰 값을 리턴해 올바른 결과를 출력한다. 최고 평점에 해당하는 객체를 찾아 JSON.stringify를 이용해 JSON 결과를 출력한다.
2-2. 다른 객체지만 같은 값으로 평가할 때의 max 연산자 사용
[코드 7-7]에서 만약 평균 평점이 같은 객체가 여러 개 있으면 지금까지 누적된 x가 현재 값인 y보다 크지 않으므로 y를 리턴한다. 소스 옵저버블에서 가장 나중에 발행하는 객체를 선택한다.
이유 -> max 연산자의 내부 구현에서 reduce 연산자의 누적자 함수의 구현이 (x, y) => x > y ? s : y 이기 때문이다. 부등호가 > 대신 >= 였다면 누적값인 x가 선택되므로 가장 먼저 나온 값을 발행한다.
/* 같은 값으로 평가하는 다른 객체 처리 */
const { from } = require('rxjs');
const { max } = require('rxjs/operators');
const movies = [
{ title: '영화 1', avg: 5.12 },
{ title: '영화 2', avg: 9.14 },
{ title: '영화 3', avg: 8.28 },
{ title: '영화 4', avg: 9.14 }
];
from(movies).pipe(max((x, y) => x.avg - y.avg))
.subscribe(x => console.log(JSON.stringify(x)));
[코드 7-8] 같은 값으로 평가하는 다른 객체 처리
실행 결과
{"title":"영화 4","avg":9.14}
'영화 2'와 '영화 4'가 평균 평점이 같다. 이럴 경우 소스 옵저버블에서 가장 마지작에 발행한 '영화 4'를 발행했다.
3. min 연산자
min 연산자는 reduce 연산자에 가장 작은 값을 누적한 누적자 함수 값을 전달한다.
min 연산자의 마블 다이어그램
연산자 원형 min<T>(comaprer?: (x: T, y: T) => number): MonoTypeOperatorFunction<T>
max 연산자의 반대 개녕믜 연산자이다.
/* min 연산자의 구현 코드 일부분 */
export function min(comparer) {
const min = (typeof comparer === 'function')
? (x, y) => comparer(x, y) < 0 ? x : y
: (x, y) => x < y ? x : y;
return this.lift(new ReduceOperator(min));
}
[코드 7-9] min 연산자의 구현 코드 일부분
max 연산자으 구현 코드와 부등호 방향만 다르다.
min 연산자도 값이 같으면 y를 리턴하므로 뒤에 있는 값이 발행된다.
/* comparer 함수를 사용하지 않는 min 연산자 예 */
const { range } = require('rxjs');
const { min } = require('rxjs/operators');
range(1, 10).pipe(min())
.subscribe(result => console.log(`result: ${result}`));
[코드 7-10] comparer 함수를 사용하지 않는 min 연산자 예
실행 결과
result: 1
소스 옵저버블이 발행하는 값이 숫자이므로 부동호 바교하여 가장 작은 값인 1을 발행한다.
/* comparer 함수를 사용하는 min 연산자 예 */
const { from } = require('rxjs');
const { min } = require('rxjs/operators');
const movies = [
{ title: '영화 1', avg: 5.12 },
{ title: '영화 2', avg: 9.14 },
{ title: '영화 3', avg: 8.28 }
];
from(movies).pipe(min((x, y) => x.avg - y.avg))
.subscribe(x => console.log(JSON.stringify(x)));
[코드 7-11] comparer 함수를 사용하는 min 연산자 예
실행 결과
{"title":"영화 1","avg":5.12}
/* 같은 값으로 평가하는 다른 객체를 min 연산자로 처리 */
const { from } = require('rxjs');
const { min } = require('rxjs/operators');
const movies = [
{ title: '영화 1', avg: 5.12 },
{ title: '영화 2', avg: 9.14 },
{ title: '영화 3', avg: 5.12 }
];
from(movies).pipe(min((x, y) => x.avg - y.avg))
.subscribe(x => console.log(JSON.stringify(x)));
[코드 7-12] 같은 값으로 평가하는 다른 객체를 min 연산자로 처리
실행 결과
{"title":"영화 3","avg":5.12}
'영화 1'과 '영화 3'이 가장 작은 평균 평점이며 값은 같다. 이럴 경우 소스 옵저버블에서 가장 마지작에 발행한 '영화 3'를 발행했다.
4. count 연산자
count 연산자는 소스 옵저버블에서 값을 발행할 때마다 개수를 내부에서 카운트한다. 구독을 완료한 후 총 몇 개인지 발행한다.