반응형

변환 연산자는 옵저버블의 값을 다른 값으로 바꾸는 옵저버블을 만든다.

1. map 연산자

map 연산자를 사용한 옵저버블이 발행하는 값 각각은 소스 옵저버블에서 발행한 값 각각에 project라는 함수를 적용한 결과를 발행하는 옵저버블로 바꾼다.

map 연산자의 마블 다이어그램

소스 옵저버블의 값에 10을 곱하는 project 함수를 적용하여 10, 20, 30을 발행하는 옵저버블로 변환되었다.

연산자 원형
map<T, R>(
       project: (value: T, index: number) => R,
       thisArg?: any
): OperatorFunction<T, R>

RxJS의 map 연산자는 다른 map 연산자와는 달리 소스 옵저버블을 map 연산자로 동작하는 새 옵저저블로 바꾸고, 이를 구독할 때 각각의 값을 발행한다.

  • thisArg? - project 함수에서 어떤 것을 정의할 때 사용
  • OperatorFunction<T, R> - 대부분 파라미터의 동작에 맞게 발행한 값을 옵저버블 형태로 리턴
/* 옵버버블을 사용한 map 연산자 */
const { from } = require('rxjs');
const { map } = require('rxjs/operators');

const source$ = from([1, 2, 3, 4, 5]);
const resultSource$ = source$.pipe(
  map(x => x + 1),
  map(x => x * 2)
);

resultSource$.subscribe(x => console.log(x));

[코드 5-1] 옵저버블을 사용한 map 연산자

/* 배열을 사용한 map 연산자 */
const sourceArray = [1, 2, 3, 4, 5];
const resultArray = sourceArray.map(x => x + 1).map(x => x * 2);

for (let i = 0; i < resultArray.length; i++) {
  console.log(resultArray[i]);
}

[코드 5-2] 배열을 사용한 map 연산자

 

실행 결과 (코드5-1, 코드5-2)

4
6
8
10
12

코드 5-1과 코드 5-2는 실행 결과는 같지만 동작은 다르다.

코드 5-1과 코드5-2의 차이점
[코드5-1] - map 연산자를 실행할 때마다 새 옵저버블의 값 각각을 바로 발행하므로 실행 결과가 빠르다.
[코드5-2] - 모든 요소를 확인한 후 새 배열을 만들기 때문에 실행 결과가 느리다.
/* 배열 요소 각각을 바로 출력하기 */
const sourceArray = [1, 2, 3, 4, 5];
const func1 = x => x + 1;
const func2 = x => x * 2;

for (let i = 0; i < sourceArray.length; i++) {
  console.log(func2(func1(sourceArray[i])));
}

[코드5-3] 배열 요소의 각각을 바로 출력하기

 

실행결과는 위와 모두 같다.

새 배열을 만들 때까지 기다리지 않고, 배열 각 요소를 for 문으로 순회하여 바로 출력한다.

/* 발행한 값이 짝수인지 홀수인지 확인 */
const { range } = require('rxjs');
const { map } = require('rxjs/operators');

const source$ = range(0, 5).pipe(
  map(x => ({ x, isEven: x % 2 === 0 }))
);

source$.subscribe(result =>
  console.log(`${result.x}은(는) ${result.isEven ? "짝수" : "홀수"} 입니다.`)
);

[코드 5-4] 발행한 값이 짝수인지 홀수인지 확인

 

실행 결과

0은(는) 짝수 입니다.
1은(는) 홀수 입니다.
2은(는) 짝수 입니다.
3은(는) 홀수 입니다.
4은(는) 짝수 입니다.

map 연산자는 다른 타입이나 객체를 리턴할 수도 있다.


2. pluck 연산자

pluck 연산자는 map 연산자처럼 동작하지만, 소스 옵저버블에서 객체를 리턴할 때 해당 객체의 속성을 기준으로 변환하는 연산자다.

pluck 연산자의 마블 다이어그램

연산자 원형
pluck<T, R>(...propertied: string[]): OperatorFunction>T, R>

...properties 는 객체의 속성이 중첩될 때 중첩된 속성의 이름을 순서대로 나열해 사용할 수 있는 파라미터다.

/* pluck 연산자 사용 예 */
const { range } = require('rxjs');
const { map, pluck } = require('rxjs/operators');

const source$ = range(0, 5).pipe(
  map(x => ({ x, isEven: x % 2 === 0 }))
);

source$.pipe(pluck('isEven')).subscribe(isEven =>
  console.log(`${isEven ? "짝수" : "홀수"} 입니다.`)
);

source$.pipe(pluck('x')).subscribe(x =>
  console.log(`${x}입니다.`)
);

[코드5-5] pluck 연산자 사용 예

 

실행 결과

짝수 입니다.
홀수 입니다.
짝수 입니다.
홀수 입니다.
짝수 입니다.
0입니다.
1입니다.
2입니다.
3입니다.
4입니다.

pluck 연산자는 소스 옵저버블에서 발행한 값을 저장하는 객체에서 꺼내고 싶은 속성 이름을 문자열로 전달받는다.

/* 중첩 속성이 있는 객체에 pluck 연산자 사용 */
const { range } = require('rxjs');
const { map, pluck } = require('rxjs/operators');

const source$ = range(0, 5).pipe(map(x =>
  ({ x, numberProperty: {
    isEven: x % 2 === 0,
  }})
));

source$.pipe(pluck('numberProperty', 'isEven'))
  .subscribe(isEven => console.log(`${isEven ? "짝수" : "홀수"} 입니다.`));

[코드 5-6] 중첩 속성이 있는 객체에 pluck 연산자 사용

 

실행 결과

짝수 입니다.
홀수 입니다.
짝수 입니다.
홀수 입니다.
짝수 입니다.
  1. 전통적인 자바스크립트 문법 - 중첩 형태의 객체를 다뤄야 할 때 각각의 값이 있는지 없는지 확인하기 위해 중첩된 if문을 사용하게나, &&, ||와 같은 논리 연산자를 사용해 처리한다.
  2. RxJS의 pluck 연산자 사용 - 인자로 여러값을 나열하면 좀 더 쉽게 값을 변환할 수 있다.

해당 속성 이름에 속하는 값이 없다면(중첩 구조를 찾다 중간에 속성 이름이 없거나 처음부터 속성 이름이 없을 때 모두) undefined를 발행한다.

3. mergeMap 연산자

mergeMap 연산자는 Observable 인스턴스를 리턴하는 project 함수를 인자로 사용해 여기서 리턴된 인스턴스를 구독하는 map 연산자다.

 

map과 mergeMap의 차이

  • map 연산자 - Observable 객체 자체를 발행하는 방식
  • mergeMap 연산자 - project 함수에서 리턴하는 Observable 객체를 구독해 값을 각각 발행하는 방식
연산자 원형
mergeMap<T, I, R>(
       project: (value: T, index: number) => ObservableInput<I>,
       resultSelector?: ((
              outerValue: T, innerValue: I,
              outerIndex: number, innerIndex: number) => R
       ) | number,
       concurrent: number = Number.POSITIVE_INFINITY
): OperatorFunction<T, I | R>
  • project 함수 - Observable 클래스의 소스 옵저버블을 리턴
  • resultSelector? - 연산자 원형에 있는 타입을 설정하는 역할 (RxJS6 -> 사용을 권장하지 않음)
  • concurrent - 동시에 구독 중인 최대 옵저버블 수를 설정

mergeMap 연산자의 마블 다이어그램

/* subscribe 함수를 중접 사용해 배열 요소의 값을 발행 */
const {timer, range} = require('rxjs');
const {map} = require('rxjs/operators');

const requests = [
  timer(Math.floor(Math.random() * 2000)).pipe(map(value => "req1")),
  timer(Math.floor(Math.random() * 1000)).pipe(map(value => "req2")),
  timer(Math.floor(Math.random() * 1500)).pipe(map(value => "req3")),
];

range(0, 3).subscribe(x => {
  requests[x].subscribe(req => console.log(`response from ${req}`));
});

[코드5-7] subscribe 함수를 중첩 사용해 배열 요소의 값을 발행

 

실행 결과

response from req1
response from req3
response from req2

timer 함수를 사용해 값 발행 간격이 다른 무작위 요청을 배열에 넣고 배열 요소 순서로 요청해 응답이 온 순서대로 출력한다.

(매번 순서가 다를 수 있다.)

