반응형

필터링 연산자는 옵저버블에서 값을 발행할 때 조건을 정한다.

  • 조건에 해당하지 않으면 값을 발행하지 않도록 할 수 있다.
  • 값을 발행하는 도중에도 일정 조건에 해당하지 않으면 값 발행을 멈출 수 있다.
  • 더 발행할 값이 없으면 complete 함수를 호출할 수도 있다.
  • 원하는 조건일 때 complete 함수를 호출할 수도 있다. 특정 조건을 명시했을 때 구독중인 옵저버블을 해제하는 코드 없이 연산자로만 옵저버블을 관리할 수도 있다.
  • 불필요한 연산을 실행하지 않도록 막을 수 있다.

1. filter 연산자

fiter 연산자는 소스 옵저버블에서 값을 발행할 때 참, 거짓을 리터하는 predicate 함수라는 파라미터가 있어 조건을 만족(true)할 때만 값을 발행하도록 한다.

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

연산자 원형
filter<T>(
       predicate: (value: T, index: number) => boolean,
       thisArg?: any,
): MonoTypeOperatorFuction<T>
  • predicate - 옵저버블이 생성한 값을 평가하는 함수 (true면 값을 발행하고 false면 값을 발행하지 않는다.)
  • index - 0 부터 구독 이후 발행한 i 번째 값을 나타냄
  • thisArg? - predicate 함수를 사용할 때 this 로 제공하는 값
  • MonoTypeOperatorFunction<T> - filter 가 적용된 옵저버블을 리턴하는 함수로, pipe 함수에서 이 함수를 호출함으로써 옵저버블 인스턴스를 얻음
/* filter 연산자로 짝수만 발행 */
const { range } = require('rxjs');
const { filter } = require('rxjs/operators');

range(1, 5).pipe(
  filter(x => x % 2 === 0)
).subscribe(x => console.log(`result: ${x}`));

[코드 4-1] filter 연산자로 짝수만 발행

 

실행 결과

result: 2
result: 4

predicate 함수 (x => x % 2 === 0)를 인자로 사용한다.

subscribe 함수를 호출하면 range 함수에 설정한 값을 predicate 함수에 전달하여 조건을 만족하는 값만 next 함수로 전달해 값을 발행한다.


2. first 연산자

first 연산자는 값 1개만 발행하고 complete 함수를 호출하는 연산자다.

  • 이벤트 기반의 프로그램에서 첫 이벤트만 받고 구독을 종료해야할 때
  • 옵저버블 스트림 특정 부분에서 처음 1개의 값만 필요할 때

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

연산자 원형
first<T>{
       predicate?: (value: T, index: number, source: Observable<T>) => boolean,
       defaultValue?: T
): MonoTypeOperatorFuction<T>

predicate 함수 안에 설정한 조건을 만족할 때의 값 하나만 발행하고 complete 함수를 호출한다.

만약, 조건을 만족하는 값이 없다면 소스 옵저버블을 계속 구독만 하다가 complete 나 error 함수를 호출할 때 같이 종료한다.

  • defaultValue? - 소스 옵저버블에 유효한 값이 없을 때 발행하는 기본값을 설정
/* 조건이 없는 first 연산자 사용 */
const { range } = require('rxjs');
const { first } = require('rxjs/operators');

range(1, 10).pipe(first())
  .subscribe(x => console.log(`result: ${x}`));

[코드4-2] 조건이 없는 first 연산자 사용

 

실행 결과

result: 1

조건이 없는 기본 동작은 첫번째 값만 다음 스트림에 발행한다.

/* 조건이 있는 first 연산자 사용 */
const { range } = require('rxjs');
const { first } = require('rxjs/operators');

range(1, 10).pipe(first(x => x >= 3))
  .subscribe(x => console.log(`result: ${x}`));

[코드 4-3] 조건이 있는 first 연산자 사용

 

실행 결과

result: 3

처음 조건을 만족하는 값만 전달해 발행하고 complete 함수를 호출한다.

 

first 연산자는 filter 연산자처럼 모든 발행 값을 대상으로 조건을 검사하지만, 소스 옵저버블의 모든 값을 발행하는 것이 아니라, 조건에 해당하는 첫 번째 값만 발행하고 구독을 중단한다.


3. last 연산자

last 연산자는 마지막 값 1개만 발행하는 연산자다. 옵저버블 내부에 next 함수로 전달한 값을 저장하다가 complete 함수를 호출할 때 마지막에 저장한 값을 발행한다.

  • 여러개의 값을 모아서 1개의 값을 만들 때
  • complete 함수를 호출하기 직전의 값을 알아야 할 때

주의 할 점 - complete 함수를 확실히 호출할 수 있는 옵저버블을 소스 옵저버블로 사용해야 한다. 아니면 무한 루프에 빠진다.

 

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

