1. tap 연산자
tap 연산자는 소스 옵저버블에서 발행하는 값을 전달 받은 후 인자로 사용하는 함수를 호출하고 소스 옵저버블에서 발행한 값을 그대로 발행한다.

연산자 원형
tap<T>(
nextOrObserver?: partialObserver<T> | ((x: T) => void),
error?: (e: any) => void,
complete?: () => void
): MonoTypeOperatorFunction<T>
- 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 콜백 함수를 호출한다.
1-3. complete 콜백 함수 사용
/* concat 연산자를 기준으로 앞과 뒤에서 값 발행 */
const { concat, range } = require('rxjs');
const { tap } = require('rxjs/operators');
concat(
range(1, 4).pipe(
tap(
x => console.log(`tap next: ${x} STREAM 1`),
err => console.error(`tap ERROR: ${err} STREAM 1`),
() => console.log('complete STREAM 1')
)
),
range(5, 3).pipe(
tap(
x => console.log(`tap next: ${x} STREAM 2`),
err => console.error(`tap ERROR: ${err} STREAM 2`),
() => console.log('complete STREAM 2')
)
)
).subscribe(
x => console.log(` result: ${x}`),
err => console.error(` subscribe ERROR: ${err}`),
() => console.log(' subscribe complete')
);
[코드 8-3] concat 연산자를 기준으로 앞과 뒤에서 값 발행
실행 결과
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
1-4. 옵저버로 콜백 함수를 묶어서 사용
/* 옵저버 객체를 사용하는 tap 연산자 */
const { range, concat } = require('rxjs');
const { tap } = require('rxjs/operators');
const observer1 = {
next: x => console.log(`tap next: ${x} STREAM 1`),
error: err => console.error(`tap ERROR: ${err} STREAM 1`),
complete: () => console.log('complete STREAM 1')
};
const observer2 = {
next: x => console.log(`tap next: ${x} STREAM 2`),
error: err => console.error(`tap ERROR: ${err} STREAM 2`),
complete: () => console.log('complete STREAM 2')
};
concat(
range(1, 4).pipe(tap(observer1)),
range(5, 3).pipe(tap(observer2))
).subscribe(
x => console.log(` result: ${x}`),
err => console.error(` subscribe ERROR: ${err}`),
() => console.log(' subscribe complete')
);
[코드 8-4] 옵저버 객체를 사용하는 tap 연산자
실행 결과
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 함수를 호출했을 때 해당 콜백 함수도 같이 호출한다.
/* finalize 연산자를 사용한 예 */
const { range } = require("rxjs");
const { finalize } = require("rxjs/operators");
range(1, 3).pipe(
finalize(() => console.log('FINALLY CALLBACK'))
).subscribe(
x => console.log(`next: ${x}`),
err => console.error(`error: ${err}`),
() => console.log('COMPLETE')
);
[코드 8-6] finalize 연산자를 사용한 예
실행 결과
next: 1
next: 2
next: 3
COMPLETE
FINALLY CALLBACK
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<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}`)
);
[코드 8-10] toArray 연산자의 사용 예
실행 결과
배열여부: true, 값: 2,4,6,8,10,12,14,16,18,20,22,24,26,28,30
소스 옵저버블에서는 총 15개 값을 발행하지만 최종 구독을 완료한 시점에 배열 1개만 발행한다는 사실을 확인할 수 있다.
5. timeout 연산자
timeout 연산자는 일정 시간 동안 소스 옵저버블에서 값을 발행하지 않으면 에러를 발생시키는 연산자다.
서버에 어떤 요청을 하거나, 상황에 따라서 기대한 시간보다 오래 걸릴 수 있는 작업에 옵저버블을 사용해야 할 때 유용하다.

연산자 원형
timeout<T>(
due: number | Date, scheduler: SchedulerLike = async
): MonoTypeOperatorFunction<T>
일정 시간 안에 응답이 오지 않으면 에러 메세지를 표시하거나, 특정 표시를 하지 않거나, 예외 처리를 할 수 있다.
/* node-fetch 라이브러리를 이용하는 timeout 연산자 사용 예 */
const { defer, timer} = require('rxjs');
const { timeout, map} = require('rxjs/operators');
const fetch = require('node-fetch');
const source$ = defer(() =>
fetch(`https://httpbin.org/delay/${parseInt(Math.random() * 5, 10)}`)
.then(x => x.json())
);
/*
const source$ = timer(Math.floor(Math.random() * 2000)).pipe(
map(x => ({ value: x }))
);
*/
source$.pipe(timeout(2000)).subscribe(
x => console.log(`${JSON.stringify(x)}`),
err => {
console.error(`ERROR: ${err}`);
process.exit(1);
}
);
[코드 8-11] node-fetch 라이브러리를 이용하는 timeout 연산자 사용 예
실행 결과
# 타임아웃 에러가 발생했을 때
ERROR: TimeoutError: Timeout has occurred
# 2초 안에 응답이 올 때 (주석 처리 부분 활성화 했을 때)
{"value":0}'RxJS' 카테고리의 다른 글
| 10장. 에러 처리 요약 (0) | 2024.08.05 |
|---|---|
| 9장. 조건 연산자 요약 (0) | 2024.08.05 |
| 7장. 수학 및 결합 연산자 요약 (0) | 2024.08.02 |
| 6장. 조합 연산자 요약 (0) | 2024.08.01 |
| 5장. 변환 연산자 요약 (0) | 2024.07.31 |