/* mergeMap 연산자로 리턴한 옵저버블을 구독해 값 발행 */
const { timer, range } = require('rxjs');
const { mergeMap, map } = require('rxjs/operators');

const requests = [
  timer(Math.floor(Math.random() * 2000)).pipe(map(value => "req1")),
  timer(Math.floor(Math.random() * 1000)).pipe(map(value => "req2")),
  timer(Math.floor(Math.random() * 1500)).pipe(map(value => "req3")),
];

range(0, 3).pipe(mergeMap(x => requests[x]))
  .subscribe(req => console.log(`response from ${req}`));

[코드 5-8] mergeMap 연산자로 리턴한 옵저버블을 구독해 값 발행

 

실행 결과

response from req3
response from req2
response from req1

[코드 5-8]은 [코드5-7]에서 이중 구독하는 requests 옵저버블을 mergeMap의 project 함수에서 리턴하도록 바꾸고, 실행 결과를 발행하는 옵저버블로 만든 예이다.

(매번 순서가 다를 수 있다.)

3-1. mergeMap 연산자에 사용하는 배열, 프로미스, 이터러블

mergeMap 연산자는 project 함수에서 리턴하는 객체를 구독할 때 subscribeToResult라는 함수를 사용한다. 객체의 타입을 검사해 적절한 동작으로 바꿔주므로 배열, 이터러블, 비동가 동작을 위한 프로미스로 리턴할 수 있다.

mergeMap 연산자에 배열이나 유사 배열 사용

mergeMap의 project 함수가 배열을 리턴하면 배열의 길이 만큼 순회하며 next 함수로 값을 발행한 후 complete 함수를 호출한다.

/* mergeMap 연산자에 배열 사용 */
const { range } = require('rxjs');
const { mergeMap } = require('rxjs/operators');

range(0, 3).pipe(mergeMap(x => [x + 1, x + 2, x + 3, x + 4]))
  .subscribe(value => console.log(`current value: ${value}`));

[코드 5-9] mergeMap 연산자에 배열 사용

 

실행 결과

current value: 1
current value: 2
current value: 3
current value: 4
current value: 2
current value: 3
current value: 4
current value: 5
current value: 3
current value: 4
current value: 5
current value: 6

range 함수에서 발행하는 0, 1, 2에 1~4를 각각더한 배열을 리턴한다.

순서대로 리턴하는 배열을 나열하면 [1, 2, 3, 4], [2, 3, 4, 5], [3, 4, 5, 6]이다.

배열이기 때문에 비동기로 순서가 바뀌지 않고 값을 발행한다.

/* 유사 배열(ArrayLike)을 배열처럼 취급 */
const { range } = require('rxjs');
const { mergeMap } = require('rxjs/operators');

range(0, 3).pipe(mergeMap(x => {
  const nextArrayLike = {
    length: 4,
    0: x + 1,
    1: x + 2,
    2: x + 3,
    3: x + 4
  };
  console.log(`typeof nextArrayLike: ${typeof nextArrayLike}`);
  return nextArrayLike;
})).subscribe(value => console.log(`current value: ${value}`));

[코드 5-10] 유사 배열(ArrayLike)을 배열처럼 취급

 

실행 결과

typeof nextArrayLike: object
current value: 1
current value: 2
current value: 3
current value: 4
typeof nextArrayLike: object
current value: 2
current value: 3
current value: 4
current value: 5
typeof nextArrayLike: object
current value: 3
current value: 4
current value: 5
current value: 6

length 속성이 Number 타입인 배열과 비슷한 객체도 배열처럼 취급해 리턴할 수 있다.

mergeMap 연산에 프로미스 사용

/* mergeMap 연산자의 프로미스의 생성과 사용 예 */
const { range } = require('rxjs');
const { mergeMap } = require('rxjs/operators');

range(0, 3).pipe(mergeMap(x =>
  new Promise(resolve => setTimeout(() => resolve(`req${x + 1}`),
  Math.floor(Math.random() * 2000)))
)).subscribe(req => console.log(`response from ${req}`));

[코드 5-11] mergeMap 연산자의 프로미스의 생성과 사용 예

 

실행 결과

response from req1
response from req2
response from req3

자바스크립트의 프로미스를 그대로 사용할 수 있으므로 별도로 옵저버블로 변환해줄 필요가 없다. 배열도 마찬가지다.

mergeMap 연산자에 이터러블 사용

subscribeToResult 함수에서 타입을 검사할 때 배열(또는 유사 배열)인지 먼저 검사하고 그 다음 이터러블인지 검사한다.

이터러블이면 해당 이터러블에서 이터레이터를 가져와서 next 함수를 호출해 구독한다.

/* mergeMap 연산자에 이터러블 사용 예 */
const { range } = require('rxjs');
const { mergeMap } = require('rxjs/operators');

range(0, 3).pipe(mergeMap(x => {
  const nextMap = new Map();
  nextMap.set("original", x);
  nextMap.set("plusOne", x + 1);
  return nextMap;
})).subscribe(entry => {
  const [key, value] = entry;
  console.log(`key is ${key}, value is ${value}`);
});

[코드 5-12] mergeMap 연산자에 이터러블 사용 예

 

실행 결과

key is original, value is 0
key is plusOne, value is 1
key is original, value is 1
key is plusOne, value is 2
key is original, value is 2
key is plusOne, value is 3

mergeMap 연산자의 project 함수에서 배열타입을 제외한 이터러블 객체를 리턴하면 이터레이터를 불러와 next 함수를 계속 호출한다. 옵저버의 next 함수로 계속 값을 발행한다. done 이 true 이면 complete 함수를 호출한다.

Map객체란?
map[Symbol, iterator]가 리턴하는 함수를 호출해서 이터레이터를 리턴받을 수 있다.

3-2. mergeMap 연산자의 최대 동시 요청 수 정하기

mergeMap 연산자에는 concurrent 인자가 있다.

동작 방식
이를 사용하면 동시에 최대 구독할 수 있는 옵저버블 수를 제한할 수 있다. 소스 옵저버블에서 발행한 값을 빠른 속도로 전달하더라도 mergeMap 연산자에서 구독 완료하지 않은 옵저버블 수가 concurrent 개수만큼이라면 해당 값을 연산자 내부에 구현해 놓은 버퍼(배열)에 잠시 저장해둔다. 옵저버블 중 하나라도 구독을 해제한다면 버퍼에 저장한 순서대로 값을 하나씩 꺼내어 project 함수에서 새 옵저버블을 만들어 구독한다.

이 방식으로 최대 concurrent 수 만큼 옵저버블 동시 구독을 유지할 수 있다.

 

장점 - 서버와 통신할 때 한 번에 너무 많은 요청을 하지 않으면서 최대 효율을 낼 수 있는 만큼 동시성 유지 가능

 

/* 최대 동시 요청 수 정하기 */
const { range } = require('rxjs');
const { mergeMap } = require('rxjs/operators');
const fetch = require('node-fetch');

const colors = [
  'blue', 'red', 'black', 'yellow', 'green',
  'brown', 'gray', 'purple', 'gold', 'white'
];
const concurrent = 5;
const maxDelayInSecs = 6;
console.time('request_color');

range(0, colors.length).pipe(mergeMap(colorIndex => {
  const currentDelay = parseInt(Math.random() * maxDelayInSecs, 10);
  console.log(
    `[Request Color]: ${colors[colorIndex]}, currentDelay: ${currentDelay}`
  );
  return fetch(
    `https://httpbin.org/delay/${currentDelay}?color_name=${colors[colorIndex]}`
  ).then(res => res.json());
}, concurrent)
).subscribe(response =>
  console.log(
    `<Response> args: ${JSON.stringify(response.args)}, url: ${response.url}`
  ),
  console.error,
  () => {
    console.log('complete!');
    console.timeEnd('request_color');
  }
);

[코드 5-13] 최대 동시 요청 수 정하기

 

실행 결과

# 첫 응답 후 다음 요청 (blue 완료 후 brown 요청)
[Request Color]: blue, currentDelay: 0
[Request Color]: red, currentDelay: 3
[Request Color]: black, currentDelay: 3
[Request Color]: yellow, currentDelay: 2
[Request Color]: green, currentDelay: 5
<Response> args: {"color_name":"blue"}, url: https://httpbin.org/delay/0?color_name=blue
[Request Color]: brown, currentDelay: 1