연산자 원형
last<T>(
       predicate?: (value: T, index: number, source: Observable<T>) => boolean,
       defaultValue?: T
): MonoTypeOperatotFunction<T>

predicate 함수의 조건에 맞는 값이 발행 될 때마다 좁저버블 객체 내부에 최신 값을 업데이트하다가 complete 함수를 호출하면 해당 시점의 최신 값을 발행한다. 만약 조건에 해당하는 값이 없다면 아무것도 발행하지 않는다.

/* 조건이 없는 last 연산자 */
const { range } = require('rxjs');
const { last } = require('rxjs/operators');

range(1, 10).pipe(last())
  .subscribe(x => console.log(`result: ${x}`));

[코드 4-4] 조건이 없는 last 연산자

 

실행 결과

result: 10

complete 함수를 호출한 시점의 마지막 값인 10을 발행한다.

/* 조건이 있는 last 연산자 사용 */
const { range } = require('rxjs');
const { last } = require('rxjs/operators');

range(1, 10).pipe(last(x => x <= 3))
  .subscribe(x => console.log(`result: ${x}`));

[코드 4-5] 조건이 있는 last 연산자 사용

 

실행 결과

result: 3

1부터 3까지는 조건을 만족하므로 옵저버블 내부에 있는 last 값이 1 -> 2 -> 3 순서로 바뀐다. 이후에는 조건을 만족하는 값이 없으므로 마지막 값인 3을 발행한다.


4. 명시적으로 구독 해제하지 않도록 돕는 연산자

  • 특정 개수만큼만 옵저버블을 구독하고 실행을 종료 할 때
  • 특정 조건을 만족하는 값을 발행한 때만 구독하다가 해당 조건을 만족하지 않을 때 구독 해제
  • 반대로 특정 조건을 만족하기 전까지 값을 발행하다가 조건을 만족하면 complete 함수를 호출해 구독 완료

complete 함수를 호출하는 시점을 제어할 수 있는 연산자에는 take 라는 접두어가 붙는다.

4-1. take 연산자

소스 옵저버블에서 정해진 개수만큼만 구독하고 구독을 해제한다.

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

연산자 원형
take<T>(count: number): MonoTypeOperatorFuction<T>

conut 는 next 함수에서 발행하는 값 중 최대값이다.

 

take 연산자는 interval 함수처럼 무한 반복 실행할 수 있는 연산자와 같이 사용하면 유용하다.

/* 1초 마다 숫자를 5회만 반복해서 발행 */
const { interval } = require('rxjs');
const { take } = require('rxjs/operators');

interval(1000).pipe(take(5))
  .subscribe(x => console.log(`result: ${x}`));

[코드 4-6] 1초 마다 숫자를 5회만 반복해서 발행

 

실행 결과

result: 0
result: 1
result: 2
result: 3
result: 4

interval 함수는 무한 반복 함수인데 take 연산자로 인해 소스 옵저버블인 interval 함수에서 발행하는 첫 5개 값만 발행하고 complete 함수를 호출하며 프로그램 실행을 종료한다. 만약 구독 해제를 하지 않았다면 interval 함수는 계속 동작했을 것이다.

4-2. takeUntil 연산자

takeUntil 연산자는 특정 이벤트가 발생할 때까지 옵저버블을 구독해야할 때 유용하다.

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

연산자 원형
takeUntil<T>(notifier: Observable<any>): MonoTypeOperatorFunction<T>

특정 이벤트가 발생했을 때 콜백 함수를 등록하여 어떤 옵저버블의 구독을 해제할 때 subscribe 함수와 unsubscribe 함수를 나눠서 작성하는 경우가 발생한다. 이 때 unsubscribe 함수를 호출할 때 발생하는 이벤트를 옵저버블로 만들어서 notifier를 사용하면 takeUntil 연산자가 있는 옵저버블을 구독할 때 notifier 를 같이 구독한다. notifier 가 처음 값을 발행하는 순간 takeUntil 연산자에 설정된 옵저버블과 notifier 의 구독을 중단한다.

<!-- 조건이 있는 takeUntil 연산자 사용 -->
<!DOCTYPE html>
<html lang="en">
<head>
  <meta charset="UTF-8">
  <title>takeUntil Example</title>
  <script src="https://unpkg.com/rxjs@^7/dist/bundles/rxjs.umd.min.js"></script>
</head>
<body>
  <div id="displayArea">
    0secs
  </div>
  <button id="stop_btn" type="button">Stop</button>
  <script>
    const { fromEvent, interval } = rxjs;
    const { take, takeUntil } = rxjs.operators;

    interval(1000).pipe(
      take(100),
      takeUntil(fromEvent(document.querySelector('#stop_btn'), 'click'))
    ).subscribe(
      x => document.querySelector('#displayArea').innerHTML = `${x+1}secs`
    );
  </script>