# yellow 완료 후 gray 요청
<Response> args: {"color_name":"yellow"}, url: https://httpbin.org/delay/2?color_name=yellow
[Request Color]: gray, currentDelay: 1

# 같은 패턴 반복
<Response> args: {"color_name":"brown"}, url: https://httpbin.org/delay/1?color_name=brown
[Request Color]: purple, currentDelay: 3
<Response> args: {"color_name":"red"}, url: https://httpbin.org/delay/3?color_name=red
[Request Color]: gold, currentDelay: 3

# 마지막 white 요청 후 나머지 응답
<Response> args: {"color_name":"black"}, url: https://httpbin.org/delay/3?color_name=black
[Request Color]: white, currentDelay: 0
<Response> args: {"color_name":"gray"}, url: https://httpbin.org/delay/1?color_name=gray
<Response> args: {"color_name":"white"}, url: https://httpbin.org/delay/0?color_name=white
<Response> args: {"color_name":"green"}, url: https://httpbin.org/delay/5?color_name=green
<Response> args: {"color_name":"purple"}, url: https://httpbin.org/delay/3?color_name=purple
<Response> args: {"color_name":"gold"}, url: https://httpbin.org/delay/3?color_name=gold
complete!
request_color: 7.706s

5개 요청을 모두 처리해야 나머지 5개를 요청하는 방식이 아니라, 5개 중 요청 하나라도 먼저 처리하면 다음 처리할 것을 요청해서 항상 최대 5개 요청을 유지한다.

참고 사항
실무에서 대량의 네트워크 요청을 처리할 때는 네트워크 상태나 요청 값에 따라 응답 시간이 다를 수 있다. 이때 요청 수 제한 없이 무작정 많은 연결을 한 번에 만들면 효율이 낮다. 그렇다고 요청 하나의 처리를 완료하고 다음 요청을 처리하면 너무 느리다. 따라서 서버에서 한번에 최대 어느 정도의 요청을 동시에 유지하면 효율이 높을지 측정한 후 적정선을 찾아야 한다.

즉, mergeMap 연산자에 적정한  concurrent 값을 설정해 사용해야 한다는 뜻이다.

 

mergeMap 연산자와 concatMap 연산자의 차이점

  • mergeMap 연산자 - 리턴하는 순서대로 옵저버블을 구독하고 응답은 먼저 전달받는 것부터 처리한다.
  • concatMap 연산자 - 리턴하는 순서대로 옵저버블을 구독하고 먼저 구독한 옵저버블에서 complete 함수를 호출해야만 그 다음 순서로 리턴한 옵저버블을 구독한다.

4. switchMap 연산자

mergeMap 연산자와 switchMap 연산자의 차이점

  • mergeMap - project 함수에서 리턴한 옵저버블을 구독하는 중 소스 옵저버블에서 발행한 값이 있다면 새로 구독하는 옵저버블을 구독한다. 이미 구독하던 옵저버블과 새로 구독하는 옵저버블 모두 함께 동작한다.
  • switchMap - project 함수에서 새 옵저버블을 리턴해 구독하기 전 기존 연산자로 구독하여 완료되지 않은 옵저버블이 있다면 해당 옵저버블의 구독을 해제하고 새 옵저버블을 구독한다.

switchMap 연산자의 마블 다이어그램

연산자 원형
switchMap<T, I, R>(
       project: (value: T, index: number) => ObservableInput<I>,
       resultSelector?: (
              outerValue: T, innerValue: I,
              outerIndex: number,
              innerIndex: number
       ) => R
): OperatorDunction<T, I | R>
  • project 함수 - Observable 클래스의 소스 옵저버블을 리턴
  • resultSelector? - 연산자 원형에 있는 타입을 설정하는 역할
/* switchMap 연산자로 옵저버블 변경 */
const { interval } = require('rxjs');
const { switchMap, take, map } = require('rxjs/operators');

interval(600).pipe(
  take(5),
  switchMap(x =>
    interval(250).pipe(
      map(y => ({x, y})),
      take(3)
    )
  )
).subscribe(result => console.log(`next x: ${result.x}, y: ${result.y}`));

[코드 5-14] switchMap 연산자로 옵저버블 변경

 

실행 결과

next x: 0, y: 0
next x: 0, y: 1
next x: 1, y: 0
next x: 1, y: 1
next x: 2, y: 0
next x: 2, y: 1
next x: 3, y: 0
next x: 3, y: 1
next x: 4, y: 0
next x: 4, y: 1
next x: 4, y: 2

코드 설명

switchMap 연산자의 project 함수에서 take(5)을 적용한 옵저버블을 리턴하면 500ms 동안 250ms 간격으로 두 번만 값을 발행한다. 그리고 600ms 이후 소스 옵저버블에서 새 옵저버블을 구독하므로 기존에 구독중인 값은 기존 구독을 해제하고 750ms 차례의 값은 발행할 수 없다. 그러나 소스 옵저버블의 마지막 값인 4를 발행할 때는 그 다음 구독할 옵저버블이 없으므로 0~2라는 y 값을 모두 발행한다.


5. concatMap 연산자

concatMap 연산자는 project 함수에서 리턴하는 옵저버블을 구독한 후 값 발행을 완료해야 다음 옵저버블을 구독하는 연산자다.

동작 방식
이미 구독중인 옵저버블의 값 발행을 완료하기 전에 다른 옵저버블에서 발행하는 값은 버퍼에 오는 순서대로 잠시 저장해둔다. 그리고 구독 중인 옵저버블의 값 발행을 완료하면 버퍼에서 저장한 값을 꺼내서 project 함수로 다음 옵저버블을 구독하는 일을 반복한다. 그러므로 값 발행 순서를 보장한다.

concatMap 연산자의 마블 다이어그램

연산자 원형
public concatMap(
       project: function(value: T, index: number): ObservableInput,
       resultSelector?: function(
              outerValue: T, innerValue: I, outerIndex: number, innerIndex: number
       ): any
): Observable
  • project 함수 - Observable 클래스의 소스 옵저버블을 리턴
  • resultSelector? - 연산자 원형에 있는 타입을 설정하는 역할
/* concatMap 의 구현 코드 일부분 */
import {mergeMap} from "rxjs/operators";

export function concatMap(project, resultSelector) {
  return mergeMap(project, resultSelector, 1);
}

[코드 5-15] concatMap 의 구현 코드 일부분

 

[코드5-15]를 보면 concatMap 연산자는 mergeMap 연산자의 concurrent 값을 1로 설정한 동작을 추상화한 연산자다. 즉, 최대 구독 할 수 있는 옵저버블 수가 1이라는 뜻이다.

/* concatMap 연산자로 값 발행 순서를 보장 */
const { timer, interval, range } = require('rxjs');
const { concatMap, take, map } = require('rxjs/operators');

const requests = [
  timer(2000).pipe(map(value => 'req1')),
  timer(1000).pipe(map(value => 'req2')),
  timer(1500).pipe(map(value => 'req3'))
];

interval(1000).pipe(take(5))
  .subscribe(x => console.log(`${x + 1} 초 경과`));

range(0, 3).pipe(concatMap(x => requests[x]))
  .subscribe(req => console.log(`response from ${req}`));

[코드 5-16] concatMap 연산자로 값 발행 순서를 보장

 

실행 결과

1 초 경과
response from req1
2 초 경과
response from req2
3 초 경과
4 초 경과
response from req3
5 초 경과

서로 다른 시간으로 동작하지만 값 발행 순서가 정해진 것을 확인할 수 있다.

 

동기 환경에서는 mergeMap 연산자나 concatMap 연산자나 동작은 같다. 기존 구독하는 옵저버블이 동기 방식으로 실행되므로 그 다음에 next 함수가 발행하는 값을 순서대로 처리할 수 있기 때문이다.

/* 비동기 처리의 코드 동작 순서 */
const { timer, interval, range } = require('rxjs');
const { startWith, tap, skip, map, take, concatMap } = require('rxjs/operators');

const FIRST_VALUE = -1;
const requests = [
  timer(2000).pipe(
    startWith(FIRST_VALUE),
    tap(x => x === FIRST_VALUE && console.log('req1 구독')),
    skip(1),
    map(value => 'req1')
  ),
  timer(1000).pipe(
    startWith(FIRST_VALUE),
    tap(x => x === FIRST_VALUE && console.log('req2 구독')),
    skip(1),
    map(value => 'req2')
  ),
  timer(1500).pipe(
    startWith(FIRST_VALUE),
    tap(x => x === FIRST_VALUE && console.log('req3 구독')),
    skip(1),
    map(value => 'req3')
  )
];

interval(1000).pipe(take(5))
  .subscribe(x => console.log(`${x + 1} 초 경과`));

range(0, 3).pipe(
  tap(x => console.log(`range next 함수 ${x}`)),
  concatMap(x =>
    console.log(`시작 - concatMap 연산자의 project function ${x}`) || requests[x]
  )
).subscribe(req => console.log(`response from ${req}`));

[코드 5-17] 비동기 처리의 코드 동작 순서

 

실행 결과

range next 함수 0
시작 - concatMap 연산자의 project function 0
req1 구독
range next 함수 1
range next 함수 2
1 초 경과
2 초 경과
response from req1
시작 - concatMap 연산자의 project function 1
req2 구독
3 초 경과
response from req2
시작 - concatMap 연산자의 project function 2
req3 구독
4 초 경과
response from req3
5 초 경과

[코드 5-17] 동작 순서

  1. startWith 연산자로 구독할 때의 첫 값을 발행하도록 함
  2. tap 연산자로 로그를 출력
  3. FIRST_VALUE는 구독 확인용 더미 값 - 구독이 필요 없어 skip(1)을 사용
  4. range 함수 안의 tap 연산자는 concatMap 연산잔의 소스 옵저버블이 언제 다음 값을 발행하는 지 확인의 목적
  5. project 함수를 언제 호출하는지 알기 위해 논리 연산자인 || 를 이용 req[x] 를 리턴하기 전에 로그 출력

주의 할 점 - concurrent 값이 1이므로 해당 값이 없거나 더 큰 mergeMap 연산자 보다는 메모리에 저장해야 할 값의 수가 많다.


6. scan 연산자

 scan 연산자와 reduce 연산자의 차이점

  • scan 연산자 - next 함수로 값을 발행할 때마다 호출해 중간에 누적된 값을 매번 발행한다.
  • reduce 연산자 - 최종 누적된 값 1개만 발행한다.
연산자 원형
scan<T, R>(
       accumulator: (acc: R, value: T, index: number) => R,
       seed?: T | R
): OperatorFunction<T, R>
  • accumulator 함수 - 값 누적을 어떻게 할지 설정하는 함수 (누적자 함수)
  • seed? - 처음 발행하는 값부터 누적하려면 초기값이 필요 (seed? = 초기값)

scan 연산자의 마블 다이어그램

첫 번째로 발행하는 값은 그냥 건너뛰고, 이후 발행하는 새로운 값을 계속 누적시켜 다음 새 값을 발행한다. 누적하는 방법은 객체 추가, 기본 숫자 타입의 사칙 연산 등 다양하다.

6-1. 초기값이 없는 scan 연산자 예

/* 초기값이 없는 scan 연산자 예 */
const { range } = require('rxjs');
const { scan } = require('rxjs/operators');

range(0, 3).pipe(
  scan((accumulation, currentValue) => {
    console.log(`accumulation: ${accumulation}, currentValue: ${currentValue}`);
    return accumulation + currentValue;
  })
).subscribe(result => console.log(`result: ${result}`));

[코드 5-18] 초기값이 없는 scan 연산자 예

 

실행 결과

result: 0
accumulation: 0, currentValue: 1
result: 1
accumulation: 1, currentValue: 2
result: 3

누적자 함수 1개만 존재, 초기값은 보통 두 번째 인자, 처음 발행하는 값은 누적자 함수를 거치지 않고 바로 발행, 이후 부터 누적자 함수로 더한 값을 누적 변수(accumulation) 에 계속 저장한다.

6-2. 초기값이 있는 scan 연산자 예

/* scan 연산자의 초기값을 0로 설정한 예 */
const { range } = require('rxjs');
const { scan } = require('rxjs/operators');

range(0, 3).pipe(
  scan((accumulation, currentValue) => {
    console.log(`accumulation: ${accumulation}, currentValue: ${currentValue}`);
    return accumulation + currentValue;
  }, 0)
).subscribe(result => console.log(`result: ${result}`));

[코드 5-19] scan 연산자의 초기값을 0로 설정한 예

 

실행 결과

accumulation: 0, currentValue: 0
result: 0
accumulation: 0, currentValue: 1
result: 1
accumulation: 1, currentValue: 2
result: 3

초기값을 설정하면 누적함수를 건너뛰지 않고 누적함수를 거쳐서 발행한다는 차이점을 볼 수 있다.

주의 할 점
초기값이 기본 타입의 값이면 변경할 수 없는 값이다. 다시 구톡할 때 초기값이 변하지 않는다.
초기값이 참조 타입의 값이면 어디서든 변경할 수 있는 값이므로 초기값이 변할 수 있다. 
/* 초기값으로 객체를 사용해 재구독하는 피보나치 수열 예 */
const { interval } = require('rxjs');
const { take, scan, pluck } = require('rxjs/operators');
const n = 7;

const source$ = interval(500).pipe(
  take(n),
  scan((accumulation, currentValue) => {
    const tempA = accumulation.a;
    accumulation.a = accumulation.b;
    accumulation.b = tempA + accumulation.b;
    return accumulation;
  }, { a: 1, b: 0 }),
  pluck('a')
);

source$.subscribe(result => console.log(`result1: ${result}`));
setTimeout(() =>
  source$.subscribe(result =>
    console.log(`result2: ${result}`)
  ), 3100
);

[코드 5-20] 초기값으로 객체를 사용해 재구독하는 피보나치 수열 예

 

실행 결과

result1: 0
result1: 1
result1: 1
result1: 2
result1: 3
result1: 5
result1: 8
result2: 13
result2: 21
result2: 34
result2: 55
result2: 89
result2: 144
result2: 233

두 번째 옵저버블을 구독했을 때 0이 아닌 13부터 발행하는 문제가 있다. 구독 할 때마다 이전 옵저버블에 영향을 받아 첫 구독하는 값이 매번 달라진다면 여러 옵저버블이 같은 객체를 참조하여 부수 효과를 유발 할 수 있다.

 

해결 방법 - 누적자 함수 안에서 구독할 때마다 초기값을 새로 생성하면 된다.

/* 팩토리 함수를 이용해 초기값 구분 */
const { interval } = require('rxjs');
const { take, scan, pluck } = require('rxjs/operators');
const n = 7;

const source$ = interval(500).pipe(
  take(n),
  scan((accumulation, currentValue) => {
    let localAccumulation = accumulation;
    if (typeof accumulation === 'function') {
      localAccumulation = localAccumulation();
    }
    const tempA = localAccumulation.a;
    localAccumulation.a = localAccumulation.b;
    localAccumulation.b = tempA + localAccumulation.b;
    return localAccumulation;
  }, () => ({ a: 1, b: 0 })),
  pluck('a')
);

source$.subscribe(result => console.log(`result1: ${result}`));
setTimeout(() =>
  source$.subscribe(result =>
    console.log(`result2: ${result}`)
  ), 3100
);

[코드 5-21] 팩토리 함수를 이용하여 초기값 구분

 

실행 결과

result1: 0
result1: 1
result1: 1
result1: 2
result1: 3
result1: 5
result1: 8
result2: 0
result2: 1
result2: 1
result2: 2
result2: 3
result2: 5
result2: 8

객체를 리턴하는 팩토리 함수로 초기값을 만들고 누적자 함수에서 전달하는 값이 함수인지 아니지로 초기값을 구분한다. 매번 리턴하는 값은 구독할 때마다 팩토리 함수로 생성해 구분한 객체가 된다. 초기값 자체는 팩토리 함수로 구독할 때마나 새로 생성할 수 있다.

 

누적자 함수는 localAccumulation를 사용해 전달받는 accumulation이 객체이면 그대로 참조하고, 처음 시작하는 팩토리 함수이면 이를 호출해서 리턴하는 새 객체를 전달 받는다.

 

문제점 - 매번 구독할 때만 다른 객체를 생성할 뿐 구독하는 중에 누적자 함수에서 리턴하는 객체는 같다. 해당 객체를 구독하는 중 수정하면 누적자 객체의 값이 의도와는 다르게 수정되는 부수 효과 발생 우려가 있다.

 

해결 방법 - 매 번 새 객체를 생성해 누적자 함수에서 리턴하면 된다. 초기값을 확인할 필요 없이 매번 새 값을 리턴한다는 점에서 변경 가능한 객체를 재사용하지 않는 장점이 있다. 하지만, 새 객체를 누적자 함수 호출 때마다 매번 생성해야한다는 단점이 있다. 