</body>
</html>

[코드 4-7] 조건이 있는 takeUntil 연산자 사용

 

실행 결과

stop_btn이라는 id 속성값이 있는 버튼을 클릭하면 구독을 해제하고 더 이상 시간을 업데이트하지 않는다. 100초가 지나면 take(100) 때문에 해당 옵저버블의 구독을 해제하므로 stop_btn 을 클릭해도 아무 일도 일어나지 않는다.

4-3. takeWhile 연산자

takeWhile 연산자는 특정 조건을 갖는 predicate 함수를 인자로 갖는다.

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

연산자 원형
takeWhile<T>(
       prdeicate: (value: T, index: number) => boolean
): MonoTypeOperatorFunction<T>

소스 옵저버블에서 발행하는 값을 조건을 만족하는 동안 값을 발행하고, 조건을 만족하지 않으면 더 값을 발행하지 않도록 구독을 해제한다.

/* takeWhile 연산자의 사용 예 */
const { interval } = require('rxjs');
const { filter, takeWhile } = require('rxjs/operators');

interval(300).pipe(
  filter(x => x >= 7 || x % 2 === 0),
  takeWhile(x => x <= 10)
).subscribe(x => console.log(`result: ${x}`));

[코드 4-8] takeWhile 연산자의 사용 예

 

실행 결과

result: 0
result: 2
result: 4
result: 6
result: 7
result: 8
result: 9
result: 10
  1. filter 연산자로 300ms 마다 7미만 짝수 (0, 2, 4, 6)만 발행
  2. takeWhile 연산자로 7이상 10이하 까지 1씩 순서대로 값을 발행

발행하는 값의 조건이 일정하지 않을 때 특정 조건에 맞는 값만 발행하는 방법으로 매우 유용하다.

4-4. takeLast 연산자

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

 

takeLast 연산자는 마지막에 발행한 값을 기준으로 인자로 설정한 수(0보다 크거나 같은 정수)만큼 값을 발행한다.

연산자 원형
takeLast<T>(count: number): MonoTypeOperatorFunction<T>

숫자 타입의 count 에 마지막에 발행한 값을 기준으로 최대 발행할 값 수를 설정한다.

/* takeLast 연산자 구현 일부분 */
_next(value) {
  const ring = this.ring;
  const total = this.total;
  const count = this.count;
  if (ring.length < total) {
    ring.push(value);
  } else {
    const index = count % total;
    ring[index] = value;
  }
}

[코드 4-9] takeLast 연산자 구현 일부분

 

구현 한 코드를 보면 값 발행 순서에 따라 옵저버블 내부에 배열을 만들고 값 각각의 인텍스롤 모듈러 연산으로 계산한 후 이를 저장한다라는 것을 알 수 있다.

/* takeLast 연산자로 4개 값을 발행 */
const { interval } = require('rxjs');
const { filter, takeWhile, takeLast } = require('rxjs/operators');

interval(300).pipe(
  filter(x => x >= 7 || x % 2 === 0),
  takeWhile(x => x <= 10),
  takeLast(4)
).subscribe(x => console.log(`result: ${x}`));

[코드 4-10] takeLast 연산자로 4개 값을 발행

 

실행 결과

result: 7
result: 8
result: 9
result: 10
  1. interval(300)로 인해 300ms * 7 시간만큼 아무런 출력 결과가 없다.
  2. takeWhile 연산자로 7이상 10이하까지 4개의 숫자를 300ms 마다 발행한 후 complete 함수를 호출한다.
  3. takeLast(4) 를 실행하면 소스 옵저버블에서 complete 함수를 호출한 이후이기 때문에 가장 마지막 발행한 값의 기준은 아는 상태다. 따라서 링 형태로 저장한 가장 최근 발행 값인 7, 8, 9, 10을 순서대로 발행하고 다시 complete 함수를 호출해 구독을 종료한다.
takeWhile 연산자를 이용해 발행하는 값 [0, 2, 4, 6, 7, 8, 9, 10]
takeLast(4) 로 만든 배열에 들어가는 최대값 개수 4
0 [0]
2 [0, 2]
4 [0, 2, 4]
6 [0, 2, 4, 6]
7 [7, 2, 4, 6]
8 [7, 8, 4, 6]
9 [ 7, 8, 9, 6]
10 [7, 8, 9, 10]

takeLast(4) 일 때 링에 값을 저장하는 순서

 

앞 배열의 시작점과 마지막을 연결하면 고리 형태가 되어 링 형태의 순환 구조라고 하는 것이다.

 

즉, 내부에서 가장 마지막에 저장한 인덱스의 다음 인덱스를 모듈러 연산으로 계산해서 출발점으로 정하고 n 개의 값을 순서대로 발행한다.


5. 필요 없는 값을 발행하지 않는 연산자