/* 새 객체를 매번 생성해 누적자 함수에서 리턴하는 예 */
const { interval } = require('rxjs');
const { take, scan, pluck } = require('rxjs/operators');
const n = 7;

const source$ = interval(500).pipe(
  take(n),
  scan((accumulation, currentValue) => ({
    a: accumulation.b,
    b: accumulation.a + accumulation.b
  }), { a: 1, b: 0 }),
  pluck('a')
);

source$.subscribe(result => console.log(`result1: ${result}`));
setTimeout(() =>
  source$.subscribe(result =>
    console.log(`result2: ${result}`)
  ), 3100
);

[코드 5-22] 새 객체를 매번 생성해 누적자 함수에서 리턴하는 예

 

실행 결과

result1: 0
result1: 1
result1: 1
result1: 2
result1: 3
result1: 5
result1: 8
result2: 0
result2: 1
result2: 1
result2: 2
result2: 3
result2: 5
result2: 8

accunulation 을 조작할 필요가 없으므로 특정 값을 저장할 임시 변수가 필요 없이 화살표 함수를 더 간결하게 사용할 수 있다는 장점도 있다. 하지만, 그만큼 객체를 생성하는 데 자원을 할당해야 한다.

6-3. 누적자 함수의 index 제공

현재 누적자 함수에서 소스 옵저버블의 몇 번째 발행한 값을 전달하는 지 알아야 한다면 누적자 함수 세 번째 파라미터로 index 를 추가하여 로그를 출력하면 된다.


7. partition 연산자

partition 연산자는 predicate 함수를 호출하면 2개의 옵저버블을 배열로 리턴한다.

  • 첫 번쩨 배열 요소 - predicate 함수의 조건을 만족하는 filter 연산자를 적용한 옵저버블
  • 두 번째 배열 요소 - predicate 함수의 조건을 만족하지 않는 filter 연산자를 적용한 옵저버블

partition 연산자의 마블 다이어그램

연산자 원형
partition<T>(
       predicate: (value: T, index: number) => boolean,
       thisArg?: any
): unaryFunction<Observable<T>, [Observable<T>, Observable<T>]>
  • thisArgs? - predicate 함수에서 리턴하는 값을 확인
  • UnaryFunction<Observable<T>, [Observable<T>, Observable<T>]> - 2개의 옵저버블을 배열로 리턴

중요한 점 - 연산자 사용 후 바로 다른 연산자를 붙이거나 구독할 수 없다. 배열에서 각 옵저버블을 꺼내 별도로 처리해주어야 한다.

/* partition 연산로 소스 옵저버블을 독립적으로 두는 예 */
const { interval } = require("rxjs");
const { partition, take, map } = require("rxjs/operators");

const [winSource$, loseSource$] = interval(500).pipe(
  partition(x => Math.random() < 0.7)
);

winSource$.pipe(
  map(x => `당첨!! (${x})`),
  take(10)
).subscribe(result => console.log(`win: ${result}`));

loseSource$.pipe(
  map(x => `꽝!! (${x})`),
  take(10)
).subscribe(result => console.log(`lose: ${result}`));

[코드 5-23] partition 연산자로 소스 옵저버블을 독립적으로 두는 예

 

실행 결과

win: 당첨!! (1)
win: 당첨!! (2)
win: 당첨!! (3)
win: 당첨!! (4)
win: 당첨!! (6)
win: 당첨!! (7)
lose: 꽝!! (7)
win: 당첨!! (8)
win: 당첨!! (9)
win: 당첨!! (11)
lose: 꽝!! (11)
win: 당첨!! (12)
lose: 꽝!! (13)
lose: 꽝!! (16)
lose: 꽝!! (17)
lose: 꽝!! (18)
lose: 꽝!! (19)
lose: 꽝!! (20)
lose: 꽝!! (30)
lose: 꽝!! (35)
  • 비구조화 할당 - 오른쪽 배열에 있는 값을 왼쪽 배열의 요소로 포함한 변수 각각에 순서대로 할당

8. groupBy 연산자

groupBy 연산자는 소스 옵저버블에서 발행하는 값을 특정 기준을 정해 같은 그룹에 속해 있는 값들을 각각의 옵저버블로 묶어서 발행한다.

  • key - 그룹을 묶는 기준
  • keySelector - 각 값에서 키 값을 만드는 함수
  • gropuedObservable - 같은 그룹에 속한 (키 값이 같은) 값들을 발행하는 옵저버블

각 groupedObservable 마다 적절한 연산자를 결합해서 사용 할 수 있다.

groupBy 연산자의 마블 다이어그램

소스 옵저버블에서 i => i % 2 라는 keySelector로 true 와 false 2개의 그룹을 만든다. 각 그룹에 해당하는 값은 값과 같은 그룹에 속한 옵저버블에서 발행하는 것을 확인 할 수 있다.

연산자 원형
groupBy<T, K, R>(
       keySelector: (value: T) => K,
       elementSelector?: ((value: T) => R) | void,
       durationSelector?: (grouped: GroupedObservable<K, R>) => Observable<any>,
       subjectSelector?: () => Subject<R>
): OperatorFunction<T, GroupedObservable<K,R>>

8-1. keySelector 함수

/* groupBy 연산자로 당첨과 꽝을 출력하는 예 */
const { interval } = require('rxjs');
const { groupBy, take, map, mergeMap } = require('rxjs/operators');

// 1. keySelector 사용
interval(500).pipe(
  take(10),
  groupBy(x => Math.random() < 0.7),
  mergeMap(
    groupedObservable =>
      groupedObservable.key === true ?
        groupedObservable.pipe(map(x => `당첨!! (${x})`)) :
        groupedObservable.pipe(map(x => `꽝!! (${x})`))
  )
).subscribe(result => console.log(result));

[코드 5-24] groupBy 연산자로 당첨과 꽝을 출력하는 예

 

실행 결과

꽝!! (0)
꽝!! (1)
당첨!! (2)
꽝!! (3)
꽝!! (4)
꽝!! (5)
당첨!! (6)
당첨!! (7)
당첨!! (8)
당첨!! (9)

keySelector 함수는 만들기에 따라 여러 개 키를 생성해 여러 그룹을 만들 수 있다.

 

groupedObservable - 키에 해당하는 옵저버블이 없어 만들어서 값을 발행하는 인스턴스

groupedObservable 옵저버블 인스턴스는 내부 this.key 에 키를 저장한다.

8-2. elementSelector 함수

elementSelector 함수는 keySelector 함수로 생성한 groupedObservable 인스턴스에 전달하는 값을 바꿔서 리턴한 후 키에 해당하는  groupedObservable 로 값을 발행한다.

/* elementSelector 함수를 추가한 예 */
const { interval } = require('rxjs');
const { take, groupBy, mergeMap, map } = require('rxjs/operators');

// 2. keySelector, elementSelector 사용
interval(500).pipe(
  take(10),
  groupBy(
    x => Math.random() < 0.7,
    x => `${x}-${x % 2 === 0 ? '짝수' : '홀수'}`
  ),
  mergeMap(
    groupedObservable =>
      groupedObservable.key === true ?
        groupedObservable.pipe(map(x => `당첨!! (${x})`)) :
        groupedObservable.pipe(map(x => `꽝!! (${x})`))
  )
).subscribe(result => console.log(result));

[코드 5-25] elementSelector 함수를 추가한 예

 

실행 결과

당첨!! (0-짝수)
당첨!! (1-홀수)
꽝!! (2-짝수)
당첨!! (3-홀수)
꽝!! (4-짝수)
꽝!! (5-홀수)
당첨!! (6-짝수)
당첨!! (7-홀수)
당첨!! (8-짝수)
당첨!! (9-홀수)

elementSelector 함수는 해당 값이 짝수인지 홀수인지 검사한 후 이를 나타내는 문자열을 추가해 리턴한다.

 

keySelector 함수는 elementSelector 함수에서 리턴하는 값이 아닌 소스 옵저버블에서 발행하는 값을 기준으로 한다. elementSelector 함수에서 리턴하는 값은 groupedObservable 에 전달하려는 용도의 값 (element) 이다.

8-3. durationSelector 함수

durationSelector 함수의 특징

  • groupedObservable를 사용하는 함수며 여기서 리턴하는 옵저버블은 키에 해당하는 새 groupedObservable을 발행할 때 함께 구독한다.
  • 여기서 어떤 값이든 발행하는 시점에 키와 groupedObservable의 맵핑을 끊고 complete 함수를 호출해 groupedObservable의 구독을 완료한다.
  • 이 때 duartionSelector 함수에서 리턴한 옵저버블의 구독도 완료한다.
  • 이 후 같은 키에 해당하는 값을 소스 옵저버블에서 발행하면 키에 관련 맵핑이 없으므로 durationSelector 함수에서 리턴하는 옵저버블을 새로 맵핑해 구독한다.
/* reduce 연산자를 추가한 당첨 확률 예 */
const { interval } = require('rxjs');
const { take, groupBy, mergeMap, reduce, map } = require('rxjs/operators');

// 3. keySelector, elementSelector, reduce 사용
interval(500).pipe(
  take(10),
  groupBy(
    x => Math.random() < 0.7,
    x => `${x}-${x % 2 === 0 ? '짝수' : '홀수'}`
  ),
  mergeMap(
    groupByObservable =>
      groupByObservable.key === true ?
        groupByObservable.pipe(
          map(x => `당첨!! (${x})`),
          reduce((acc, curr) => [...acc, curr], [])
        ) :
        groupByObservable.pipe(
          map(x => `꽝!! (${x})`),
          reduce((acc, curr) => [...acc, curr], [])
        )
  )
).subscribe(result => console.log(result));

[코드 5-26] reduce 연산자를 추가한 당첨 확률 예

 

실행 결과

[
  '당첨!! (0-짝수)',
  '당첨!! (1-홀수)',
  '당첨!! (2-짝수)',
  '당첨!! (3-홀수)',
  '당첨!! (6-짝수)',
  '당첨!! (7-홀수)',
  '당첨!! (8-짝수)',
  '당첨!! (9-홀수)'
]
[ '꽝!! (4-짝수)', '꽝!! (5-홀수)' ]

500ms 간격으로 10번 값을 발행하는 시간인 5초를 기다려야 최종 결과를 확인할 수 있다.

mergeMap 연산자에서 옵저버블을 구독할 때 배열로 모든 값을 누적하는 reduce 연산자를 사용했다.

따라서, 당첨은 당첨끼리 모아 배열 하나를 만들고, 꽝은 꽝대로 모아 배열 하나를 만든다.

/* durationSelector 함수를 사용하는 예 */
const { interval } = require('rxjs');
const { take, groupBy, tap, mergeMap, map, reduce } = require('rxjs/operators');

// 4. keySelector, elementSelector, durationSelector, reduce 사용
interval(500).pipe(
  take(10),
  groupBy(
    x => Math.random() < 0.7,
    x => `${x}-${x % 2 === 0 ? '짝수' : '홀수'}`,
    groupedObservable =>
      groupedObservable.key === true ?
        interval(600).pipe(
          tap(x => console.log(`당첨 duration ${x}`))
        ) :
        interval(2000).pipe(
          tap(x => console.log(`꽝 duration ${x}`))
        )
  ),
  mergeMap(
    groupByObservable =>
      groupByObservable.key === true ?
        groupByObservable.pipe(
          map(x => `당첨!! (${x})`),
          reduce((acc, curr) => [...acc, curr], [])
        ) :
        groupByObservable.pipe(
          map(x => `꽝!! (${x})`),
          reduce((acc, curr) => [...acc, curr], [])
        )
  )
).subscribe(result => console.log(result));

[코드 5-27] durationSelector 함수를 사용하는 예

 

실행 결과

당첨 duration 0
[ '당첨!! (0-짝수)', '당첨!! (1-홀수)' ]
당첨 duration 0
[ '당첨!! (2-짝수)' ]
당첨 duration 0
[ '당첨!! (4-짝수)', '당첨!! (5-홀수)' ]
꽝 duration 0
[ '꽝!! (3-홀수)' ]
당첨 duration 0
[ '당첨!! (6-짝수)' ]
[ '꽝!! (7-홀수)', '꽝!! (8-짝수)', '꽝!! (9-홀수)' ]

durationSelector 함수는 당첨 결과를 600ms 간격으로, 꽝인 결과를 2000ms 간격으로 매핑을 끊도록 옵저버블을 리턴한다. 단, 해당 옵저버블에서 값을 잘 발행했다는 것을 확인하려고 tap 연산자를 사용하여 소스 옵저버블에서 값을 전달받아 로그를 출력했다.

 

각 durationselector 함수에서 리턴한 옵저버블이 값을 발행할 때마다 duartion이 있는 로그를 출력한다. 해당 옵저버블의 값ㄷ고 매번 0이다. interval 함수임에도 매번 옵저버블을 새로 구독하기 때문이다.

 

키 포인트 - durationSelector 함수가 각 groupObservable 을 일정 주기마다 새로 그룹으로 묶어 구독해야 할때 사용하는 기능이다.

8-4. subjectSelector 함수

/* subjectSelector 함수의 사용 예 */
const { interval, BehaviorSubject } = require('rxjs');
const { take, groupBy, tap, mergeMap, map, reduce } = require('rxjs/operators');

// 5. keySelector, elementSelector, durationSelector, subjectSelector, reduce 사용
interval(500).pipe(
  take(10),
  groupBy(
    x => Math.random() < 0.7,
    x => `${x}-${x % 2 === 0 ? '짝수' : '홀수'}`,
    groupedObservable =>
      groupedObservable.key === true ?
        interval(600).pipe(
          tap(x => console.log(`당첨 duration ${x}`))
        ) :
        interval(2000).pipe(
          tap(x => console.log(`꽝 duration ${x}`))
        ),
    () => new BehaviorSubject('GROUP START')
  ),
  mergeMap(
    groupByObservable =>
      groupByObservable.key === true ?
        groupByObservable.pipe(
          map(x => `당첨!! (${x})`),
          reduce((acc, curr) => [...acc, curr], [])
        ) :
        groupByObservable.pipe(
          map(x => `꽝!! (${x})`),
          reduce((acc, curr) => [...acc, curr], [])
        )
  )
).subscribe(result => console.log(result));

[코드 5-28] subjectSelector 함수의 사용 예

 

실행 결과

당첨 duration 0
[ '당첨!! (GROUP START)', '당첨!! (1-홀수)', '당첨!! (2-짝수)' ]
꽝 duration 0
[ '꽝!! (GROUP START)', '꽝!! (0-짝수)' ]
당첨 duration 0
[ '당첨!! (GROUP START)', '당첨!! (3-홀수)' ]
당첨 duration 0
[ '당첨!! (GROUP START)', '당첨!! (5-홀수)', '당첨!! (6-짝수)' ]
꽝 duration 0
[ '꽝!! (GROUP START)', '꽝!! (4-짝수)', '꽝!! (7-홀수)' ]
[ '당첨!! (GROUP START)', '당첨!! (8-짝수)', '당첨!! (9-홀수)' ]

subjectSelector 함수는 RxJS 에서 제공하는 서브젝트 타입의 인스턴스를 리턴할 수 있는 함수다.

GROUP START 라는 값을 생성자에게 전달받는 BehaviorSubject의 객체를 생성하도록 했다.
BehaviorSubject 는 초기값을 인자로 사용할 수 있다.

groupBy 연산자는 groupedObservable을 구독할 때 매핑하는 서브젝트도 같이 구독하며, 이 서브젝트에서 발행하는 값을 groupedObservable로 전달한다. 이 때 BehaviorSubject는 최초 옵저버블 구독 시점에 초기값을 발행하므로 당첨이든, 꽝이든 초기값으로 GROUP START를 무조건 발행한다. 실행 결과의 모든 배열마다 GROUP START를 처음에 포함한 것을 확인 할 수 있다.

Subject 객체의 특성
1. 멀티 캐스팅을 지원
2. Subject 객체를 생성해 groupedObservable 과 연결

키 포인트 - subjectSelector 함수에서 리턴하는 Subject 객체와 연결하는 특징이 있다.

8-5. groupObservable 첫 생성 시 최초 발행 값을 전달하는 시점

의문점 - 최초로 해당 옵저버블의 값을 전달 받아서 구독하기 전에 값을 발행하면 이 값을 전달 받을 수 없기 때문에 최초 발행 값을 전달하는 시점이 의문

 