옵저버블을 구독하다 보면 특정 조건을 만족하는 값을 발행할 필요 없이 버려야 할 때고 있을 때 skip 으로 시작하는 연산자를 사용해서 값을 발행하지 않을 수 있다.

5-1. skip 연산자

skip 연산자는 소스 옵저버블에서 발행하는 값을 인자로 설정한 개수만큼 건너뛰고 그 다음 값부터 발행하는 연산자다.

 

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

연산자 원형
skip<T>(count: number): MonoTypeOperatorFunction<T>

숫자 타입의 count 를 사용하는 것은 takeLast 연산자와 같지만 건너뛸 값 개수를 설정한다는 점은 다르다.

/* take 연산자와 skip 연산자를 조합한 예 */
const { interval } = require('rxjs');
const { take, skip } = require('rxjs/operators');

interval(300).pipe(
  skip(3),
  take(2)
).subscribe(x => console.log(x));

[코드 4-11] take 연산자와 skip 연산자를 조합한 예

 

실행 결과

3
4

interval(300)에서 발행하는 처음 3개의 값인 0, 1, 2는 건너뛰고, 그 이후는 take(2)로 3, 4라는 2개의 값만 발행한다.

5-2. skipUntil 연산자

takeUntil 연산자와는 반대로 인자로 사용한 옵저버블의 값 발행을 시작할 때까지 소스 옵저버블에서 발행하는 값을 건너뛴다.

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

연산자 원형
skipUntil<T>(motifier: Observable<any>): MonoTypeOperatorFunction<T>

notifier는 소스 옵저버블에서 발행한 값 중 별도로 결과값을 발행할 두 번째 옵저버블을 설정한다. 즉, 해당 연산자에서 두 번째 옵저버블을 이용해 결과값을 발행한다.

 

/* skipUntil 연산자로 인자의 옵저버블 값만 발행 */
const { interval } = require('rxjs');
const { skipUntil, take } = require('rxjs/operators');

const sourceIntervalTime = 300;
interval(sourceIntervalTime).pipe(
  skipUntil(interval(sourceIntervalTime * 5)),
  take(3)
).subscribe(x => console.log(`result: ${x}`));

[코드 4-12] skipUntil 연산자로 인자의 옵저버블 값만 발행

 

실행 결과

result: 4
result: 5
result: 6

0, 1, 2, 3, 까지 네 번 값을 건너뛰고 다섯 번째 값인 4부터 값을 발행한다.

 

중요한점 - interval(sourceIntervalTime * 5) 옵저버블을 소스 옵저버블 보다 먼저 구독한다. 그렇기 때문에 아래와 같은 순서로 값을 발행한다. (아래 구현 코드의 일부를 확인하자.)

  1. interval(sourceIntervalTime * 1) - 1회
  2. 0
  3. interval(sourceIntervalTime * 2) - 2회
  4. 1
  5. interval(sourceIntervalTime * 3) -3회
  6. 2
  7. interval(sourceIntervalTime * 4) - 4회
  8. 3
  9. interval(sourceIntervalTime * 5) - 5회
  10. 4
  11. 5
  12. 6
/* skipUntil 연산자 구현 중 일부분 */
class SkipUntilOperator {
  // ...생략...
  call(subscriber, source) {
    return source.subscribe(new SkipUntilSubscriber(subscriber, this.notifier));
  }
}

class SkipUntilSubscriber extends OuterSubscriber {
  constructor(destination, notifier) {
    super(destination);
    this.hasValue = false;
    this.isInnerStopped = false;
    this.add(subscribeToResult(this, notifier));
  }
  // ...생략...
}

[코드 4-13] skipUntil 연산자 구현 중 일부분

 

생성자로 인스턴스를 생성하면서 notifier 를 구독한 후 call 메서드 안 source, subscribe 를 이용해 소스 옵저버블를 구독하므로 notifier 를 소스 옵저버블 보다 먼저 구독한는 것이다.

5-3. skipWhile 연산자

skipWhile 얀산자는 predicate 함수가 조건을 만족할 때 값을 건너 뛴다. 조건을 만족하지 않는 순간부터는 조건과 상관없이 계속 값을 발행한다.

 

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

연산자 원형
skipWhile<T>{
       predicate: (value: T, index: number) => boolean
): MonoTypeOperatorFunction<T>
/* skipWhile 연산자를 이용한 4 이상의 값 3개 발행 */
const { interval } = require('rxjs');
const { skipWhile, take } = require('rxjs/operators');

interval(300).pipe(
  skipWhile(x => x < 4),
  take(3)
).subscribe(x => console.log(`result: ${x}`));

[코드 4-14] skipWhile 연산자를 이용한 4 이상의 값 3개 발행

 

실행 결과

result: 4
result: 5
result: 6