의문점에 대한 해답 - 최초 키에 해당하는 옵저버블은 처리를 완료한 후 소스 옵저버블의 값을 전달한다. 따라서, 처음 옵저버블에 발행한 값을 전달받아 이를 처리하는 단계에서 구독한다면 그룹에 해당하는 옵저버블을 구독하기 전 값을 전달받는 일은 발생하지 않는다. 아래 구현코드의 일부를 확인해 보자.

/* groupBy 연산자의 구현 코드 일부분 */
// 생략
if (!groups) {
  groups = this.groups = typeof key === 'string' ? new FastMap() : new Map();
}
let group = groups.get(key);

// 생략

if (!group) {
  group = this.subjectSelector ? this.subjectSelector() : new Subject();
  groups.set(key, group);
  const groupedObservable = new GroupedObservable(key, group, this);
  this.destination.next(groupedObservable);
  // 생략
}

if (!group.closed) {
  group.next(element);
}

[코드 5-29] groupBy 연산자의 구현 코드 일부분


9. buffer 연산자

buffer 연산자는 소스 옵저버블에서 발행하는 값을 순서대로 일정 기준으로 묶어서 하나의 배열로 발행하는 연산자다. 묶을 기준에 해당하는 시점에 배열을 만들어 발행하는 값들을 해당 배열에 쌓아둔 후, 일정 조건을 충족했을 때 쌓아둔 배열을 발행하고 다시 개로 배열을 쌓고 발행하는 일을 반복한다.

buffer 연산자의 마블 다이어그램

연산자 원형
buffer<T>(closingNotifier: Observable<any>): OperatorFunction<T, T[]>

closingNotifier 옵저버블을 구독한 후 closingNotifier 에서 값을 발행할 때까지 소스 옵저버블에서 발행하는 값들을 배열에 순서대로 값을 쌓아둔다. closingNotifier 에서 값을 발행할 때 그 동안 값을 쌓아둔 배열을 발행하고 이후 소소 옵저버블에서 발해되는 값은 새로운 배열에 다시 누적한다.

 

옵저버블 스트림에서 발행하는 값을 묶어서 적정한 때 처리할 때 유용한다.

/* 메시지를 묶어 배열로 출력 */
const { interval } = require('rxjs');
const { take, map, buffer } = require('rxjs/operators');
const message = '안녕하세요? RxJS 테스트 입니다';

interval(90).pipe(
  take(message.length),
  map(x => {
    const character = message.charAt(x);
    console.log(character);
    return character;
  }),
  buffer(interval(500))
).subscribe(x => console.log(`buffer: ${x}`));

[코드 5-30] 메시지를 묶어 배열로 출력

 

실행 결과

안
녕
하
세
요
buffer: 안,녕,하,세,요
?
 
R
x
J
buffer: ?, ,R,x,J
S
 
테
스
트
 
buffer: S, ,테,스,트, 
입
니
다

'입니다'는 take 연산자로 다음 500ms 가 지나기 전 complete 함수를 호출하므로 버퍼에 저장하지 못한 것을 확인 할 수 있다.

 

소스 옵저버블에서 complete 함수를 호출하면 소스 옵저버블의 값을 누적해 버퍼에 저장한 값을 발행하지 않는다.

10. bufferCount 연산자

  • 값 각각을 옵저버블 스트림으로 전달할 때 이를 일정 개수만큼 묶은 후 서버에 한 번의 요청을 하는 상황
  • 묶은 값을 풀어서 처리해야 할 상황

bufferCount 연산자의 마블 다이어그램

buffer(3, 2) 는 3개씩 값을 묶어 2칸 시프트 이동한다는 뜻이다. 구독 완료 후 [i]를 발행한다.

 

bufferCount 연산자는 buffer 연산자와는 달리 소스 옵저버블에서 complete 함수를 호출해도 버퍼에 저장한 값이 있으면 발행한다.

연산자 원형
bufferCount<T>(
       bufferSize: number,
       startBufferEvery: number = null
): OperatorDunction<T, T[]>
  • bufferSize - 몇 개 값을 버퍼에 저장해 묶을지 정하는 정수값을 설정
  • startBufferEvery - 버퍼 크기 만큼 값을 저장한 후 얼마나 시프트로 이동해서 값을 자를 것인지 설정 (null 이면 버퍼 크기 만큼 시프트)
/* 특정 수의 값을 정확하게 묶는 bufferCount 연산자 */
const { interval } = require('rxjs');
const { take, map, bufferCount } = require('rxjs/operators');
const message = '안녕하세요? RxJS 테스트 입니다';

interval(90).pipe(
  take(message.length),
  map(x => {
    const character = message.charAt(x);
    console.log(character);
    return character;
  }),
  bufferCount(5)
).subscribe(x => console.log(`buffer: ${x}`));

[코드 5-31] 특정 수의 값을 정확하게 묶는 bufferCount 연산자

 

실행 결과

안
녕
하
세
요
buffer: 안,녕,하,세,요
?
 
R
x
J
buffer: ?, ,R,x,J
S
 
테
스
트
buffer: S, ,테,스,트
 
입
니
다
buffer:  ,입,니,다

bufferCount(5)를 적용하여 다섯 글자씩을 묶었으며 소스 옵저버블에서 complete 함수를 호출한 후 버퍼에 저장한 값이 5개가 안되어도 버퍼에 저장한 값을 발행한다.

10-1. startBufferEvery 파라미터

/* 찾으려는 단어 수만큼 묶고 한 칸씩 시프트 이동 */
const { from } = require('rxjs');
const { bufferCount, filter, map } = require('rxjs/operators');
const message = '간장공장공장장은강공장장이고공공장공장장은장공장장이다';
const targetWord = '공장장';

from(message).pipe(
  bufferCount(targetWord.length, 1),
  filter(buffer => buffer.length === targetWord.length),
  map(buffer => {
    const bufferedWord = buffer.join('');
    console.log(`buffer: ${bufferedWord}`);
    return bufferedWord;
  }),
  filter(word => word === targetWord)
).subscribe(word => console.log(`${word} 발견`));

[코드 5-32] 찾으려는 단어 수만큼 묶고 한 칸씩 시프트 이동

 

실행 결과

buffer: 간장공
buffer: 장공장
buffer: 공장공
buffer: 장공장
buffer: 공장장
공장장 발견
buffer: 장장은
buffer: 장은강
buffer: 은강공
buffer: 강공장
buffer: 공장장
공장장 발견
buffer: 장장이
buffer: 장이고
buffer: 이고공
buffer: 고공공
buffer: 공공장
buffer: 공장공
buffer: 장공장
buffer: 공장장
공장장 발견
buffer: 장장은
buffer: 장은장
buffer: 은장공
buffer: 장공장
buffer: 공장장
공장장 발견
buffer: 장장이
buffer: 장이다
  1. from 함수 - message를 한 글자씩 나눠 발행하도록 한다.
  2. bufferCount 연산자 - 찾는 글자 수 (targetWord.length) 만큼 묶어 1칸식 시프트 이동하도록 한다.
  3. 첫번째 filter 연산자 - targetWord와 글자 수가 같은 것만 처리한다.
  4. join 연산자 -  묶인 배열을 문자열로 만든다.
  5. map 연산자 - 변환, 해당 버퍼 값을 로그로 출력한다.
  6. 두번째 filter 연산자 - 문자열로 만든 버퍼 값과 targetWord 가 같은 값만 처리한다.
  7. 구독하며 찾고자 하는 단어가 있으면 로그 출력한다.

bufferCount 연산자로 설정한 버퍼 크기와 숫자만큼의 칸을 시ㅍ프트로 이동할 수 있는 패턴은 슬라이딩 윈도우를 구현할 때 매우 유용

슬라이딩 윈도우 란?
일정 크기의 윈도우를 만들어 특정 수만큼 움직이며 검색할 수 있는 알고리즘이다. 전체 영역에서 특정 구간의 가장 큰 값이나 가장 작은 값 또는 특정 조건을 만족하는 영역을 찾을 때 등에 사용할 수 있다.

11. window 연산자

window 연산자는 buffer 연산자의 배열 대신 중첩 옵저버블로 값을 발행하는 연산자다.

연산자 원형
window<T>(windowBoundaries: Observable<any>): OperatorFunction<T, Observable<T>>

windowBoundaries 는 묶는 단위의 값을 발행하는 옵저버블이다.

window 연산자의 마블 다이어그램