4보다 작은 0부터 3까지는 계속 건너뛰다가 조건을 처음 만족하지 않는 4부터 계속 발행한다. take(3) 으로 3개의 값만 발행했다.


6. 값 발행 후 일정 시간을 기다리는 연산자

  • 소스 옵저버블의 값을 바로 발행하지 않고 일정 시간을 기다린다.
  • 일정 시간 동안 소스 옵저버블에서 새 값을 발행하지 않으면 조건에 따라 특정 값을 발행한다.
  • 일정 시간 동안 새 값을 발행하면 다시 일정 시간 동안 발행하는 값이 없는지 기다린다.
  • 조건을 만족할 때는 값 발행을 건너뛰다가 처음으로 조건을 만족하지 않는 순간부터 계속 값을 발행한다.

빠른 비동기 요청에서 오는 응답이나 이벤트(빠른 속도의 키보드 타이핑이나 마우스 클릭) 중 일정 시간 안에 발생한 것만 처리할 때 유용하다.

6-1. debounce 연산자

  • 선택자 함수로 소스 옵저버블에서 발행하는 값을 인자로 사용
  • 해당 선택자 함수에서 리턴하는 옵저버블이나 프로미스는 소스 옵저버블에서 발행한 다음 값을 전달받지 않으면 값을 발행
  • 선택자 함수는 소스 옵저버블에서 어떤 값을 전달받느냐에 따라 그에 상응하는 옵저버블이나 프로미스를 리턴할 수 있음

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

연산자 원형
debounce<T>(
       durationSelector: (value: T) => SubscribableOrPromise<any>
): MonoTypeOperatorFunction<T>

durationSelector 는 옵저버블 또는 프로미스가 전달받는 각 소스 옵저버블 값의 발행을 얼마나 대기시킬지 설정하는 타임아웃을 계산하는 함수다.

/* debounce 연산자로 값 발행 시간 간격을 선택 */
const { interval } = require('rxjs');
const { debounce, take, tap } = require('rxjs/operators');

const sourceInterval = 400;

interval(sourceInterval).pipe(
  take(4),
  debounce(srcVal => interval(
    srcVal % 2 === 0 ? sourceInterval * 1.2 : sourceInterval * 0.8
  ).pipe(
    tap(innerVal => console.log(
      `sourceInterval value: ${srcVal}, innerInterval value: ${innerVal}`
    ))
  ))
).subscribe(x => console.log(`result: ${x}`));

[코드 4-15] debounce 연산자로 값 발행 시간 간격을 선택

 

실행 결과

sourceInterval value: 1, innerInterval value: 0
result: 1
result: 3

 

[코드 4-15] 설명

소스 옵저버블의 값 발행 간격이 400ms이고, debounce 연산자의 인자로 선택자 함수를 사용해 값을 발행하는 예이다.

선택자 함수는 소스 옵저버블에서 발행되는 값이 짝수면 소스 옵저버블보다 긴 시간을, 홀수면 더 짧은 시간을 기다리는 옵저버블을 리턴한다.

선택자 함수에서 리턴하는 옵저버블에서 값을 발행해야만 tap 연산자를 호출하여 log를 출력한다.

take(4)로 인해 소스 옵저버블은 0, 1, 2, 3이라는 4개의 값을 발행하지만 debounce 연산자의 선택자로 인해 홀수 값인 1, 3만 발행한다.

 

이유 -> 선택자 함수에서 홀수 값을 리턴하는 읍저버블의 발행 속도를 소스 옵저버블보다 더 쁘르게 설정했기 때문이다.

 

1이란 값을 발행할 때는 로그를 출력했지만, 3을 발행할 때는 로그를 출력하지 않는다.

 

이유 -> 소스 옵저버블이 3을 발행한 후 take(4) 때문에 complete 함수를 호출하기 때문이다. debounce 연산자는 complete 함수가 호출되면 더 기다릴 필요 없이 값을 발행한다.

소스 옵저버블에서 complete 함수를 호출하면 선택자 함수에서 리턴해 구독 중인 가장 최근 옵저버블이나 프로미스를 구독 해제한다. 그리고 complete 함수 호출 직전 소스 옵저버블 값을 발행한다.

 

6-2. debounceTime 연산자

debounceTime 연산자는 값 발행을 기다리는 일정 시간을 인자로 설정한 후 해당 시간 안에 소스 옵저버블에서 발행한 다음 값을 전달받지 않으면 최근 값을 그대로 발행하는 연산자다.

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

연산자 원형
debounceTime<T>(
       dueTime: number,
       scheduler: SchedulerLike = async
): MonoTypeOperatorFunction<T>
  • dueTime - 최근 값을 발행하기 전 값 발행을 얼마나 대기시킬지 설정하는 타임아웃을 계산하는 함수
  • scheduler - 각 값의 발행 대기 시간을 관리하는 스케줄러