window 연산자는 인자로 사용하는 옵버버블에서 그때그때 값을 발행한다. 따라서 소스 옵저버블의 구독을 완료하기 전에도 발행한 값을 전달 받을 방법이 있다. 인자로 사용하는 새 옵저버블은 소스 옵저버블의 구독을 완료할 때까지 값을 발행할 수 있는 옵저버블이기 때문이다.

/* window 연산자의 사용 예 */
const { interval } = require('rxjs');
const { take, map, window, concatMap, filter, scan, last } = require('rxjs/operators');
const message = '안녕하세요. RxJS 테스트 입니다';

interval(90).pipe(
  take(message.length),
  map(x => {
    const character = message.charAt(x);
    console.log(character);
    return character;
  }),
  window(interval(500)),
  concatMap(windowObservable => {
    console.log('windowObservable 전달받음');
    return windowObservable.pipe(
      filter(x => x != ' '),
      take(3),
      scan((accString, current) => accString + current, ''),
      last()
    );
  })
).subscribe(string => console.log(`결과: ${string}`));

[코드 5-33] window 연산자의 사용 예

 

실행 결과

windowObservable 전달받음
안
녕
하
결과: 안녕하
세
요
windowObservable 전달받음
.
 
R
x
결과: .Rx
J
S
windowObservable 전달받음
 
테
스
트
결과: 테스트
 
windowObservable 전달받음
입
니
다
결과: 입니다

window 연산자를 활용하면 인자로 사용하는 옵저버블에서 값을 발행할 때까지 기다릴 필요 없이 바로 소스 옵저버블에서 값을 전달 받을수 있고 원하는 결과가 일찍 나오면 take 같은 연산자를 이용해서 값 발행을 중지할 수도 있다. windowObservable 을 먼저 제공하므로 소스 옵저버블의 구독을 완료할 때까지 값을 다 전달받아 볼 수 있다는 장점이 있다.

12. windowCount 연산자

windowCount 연산자의 마블 다이어그램

연산자 원형
windowCount<T>(
       windowSize: number,
       startWindowEvery: number = 0
): OperatorFunction<T, Observable<T>>
  • windowSize - 설정된 숫자만큼 중첩한 옵저버블 각각의 값을 발행함다.
  • startWindowEvery - 설정한 값만큼 건너뛰며 중첩 옵저버블 각각의 값을 발행한다.
/* windowCount 연산자의 사용 예 */
const { interval } = require('rxjs');
const { take, map, windowCount, concatMap, scan, last, filter } = require('rxjs/operators');
const message = '안녕하세요. RxJS 테스트 입니다';

interval(90).pipe(
  take(message.length),
  map(x => {
    const character = message.charAt(x);
    console.log(character);
    return character;
  }),
  windowCount(5),
  concatMap(windowObservable => {
    console.log('windowObservable 전달받음');
    return windowObservable.pipe(
      filter(x => x != ' '),
      take(3),
      scan((accString, current) => accString + current, ''),
      last()
    );
  })
).subscribe(string => console.log(`결과: ${string}`));

[코드 5-34] windowCount 연산자의 사용 예

 

실행 결과

windowObservable 전달받음
안
녕
하
결과: 안녕하
세
요
windowObservable 전달받음
.
 
R
x
결과: .Rx
J
windowObservable 전달받음
S
 
테
스
결과: S테스
트
windowObservable 전달받음
 
입
니
다
결과: 입니다

5개 글자 안에서 공백을 제외한 3개의 글자가 있으면 바로 결과를 출력한다.

window 연산자나 buffer 연산자에 interval 함수를 사용하면 개발 환경에 따라서 5개 혹은 6개 글자를 나눈다.
windowCount 연산자나 bufferCount 연산자처럼 Count 가 붙은 연산자는 정확한 개수 단위로 구간을 나눈다.

12-1. windowCount 연산자에서 windowObservable의 값 발행

windowCount 연산자의 startWindowEvery 파라미터는 해당 값만큼 시프트로 이동하는 용도다. 생략하면 windowSize와 같은 숫자를 설정해서 값을 건너뛴다.

/* windowCount 연산자에 startWindowEvery 파라미터를 사용한 예 */
const { interval } = require('rxjs');
const { take, map, windowCount, mergeMap, defaultIfEmpty, scan, last, filter } = require('rxjs/operators');
const message = '간장공장공장장은강공장장이고공공장공장장은장공장장이다';
const targetWord = '공장장';

interval(10).pipe(
  take(message.length),
  map(charIndex => {
    const character = message.charAt(charIndex);
    console.log(`${character}`);
    return character;
  }),
  windowCount(targetWord.length, 1),
  mergeMap(windowObservable => {
    console.log('windowObservable 전달받음');
    return windowObservable.pipe(
      defaultIfEmpty({ empty: true }),
      scan((accString, current) => current.empty ? current : accString + current, ''),
      last()
    );
  }),
  filter(word => {
    if (typeof word === 'string') {
      console.log(`현재 단어: ${word}`);
      return word === targetWord;
    }
    return false;
  })
).subscribe(word => console.log(`${word} 발견`));

[코드 5-35] windowCount 연산자에 startWindowEvery 파라미터를 사용한 예

 

실행 결과

windowObservable 전달받음
간
windowObservable 전달받음
장
windowObservable 전달받음
공
현재 단어: 간장공
windowObservable 전달받음
장
현재 단어: 장공장
windowObservable 전달받음
공
현재 단어: 공장공
windowObservable 전달받음
장
현재 단어: 장공장
windowObservable 전달받음
장
현재 단어: 공장장
공장장 발견
windowObservable 전달받음
은
현재 단어: 장장은
windowObservable 전달받음
강
현재 단어: 장은강
windowObservable 전달받음
공
현재 단어: 은강공
windowObservable 전달받음
장
현재 단어: 강공장
windowObservable 전달받음
장
현재 단어: 공장장
공장장 발견
windowObservable 전달받음
이
현재 단어: 장장이
windowObservable 전달받음
고
현재 단어: 장이고
windowObservable 전달받음
공
현재 단어: 이고공
windowObservable 전달받음
공
현재 단어: 고공공
windowObservable 전달받음
장
현재 단어: 공공장
windowObservable 전달받음
공
현재 단어: 공장공
windowObservable 전달받음
장
현재 단어: 장공장
windowObservable 전달받음
장
현재 단어: 공장장
공장장 발견
windowObservable 전달받음
은
현재 단어: 장장은
windowObservable 전달받음
장
현재 단어: 장은장
windowObservable 전달받음
공
현재 단어: 은장공
windowObservable 전달받음
장
현재 단어: 장공장
windowObservable 전달받음
장
현재 단어: 공장장
공장장 발견
windowObservable 전달받음
이
현재 단어: 장장이
windowObservable 전달받음
다
현재 단어: 장이다
windowObservable 전달받음
현재 단어: 이다
현재 단어: 다
  • mergeMap 연산자 - concatMap 연산자를 사용하면 같은 값을 소스 옵저버블에서 발행할 때 여러 windowObservable처럼 값을 전달 받을 수 없어 mergeMap 연산자를 사용했다.
  • last 연산자 - 한 칸씩 움직이다 보면 마지막 옵저버블은 빈 옵저버블이 될 수 있다. last 연산자는 아무것도 발행하지 않고 완료하는 빈 옵저버블이 있으며 마지막 값이 없으므로 에러를 발생 시킨다.
  • defaultIfEmpty 연산자 - 소스 옵저버블이 빈 옵저버블이면 값을 발행할 수 있도록 했다.
  • filter 연산자 - 필요 없는 값들 거르게 했다.

마지막 '다'를 발행하고 한 칸 시프트 이동하면 emptyObservable 은 defaultIfEmpty 연산자의 인자에서 제공하는 값을 발행한 후 scan 연산자에서 문자열이 아닌 상태로 건너뛴다. 따라서, 로그도 출력하지 않고 filter 연산자의 조건을 만족하지 않는 형태로 값을 발행하지 않는다.

windowCount(3, 1)을 사용해 1개 값을 발행하는 windowObservable 로 값을 깍각 3개씩 발행

반응형

'RxJS' 카테고리의 다른 글

7장. 수학 및 결합 연산자 요약  (0) 2024.08.02
6장. 조합 연산자 요약  (0) 2024.08.01
4장. 필터링 연산자 요약  (0) 2024.07.29
3장. 생성 함수 요약  (0) 2024.07.25
2장. RxJS의 기본 개념 요약  (2) 2024.07.23

+ Recent posts