/* debounceTime 연산자로 발행 간격에 따라 값 발행 */
const { interval } = require('rxjs');
const { debounceTime, take } = require('rxjs/operators');

interval(400).pipe(take(4), debounceTime(300))
  .subscribe(x => console.log(
    `- interval(400).pipe((take(4), debounceTime(300)) next: ${x}`
  )
);

interval(400).pipe(take(4), debounceTime(500))
  .subscribe(
    x => console.log(
      `-- interval(400).pipe((take(4), debounceTime(500)) next: ${x}`
    )
);

[코드 4-16] debounceTime 연산자로 발행 간격에 따라 값 발행

 

실행 결과

- interval(400).pipe((take(4), debounceTime(300)) next: 0
- interval(400).pipe((take(4), debounceTime(300)) next: 1
- interval(400).pipe((take(4), debounceTime(300)) next: 2
- interval(400).pipe((take(4), debounceTime(300)) next: 3
-- interval(400).pipe((take(4), debounceTime(500)) next: 3

소스 옵저버블보다 값 발행 간격이 짧으면 소스 옵저버블의 모든 값을 발행한다.

소스 옵저버블보다 값 발행 간격이 길면 마지막 1개의 값만 발행한다.

 

이유 -> complete 함수 호출 직전의 마지막 값은 이후에 다음 값이 없으므로 complete 함수 호출과 동시에 무조건 발행된다.


7. 중복 값을 발행하지 않는 연산자

  • 스트림에서 발행하는 값을 전달 받을 때 전체 스트림 중 중복 값을 더 전달 받고 싶지 않을 때
  • 연속된 같은 값들은 처음에 값 하나만 발행해서 중복 발행을 피하고 싶을 때

7-1. distinct 연산자

distinct 연산자는 내부에서 자체 구현한 Set 자료 구조로 이미 발행된 값을 중복 없이 저장했다가 같은 값을 전달받으면 발행하지 않는다.

해당 자료구조는 배열 형태이며 해당 배열의 indexOf 함수를 사용하여 얻은 index 값이 -1과 다른지 검사하여 값의 유무를 확인한다.

(엄격한 동등성 비교)

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

연산자 원형
distinct<T, K>(
       keySelector?: (value: T) => K,
       flushes?: Observable<any>
): MonoTypeOperatorFunction<T>
/* distinct 연산자로 중복 값을 발행하지 않는 예 */
const { of } = require('rxjs');
const { distinct } = require('rxjs/operators');

of(1, 6, 7, 7, 2, 5, 5, 2, 6).pipe(distinct())
  .subscribe(x => console.log(x));

[코드 4-17] distinct 연산자로 중복 값을 발행하지 않는 예

 

실행 결과

1
6
7
2
5

중복 값 7, 2, 5를 한 번만 발행한다.

KeySelector 함수를 이용한 객체 타입 값의 중복 확인

/* 객체 타입 값의 중복 확인 */
const { of } = require('rxjs');
const { distinct, map } = require('rxjs/operators');

of(
  { id: 1, value: 20 },
  { id: 2, value: 40 },
  { id: 3, value: 70 },
  { id: 1, value: 20 },
  { id: 2, value: 40 },
  { id: 3, value: 70 }
).pipe(distinct(), map(x => x.value)).subscribe(x => console.log(x));

[코드 4-18] 객체 타입 값의 중복 확인

 

실행 결과

20
40
70
20
40
70

중복 값을 발행했다.

 

이유 => 객체 타입은 동등성을 검사할 때 객체의 참조값을 비교하기 때문이다. 각각 새로 생성한 객체이기 때문에 다른 값으로 판단한다.

 

/* 키 값을 찾는 함수를 인자로 사용 */
const { of } = require('rxjs');
const { distinct, map } = require('rxjs/operators');

of(
  { id: 1, value: 20 },
  { id: 2, value: 40 },
  { id: 3, value: 70 },
  { id: 1, value: 20 },
  { id: 2, value: 40 },
  { id: 3, value: 70 }
).pipe(distinct(obj => obj.id), map(x => x.value)).subscribe(x => console.log(x));

[코드 4-19] 키 값을 찾는 함수를 인자로 사용

 

실행 결과

20
40
70

첫 번째 인자로 keySelector 함수(obj => obj.id)를 사용하였다. 비교하는 기준인 key값(여기서는 객체의 id)으로 동등성 비교하여 중복 값을 발행하지 않는다.

중복 값 검사를 초기화하는 flush

distinct 연산자는 두번째 인자로 옵저버블을 사용하는 flush를 제공한다. distinct 연산자를 사용하는 옵저버븝을 구독할 때 같이 구독하며 값을 발행할 때마다 소스 옵저버블에서 중복 값을 검사하는 Set 자료 구조를 초기화 한다.

/* flush 옵저버블 사용 예 */
const { interval } = require('rxjs');
const { take, map, distinct } = require('rxjs/operators');

interval(200).pipe(
  take(25),
  map(x => ({ original: x, value: x % 5 })),
  distinct(x => x.value, interval(2100))
).subscribe(x => console.log(JSON.stringify(x)));

[코드 4-20] flush 옵저버블 사용 예

 

실행 결과

{"original":0,"value":0}
{"original":1,"value":1}
{"original":2,"value":2}
{"original":3,"value":3}
{"original":4,"value":4}
{"original":10,"value":0}
{"original":11,"value":1}
{"original":12,"value":2}
{"original":13,"value":3}
{"original":14,"value":4}
{"original":20,"value":0}
{"original":21,"value":1}
{"original":22,"value":2}
{"original":23,"value":3}
{"original":24,"value":4}

 

[코드 4-20 실행 설명]

  1. 200ms 마다 값을 발행하는 소스 옵저버블은 객체로 바꿔 줄 때 original에 원래 값을 담았다.
  2. 나머지 연산자르 사용하여 value에 중복값 검사의 기준이 될 0, 1, 2, 3, 4를 반복해서 발행하도록 했다.
  3. flush 옵저버블은 2초를 조금 넘는 시간(2100ms) 마다 값을 발행한다.
  4. 첫 1초는 값을 발행하고, 두 번째 1초는 deistinct 연산자에서 중복 값을 검사하므로 값을 발행하지 않는다.

7-2. distinctUntilChaged 연산자

distinctUntilChanged 연산자는 같은 값이 연속으로 있는 지를 검사하는 연산자다. 연속해서 중복 값이 있다면 최초 값 1개만 발행하며 그 이외는 정상적으로 발행한다. 단, 연속해서 중복 값이 있는 것이 아니라면 중복 값 발행을 허용한다. (동등성 비교)

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

연산자 원형
distinctUntilChanged<T, K>(
       compare?: (x: K, y: K) => boolean,
       keySelector?: (x: T) => K
): MonoTypeOperatorFunction<T>
/* distinctUntilChanged 연산자의 사용 예 */
const { of } = require('rxjs');
const { distinctUntilChanged } = require('rxjs/operators');

of(1, 6, 7, 7, 2, 5, 5, 2, 6).pipe(distinctUntilChanged())
  .subscribe(x => console.log(x));

[코드 4-21] distinctUntilChanged 연산자의 사용 예

 

실행 결과

1
6
7
2
5
2
6

compare 함수로 값 비교

compar 함수 - 같은 값인지 어떻게 비교할 지 정하는 비교 함수를 인자로 사용

/* compare 함수로 값 비교하기 */
const { of } = require('rxjs');
const { distinctUntilChanged } = require('rxjs/operators');

of(
  { a: 1, b: 20 },
  { a: 1, b: 20 },
  { a: 2, b: 40 },
  { a: 3, b: 70 },
  { a: 3, b: 70 },
  { a: 2, b: 40 }
).pipe(distinctUntilChanged((o1, o2) => o1.a === o2.a && o1.b === o2.b))
  .subscribe(x => console.log(JSON.stringify(x)));

[코드 4-22] compare 함수로 값 비교하기

 

실행 결과

{"a":1,"b":20}
{"a":2,"b":40}
{"a":3,"b":70}
{"a":2,"b":40}

a와 b의 값이 같고, 연속해서 위치했다면 값을 발행하지 않는다.

keySelector 함수로 값 비교

두 번째 인자로 ketSelector를 사용할 수 있다. 키 값으로 발행할 값을 비교한다. 보통 compare 함수 없이 keySelector 함수를 사용하면 동등성 비교한다.

/* compare 함수 없이 keySelector 함수 사용 */
const { of } = require('rxjs');
const { distinctUntilChanged } = require('rxjs/operators');

of(
  { a: 1, b: 20 },
  { a: 1, b: 20 },
  { a: 2, b: 40 },
  { a: 3, b: 70 },
  { a: 3, b: 70 },
  { a: 2, b: 40 }
).pipe(distinctUntilChanged(null, x => x.a))
  .subscribe(x => console.log(JSON.stringify(x)));

[코드 4-23] compare 함수 없이 keySelector 함수 사용

 

실행 결과

{"a":1,"b":20}
{"a":2,"b":40}
{"a":3,"b":70}
{"a":2,"b":40}
/* keySelector 와 compare 함수를 함계 사용한 예 */
const { of } = require('rxjs');
const { distinctUntilChanged } = require('rxjs/operators');

of(
  { objkey: { a: 1, b: 20 } },
  { objkey: { a: 1, b: 20 } },
  { objkey: { a: 2, b: 40 } },
  { objkey: { a: 3, b: 70 } },
  { objkey: { a: 3, b: 70 } },
  { objkey: { a: 2, b: 40 } }
).pipe(distinctUntilChanged(
  (o1, o2) => o1.a === o2.a && o1.b === o2.b, // compare 함수
  x => x.objkey // keySelector 함수
)).subscribe(x => console.log(JSON.stringify(x)));

[코드 4-24] keySelector 와 compare 함수를 함께 사용한 예

 

실행 결과

{"objkey":{"a":1,"b":20}}
{"objkey":{"a":2,"b":40}}
{"objkey":{"a":3,"b":70}}
{"objkey":{"a":2,"b":40}}
  •  keySelector -> objkey 키 값의 객체를 불러온다.
  • compare -> 동등성을 검사한다.

a와 b의 값이 같으면서 연속해서 등장하는 객체는 발행하지 않도록 한 것이다.


8. 샘플링 연산자

  • 스트림에서 발행하는 값 중 모든 값이 필요하지 않고 일부 샘플만 있어도 충분할 때

8-1. sample 연산자

sample 연산자는 notifier라는 옵저버블을 인자로 사용해 notifier 옵저버블에서 값을 발행할 때마다 소스 옵저버블의 가장 최근 값을 발행한다.

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

sample 연산자로 값을 발행한 후 소스 옵저버블에서 다음 값을 발행하기 전 notifier 에서 또 값을 발행해도 최근 값을 중복해서 발행하지 않는다.

연산자 원형
sample<T>(notifier: Observable<any>): MonoTypeOperatorFunction<T>

notifier 에서 complete 함수를 호출해도 소스 옵저버블의 가장 최근 값은 발행한다. 단, 이후 소스 옵저버븛에서 값을 발행해도 해당 값을 발행하지 않는다.

/* sample 연산자의 사용 예 */
const { interval, timer } = require('rxjs');
const { sample, take } = require('rxjs/operators');

// source: 0(200ms),......,3(800ms),......,6(1400ms),......,9(2000ms)
// sample:        0(300ms),       1(900ms),       2(1500ms),       3(2100ms)
const sampleSize = 3;
const sourceInterval = 200;
const sampleDelay = 100;

interval(sourceInterval)  // 200ms
  .pipe(sample(timer(
    sourceInterval + sampleDelay, // 300ms
    sourceInterval * sampleSize // 600ms
    )), take(4))
  .subscribe(result => console.log(result));

[코드 4-25] sample 연산자의 사용 예

 

실행 결과

0
3
6
9
  • source 옵저버블은 200ms 마다 값을 발행한다.
  • sample 연산자의 notifier 옵저버블은  timer 함수를 사용해 100ms 정도의 차이을 두고 600ms 간격 (200ms * 3)마다 값을 발행하도록 한다.

소스 옵저버블에서 첫 값인 0을 발행한 지 300ms 후에 sample 연산자가 0 값을 발행하며, 이후 600ms 마다 3, 6, 9 총 4개 값을 발횅한다.

8-2. sampleTime 연산자

sampleTime 연산자는 ms 단위의 발행 간격을 인자로 설정한 후 해당 발행 간격 사이에 있는 소스 옵저버블의 최근 값을 발행하는 연산자다. 연속해서 발행하는 이벤트 중 일정 간격으로 가장 최근 값 하나만 뽑아 처리할 때 유용하다.

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

연산자 원형
sampleTime<T>{
       period: number,
       scheduler: SchedulerLike = async
): MonoTypeOperatorFuntion<T>
  • period -> 일정 간격을 설정하는 것
  • scheduler -> sampleTime 연산자의 실행 시점을 관리하는 스케줄러
/* sampleTime 연산자의 사용 예 */
const { timer } = require('rxjs');
const { sampleTime, take } = require('rxjs/operators');

// source: 0(300ms)  1(700ms)  2(1100ms)  3(1500ms)  4(1900ms)  5(2300ms)
// sample:                 800ms                1600ms                 2400ms
const sourcePoint = 300;
const sourceDelay = 400;
const sampleCount = 2;
const samplePeriod = sourceDelay * sampleCount; // 800ms

timer(sourcePoint, sourceDelay) // 300ms, 400ms
  .pipe(
    sampleTime(samplePeriod),  // 800ms
    take(3)
  )
  .subscribe(result => console.log(result));

[코드 4-26] sampleTime 연산자의 사용 예

 

실행 결과

1
3
5

timer 함수로 소스 옵저버블을 만든 후 처음 간격과 이후 간격을 다르게 실행한다.

반응형

'RxJS' 카테고리의 다른 글

6장. 조합 연산자 요약  (0) 2024.08.01
5장. 변환 연산자 요약  (0) 2024.07.31
3장. 생성 함수 요약  (0) 2024.07.25
2장. RxJS의 기본 개념 요약  (2) 2024.07.23
1장. RxJS 소개와 개발 환경 구축 요약  (1) 2024.07.10

+ Recent posts