반응형

1. 프로그래밍 언어란 무엇인가?

1.1.1 프로그래밍 언어 정의

언어란?

언어는 의사 전달 도구, 다른 사람에게 자신의 감정이나 뜻을 전하고 싶을 때 우리는 언어를 사용하여 말을 하고 글을 쓴다.

 

프로그래밍 언어란?

프로그램 작성에 사용되는 언어이다.

 

프로그램이란?

컴퓨터가 수행할 명령어를 순서대로 나열해 둔 것을 뜻한다.

 

자연어란?

프로그램 언어와 대비하여 사용하는 우리말이나 영어 등을 일컫는다.

 

계산이란?

주어진 입력으로부터 원하는 답을 찾기 위해 수행해야 하는 명확한 절차를 뜻한다.

프로그래밍 언어는 컴퓨터가 수행할 수 있고 사람이 읽을 수 있는 형태로 계싼을 나타내는 표기 체계이다.

1.1.2 프로그래밍 언어의 특징

  1. 프로그래밍 언어는 말이 아닌 글 형태로 사용된다.
  2. 프로그래밍 언어는 엄밀한 규칙을 따르고 있다.
  3. 프로그래밍 언어는 기계에 명령을 전달하는 단방향성을 띠고 있다.

특징 1의 예외로 시각적 언어인 MIT의 스크래치나, 내셔널 인스트루먼트의 랩뷰등을 들수 있다.

특징 3의 예외로 프로그램을 읽기 쉬운 특성 가독성을 고려하면 프로그래밍 언어는 양방향의 특성도 있다.

 

1.1.3 프로그래밍 언어를 배워야 하는 이유

프로그래밍이 체계적으로 생각하는 방법을 가르쳐 준다.


1.2 프로그래밍 언어의 기능

1.2.1 프로그래밍 언어의 기본 기능

  1. 작성력 : 프로그래머의 의도를 나타낼 수 있게 하는 기능
  2. 가독성 : 프로그램을 쉽게 해독할 수 있게 하는 기능
  3. 실행 가능성 : 컴퓨터에서 실행될 수 있게 하는 기능
char *loop(char *b, char *a) {
    char *s = b;
    while ((*b++ = *a++));
    return s;
}

[프로그램 loop.c]

char *strcpy(char dst[], const char src[]) {
    int i;
    for (i = 0; src[i] != '\0'; i++) {
        dst[i] = src[i];
    }
    dst[i] = '\0';
    return dst;
}

[프로그램 strcpy.c]

 

프로그램 strcpy.c는 loop.c의 가독성을 높인 버전이다.

1.2.2 프로그래밍 언어의 부가 기능

  1. 추상화 - 어떤 대상을 간략하게 추려 나타내는 방법
  2. 모듈화 - 복잡한 대상을 나누어 구성할 수 있는 방법

1.2.3 프로그래밍 언어의 특성

  1. 기계적 - 프로그램은 기계적으로 처리할 수 있어야 한다.
  2. 구조적 - 프로그램은 복잡한 구조를 나타낼 수 있어야 한다.
  3. 가변적 - 프로그래밍 언어는 시대의 필요에 따라 바뀔 수 있다.

1.2.4 프로그래밍 언어의 스펙트럼

추상화 수준은 고급언어와 저급언어를 구분짓는 잣대이다.


1.3 프로그래밍 언어의 구성 요소

1.3.1 데이터

데이터 - 어떤 자료를 프로그램이 처리할 수 있는 형태로 나타낸 것이다.

  • 이진 데이터 - 이진수의 나열로 이루어진 데이터
  • 텍스트 데이터 - 문자열을 나타내는 데이터
코드 - 문자와 이진 표현 사이의 대응관계

 

ASCII - 특수기호, 구두점, 숫자, 영어 대소문자를 나타내기 위한 코드 (영문만 나타낼 수 있는 한계점)

유니코드 - ASCII 코드의 한계를 극복하기 위한 코드 (문자집합과 인코딩의 조합으로 코드를 나타낸다.)

 

음성 및 영상 인코딩한 데이터를 해석하는 프로그램을 코덱이라고 한다.

1.3.2 연산

  • 연산 - 데이터를 처리하는 방법
  • 연산자 - 특별한 연산을 수행하는 함수
  • 변수 - 연산 결과를 저장하는 이름
  • 이름 & 식별자 - 변수 중 값이 바뀔 수 없는 상수
  • 대입 연산 - 값을 지정하기 위해 값을 저장하는 연산
det = b * b - 4 * a * c;

대입 연산자의 우측 값을 구하여 좌측 변수에 대입한다는 뜻이다. (=의 수학적 의미와는 다르다.)

  1. 수식 - 값을 나타내는 표현
  2. 문장 - 처리를 나타내는 표현

1.3.3 명령어

  • 명령어 - if 나 while 처럼 특정 작업을 지시하는 단어를 의미한다.
  • 사용자정의연산 - 프로그래머가 추가로 연산을 정의하여 사용 할 수 있다.
  • 원시연산 - 프로그래밍 언어가 기본적으로 제공하고 있는 연산
  • 라이브러리 - 원시연산에 포함되지 않지만, 사용자가 자주 사용할 만한 연산을 미리 정의 해 둔 것
  • 표준 라이브러리 - 프로그래밍 언어에서 기본적으로 제공하는 라이브러리

1.3.4 서브프로그램

서브프로그램 - 전체 프로그램을 이루는 작은 코드 블록에 이름을 붙인 것(프로그램 구조화 장점)

  • 함수 - 연산 수행 결과 값을 반환하는 서브프로그램
  • 프로시저 - 결과 값을 반환하지 않는 서브프로그램

1.3.5 타입

타입 - 데이터 집합과 연산 집합을 합친 개념

  1. 강타입 언어 - 타입 오류를 모두 검출하는 언어
  2. 약타입 언어 - 타입 오류를 검출하지만 일부 타입 오류를 허용하는 언어
  3. 무타입 언어 - 타입 선언문도 없고 어떤 대상의 타입이 계속 변경 될 수 있는 언어

1.3.6 모듈

모듈 - 독립적인 프로그램 구성 단위

모듈은 컴파일 단위에 해당하거나 컴파일 단위의 상위 단위에 해당한다.

스텁
인터페이스에 구현하는 서브프로그램이 완성되지 않은 경우에 간단한 임시 코드(dummy code)로 구현을 대신하는 경우에 이 임시코드를 스텁이라고 한다.

1.4 프로그래밍 언어의 학습 방법

1.4.1 프로그래밍 언어의 두 가지 측면

구문론 - 프로그래밍 언어가 일차적으로 규정하고 있는 프로그램 형태에 관한 이론

의미론 - 제대로 작성된 프로그램에 대한 수행 의미를 정의 하고 있고 이에 관한 이론

프로세스 - 수행 중인 프로그램

1.4.2 어떻게 프로그래밍 언어를 배워야 하나

프로그래밍 언어를 배우기 위해서는 프로그램을 많이 읽고, 많이 작성해 보고, 많이 생각하는 것이 중요하다.

1.4.3 프로그래밍 언어의 선택 방법

  • 자신이 조금이라도 아는 프로그래밍 언어를 선택해야 한다.
  • 사용해 볼 수 있는 프로그래밍 언어를 선택해야 한다.
  • 주위에서 정보를 얻을 수 있는 프로그래밍 언어를 선택해야 한다.
  • 프로그램을 관리하기 쉬운 언어여야 한다.

1.4.4 프로그래밍 언어의 학습 요령

  • 눈보다 손
  • 그림으로 생각
  • 점진적 변경
프로토타입 작성 시 주의 할 점
처음 작성한 프로토타입은 버릴 것을 전제로 작성해야 한다. 문제를 이해하기 위해 프로그램을 활용하는 것일 뿐이지 이를 프로그램의 모태로 삼는 것을 전제하는 것은 아니다.

1.4.5 왜 프로그래밍 언어론을 배워야 하나?

  • 새로운 프로그래밍 언어를 쉽게 습득하기 위해서
  • 내가 사용하는 언어를 더 잘 이해하기 위해서
  • 새로운 프로그래밍 언어를 설계하기 위해서
#include <stdio.h>

int main() {
    printf("%d students.\n",
           printf("Good ") + printf("morning "));
    return 0;
}

[예제 코드 - 1]

 

실행 결과

Good morning 12 students.
#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;
}

[예제코드 - 2]

 

실행 결과

Good morning 13 students.
반응형
반응형

스케줄러란?

옵저버가 옵저버블을 구독할 때 값을 전달받는 순서와 실행 컨텍스트를 관리하는 역할을 하는 자료구조다.

 


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 의 상속 예 */
const { asyncScheduler } = require('rxjs');

asyncScheduler.schedule(function work(value) {
  value = value || 0;
  console.log('value: ' + value);
  const selfAction = this;
  selfAction.schedule(value + 1, 1000);
}, 1000);

[코드 13-3] AsyncScheduler 의 상속 예

 

실행 결과

value: 0
value: 1
... 이하 생략

1초마다 1씩 증가라는 값을 발행한다.

 

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) 을 전달해 동작을 실행한다.
/* schedule 함수 호출 */
const {Immediate} = require("rxjs/internal-compatibility");

requestAsyncId(scheduler, id, delay = 0) {
  if (delay !== null && delay > 0) {
    return super.requestAsyncId(scheduler, id, delay);
  }
  scheduler.actions.push(this);
  return scheduler.scheduled
    || (scheduler.scheduled = Immediate.setImmediate(
      scheduler.flush.bind(scheduler, null)
    ));
}

[코드 13-5] schedule 함수 호출

 

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 사용 - 연산자 안에서 동기 방식 및 콜스택이 아닌 반복문으로 꼬리 재귀를 호출해야 하거나 큐에 넣어 순서를 맞춰야 할 때 사용하면 좋다.

/* QueueScheduler 를 사용하는 피보나치 수열 */
const { queueScheduler } = require('rxjs');

const n = 6;

queueScheduler.schedule(function (state) {
  console.log(`fibonacci[${state.index}]: ${state.a}`);
  if (state.index < n) {
    this.schedule({
      index: state.index + 1,
      a: state.b,
      b: state.a + state.b
    });
  }
}, null, { index: 0, a: 0, b: 1 });

[코드 13-14] QueueScheduler 를 사용하는 피보나치 수열

 

실행 결과

fibonacci[0]: 0
fibonacci[1]: 1
fibonacci[2]: 1
fibonacci[3]: 2
fibonacci[4]: 3
fibonacci[5]: 5
fibonacci[6]: 8

두번째 인자 null 은 delay 값을 지정하지 않겠다는 의미이다. 이유는 QueueScheduler 도 delay 값을 지정하면 AsyncScheduler 를 상속받아 동작하므로 null 로 지정했다. 세번째 인자는 초기값을 넣었다.index 값을 1씩 증가시켜 재귀로 a, b 의 값을 누적 시킨다.

 

큐에서 하나식 꺼내서 순차적으로 동작함을 확인할 수 있다.


4. 스케줄러에서 사용하는 연산자

4-1. subscribeOn 연산자

subscribeOn 연산자는 구독하는 옵저버블 자체를 인자로 사용할 스케줄러로 바꿔준다.

/* subscribeOn 연산자의 사용 예 */
const { Observable, asyncScheduler } = require('rxjs');
const { subscribeOn } = require('rxjs/operators');

const source$ = Observable.create(observer => {
  console.log("BEGIN source");
  observer.next(1);
  observer.next(2);
  observer.next(3);
  observer.complete();
  console.log("END source");
});

console.log("before subscribe");
source$.pipe(subscribeOn(asyncScheduler, 1000)).subscribe(x => console.log(x));
console.log("after subscribe");

[코드 13-15] subscribeOn 연산자의 사용 예

 

실행 결과

before subscribe
after subscribe
BEGIN source
1
2
3
END source

'after subscribe' 를 출력할 때까지 동기 방식으로 실행되며, 1초 후 그 다음에는 비동기로 나머지 결과를 출력한다.

구독할 때 맨 앞 옵저버블의 작업을 subscribeOn 연산자에서 지정한 스케줄러로 실행하는 것이다.

연산자 원형
subscribeOn<T>(
       scheduler: SchedulerLike, delay: number = 0
): MonoTypeOperatorFunction<T>
  • scheduler - 구독 작업을 실행하는 스케줄러를 설정
  • delay - 숫자 타입으로 지연 시간을 설정

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

/* subscribeOn 연산자의 구현 코드 일부 */
const {SubscribeOnObservable} = require("rxjs/internal-compatibility");
call(subscriber, source) {
  return new SubscribeOnObservable(source, this.delay, this.scheduler)
    .subscribe(subscriber);
}

[코드 13-16] subscribeOn 연산자의 구현 코드 일부

/* SubscribeOnObservable 의 구현 코드 일부 */
static dispatch(arg) {
  const {source, subscriber} = arg;
  return this.add(source.subscribe(subscriber));
}

_subscribe(subscriber) {
  const delay = this.delayTime;
  const source = this.source;
  const scheduler = this.scheduler;
  return scheduler.schedule(SubscribeOnObservable.dispatch, delay, {source, subscriber});
}

[코드 13-17] SubscribeOnObservable 의 구현 코드 일부

 

스케줄러 안에서 소스 옵저버블 구독 부분을 실행한다는 것을 알수 있다. 소스 옵저버블뿐만 아니라 그 아래 영향을 받는 다은 연산자에서 파생한 옵저버블도 해당 스케줄러의 영향을 받는다.

4-2. observeOn 연산자

subscribeOn 연산자를 사용하면 모든 스트림이 해당 스케줄러로 바뀐다. 이 때 observeOn 연산자를 사용하면 observeOn 이 호출된 이후부터 스케줄러를 바꿔 실행할 수 있다.

 

특징 - 여러 연산자가 연속으로 연결되었을 때 가장 마지막에 호출된 연산자의 스케줄러가 우선 적용된다.

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

연산자에 있는 세모 모양의 스케줄러를 이용해서 구독할 때 next, error, complete 함수를 destination 이라는 생성자의 파라미터로 전달한다.

옵저버블의 연산자 체인

subscribeOn 연산자가 위치한 부분 가장 위 옵저버블이 시작하는 부분부터 사용할 스케줄러를 지정한다.

observeOn 연산자는 이 연산자가 호출한 지점부터 아래에 있는 모든 연산자가 지정한 스케줄러로 동작하도록 영향을 준다.

그렇기 때문에 중간에 observeOn 을 여러번 호출해 연산자별로 다른 스케줄러를 사용할 수 있는 것이다.

 

간단히 정리하자면 subscribeOn 는 밑에서 위(bottom up)로 observeOn 은 위에서 아래(top down)방향으로 영향을 준다고 이해하면 된다.

 

각 연산자가 연속으로 여러 개 연결되었을 때 subscribeOn 은 제일 먼저 연결된 연산자의 스케줄러가 적용되고, observeOn 은 가장 나중에 연결된 연산자의 스케줄러를 사용한다고 이해하면 된다. 해당 방향으로 덮어쓴다는 개념이다.

 연산자 원형
observeOn<T>(
       scheduler: SchedulerLike, delay: number = 0
): MonoTypeOperatorFuntion<T>
  • scheduler - 구독 작업을 실행하는 스케줄러를 설정
  • delay - 숫자 타입으로 지연 시간을 설정
/* observeOn 연산자의 사용 예 */
const { Observable, asyncScheduler } = require('rxjs');
const { observeOn } = require('rxjs/operators');

const source$ = Observable.create(observer => {
  console.log("BEGIN source");
  observer.next(1);
  observer.next(2);
  observer.next(3);
  observer.complete();
  console.log("END source");
});

console.log("before subscribe");
source$.pipe(observeOn(asyncScheduler, 1000)).subscribe(x => console.log(x));
console.log("after subscribe");

[코드 13-18] observeOn 연산자의 사용 예

 

실행 결과

before subscribe
BEGIN source
END source
after subscribe
1
2
3

'after subscribe' 출력까지는 동기 방식으로 실행되고 observeOn 연산자 다음으로 생성되는 옵저버블은 1초 후 스케줄러를 이용해서 실행된다. 즉, subscribe 함수 안에 있는 next 함수의 동작이 스케줄러의 영향을 받아 1부터 3까지 출력만 1초 후 비동기로 처리한다.

 

observeOn 연산자 다음에 바로 subscribe 함수를 호출하지 않고 다른 연산자를 추가했어도 그 다음에 추가하는 연산자부터는 스케줄러를 이용해 옵저버블을 실행한다.

/* observeOn 연산자의 구현 코드 일부 */
import {Subscriber} from "rxjs";

export class ObserveOnSubscriber extends Subscriber {
  constructor(destination, scheduler, delay = 0) {
    super(destination);
    this.scheduler = scheduler;
    this.delay = delay;
  }
  
  static dispatch(arg) {
    const {notification, destination} = arg;
    notification.observe(destination);
    this.unsubscribe();
  }
  
  scheduleMessage(notification) {
    this.add(this.scheduler.schedule(
      ObserveOnSubscriber.dispatch,
      this.delay,
      new ObserveOnMessage(notification, this.destination)
    ));
  }
  
  _next(value) {
    this.scheduleMessage(Notification.createNext(value));
  }
  
  _error(err) {
    this.scheduleMessage(Notification.createError(err));
  }
  
  _complete() {
    this.scheduleMessage(Notification.createComplete());
  }
}

export class ObserveOnMessage {
  constructor(notification, destination) {
    this.notification = notification;
    this.destination = destination;
  }
  
}

[코드 13-19] observeOn 연산자의 구현 코드 일부

 

dispatch 가 스케줄러의 work 함수로 동작한다. 이 때 Notification 은 객체에서 next, error, complete 함수 중 무엇을 실행할지 정해서 observeOn 함수에 있는 destination 에 전달한다.

4-3. observeOn 연산자 안 AsyncScheduler 사용

  • subscribeOn 연산자 - 소스 옵저버블을 구독하는 동작 하나만 스케줄러에서 실행
  • observeOn 연산자 - 스케줄러에서 next 함수로 여러 값을 전달하는 동작을 담당
/* observeOn 연산자 안 AsyncScheduler 사용 */
const { of, asyncScheduler } = require('rxjs');
const { observeOn } = require('rxjs/operators');

console.log('start');
of(1, 2, 3).pipe(observeOn(asyncScheduler, 1000)).subscribe(x => console.log(x));
console.log(`actions length : ${asyncScheduler.actions.length}`);
console.log('end');

[코드 13-20] observeOn 연산자 안 AsyncScheduler 사용

 

AsyncScheduler 는 각각의 action 을 setInterval 함수로 지정된 시간 뒤에 실행되도록 구현되어 있다.

action - next 함수 3개 + complete 함수 = 4개

'end'까지는 동기 방식으로 실행되고 다음부터는 1초후 스케줄러를 비동기 방식으로 살행된다.

4-4. observeOn 연산자 안 AsapScheduler 사용

  • AsyncScheduler - next 함수를 호출할 때마다 setInterval 함수도 매번 호출한다.
  • AsapScheduler- 스케줄러 안에 있는 actions 배열에 해당 액션을 푸시하는 동작만 한다.
/* observeOn 연산자 안 AsapScheduler 사용 */
const { of, asapScheduler } = require('rxjs');
const { observeOn } = require('rxjs/operators');

console.log('start');
of(1, 2, 3).pipe(observeOn(asapScheduler)).subscribe(x => console.log(x));
console.log(`actions length : ${asapScheduler.actions.length}`);
console.log('end');

[코드 13-21] observeOn 연산자 안 AsapScheduler 사용

 

실행 결과

start
actions length : 4
end
1
2
3

1~3의 출력 부분은 비동기로 실행된다. actions 배열의 4개 액션은 동기 방식으로 푸시한다.

4-5. observeOn 연산자 안 QueueScheduler 사용

/* observeOn 연산자 안 QueueScheduler 사용 */
const { of, queueScheduler } = require('rxjs');
const { observeOn } = require('rxjs/operators');

console.log('start');
of(1, 2, 3).pipe(observeOn(queueScheduler)).subscribe(x => console.log(x));
console.log(`actions length : ${queueScheduler.actions.length}`);
console.log('end');

[코드 13-22] observeOn 연산자 안 QueueScheduler 사용

 

실행 결과

start
1
2
3
actions length : 0
end
/* range 함수 안에 QueueScheduler 사용 */
const { range, queueScheduler } = require('rxjs');
const { mergeMap, observeOn } = require('rxjs/operators');

console.log('start queue');
range(0, 3, queueScheduler).pipe(mergeMap(x => range(x, 3, queueScheduler)))
  .subscribe(x => console.log(x));
console.log('end queue');

console.log('start without queue');
range(0, 3).pipe(mergeMap(x => range(x, 3)))
  .subscribe(x => console.log(x));
console.log('end without queue');

[코드 13-23] range 함수 안에 QueueScheduler 사용

 

실행 결과

start queue
0
1
1
2
2
2
3
3
4
end queue
start without queue
0
1
2
1
2
3
2
3
4
end without queue
반응형

'RxJS' 카테고리의 다른 글

12장. 멀티캐스팅 연산자 요약  (0) 2024.08.07
11장. 서브젝트 요약  (0) 2024.08.06
10장. 에러 처리 요약  (0) 2024.08.05
9장. 조건 연산자 요약  (0) 2024.08.05
8장. 유틸리티 연산자 요약  (0) 2024.08.02
반응형

1. 핫 옵저버블과 콜드 옵저버블

1-1. 핫 / 콜드 옵저버블

핫 옵저버블이란?

옵저버블이 푸시하는 값을 여러 옵저버에 멀티캐스팅하는 옵저버블이다. 옵저버블이므로 next, error, complete 함수를 제공하지 않고 옵저버블 내부에서 멀티캐스팅 할 값을 푸시한다.

 

콜드 옵저버블이란?

멀티캐스팅을 지원하지 않는 옵저버블이다. 여러 옵저버가 어떤 옵저버블을 구독하든 각 구독은 독립적으로 동작하며 옵저버블에서 푸시하는 값이 여러 옵저버에 공유되지 않는다.

/* 콜드 옵저버블 예 */
const { interval } = require('rxjs');
const { take } = require('rxjs/operators');

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 intervalSource$ = interval(500).pipe(take(5));

intervalSource$.subscribe(observerA);
setTimeout(() => intervalSource$.subscribe(observerB), 1000);

[코드 12-1] 콜드 옵저버블 예

 

실행 결과

observerA: 0
observerA: 1
observerB: 0
observerA: 2
observerB: 1
observerA: 3
observerB: 2
observerA: 4
observerA: complete
observerB: 3
observerB: 4
observerB: complete

콜드 옵저버블의 구독 각각은 독립적으로 동작할 뿐 멀티캐스팅으로 값을 공유하지 않는다.

1-2. 서브젝트와 연결하여 핫 옵저버블 흉내내기

/* connect 연산자를 이용해 서브젝트와 연결 */
const { interval, Subject } = require('rxjs');
const { take, tap } = require('rxjs/operators');

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')
};

function createHotObservable(sourceObservable, subject) {
  return {
    connect: () => sourceObservable.subscribe(subject),
    subscribe: subject.subscribe.bind(subject)
  };
}

const sourceObservable$ = interval(500).pipe(
  take(5),
  tap(x => console.log(`tap ${x}`))
);

const hotObservableExample = createHotObservable(sourceObservable$, new Subject());

hotObservableExample.subscribe(observerA);
console.log('observerA subscribe');
hotObservableExample.subscribe(observerB);
console.log('observerB subscribe');

hotObservableExample.connect();
console.log('connect called');

setTimeout(() => {
  console.log('1000ms...');
  hotObservableExample.subscribe(observerC);
  console.log('observerC subscribe');
}, 1000);

[코드 12-2] connect 연산자를 이용해 서브젝트와 연결

 

실행 결과

observerA subscribe
observerB subscribe
connect called
tap 0
observerA: 0
observerB: 0
1000ms...
observerC subscribe
tap 1
observerA: 1
observerB: 1
observerC: 1
tap 2
observerA: 2
observerB: 2
observerC: 2
tap 3
observerA: 3
observerB: 3
observerC: 3
tap 4
observerA: 4
observerB: 4
observerC: 4
observerA: complete
observerB: complete
observerC: complete

앞으로 소개할 연산자는 기존 옵저버블을 ConnectableObservable 로 변환한다. 이 옵저버블에는 connect 함수가 있고, 소스 옵저버블이 내부에 있는 서브젝트를 구독하도록 동작한다.

 

[코드 12-2]의 subscribe 함수는 내부에 있는 subject 를 구독하는 것이고, connect 함수는 소스 옵저버블에서 값을 발행해 subject 로 보낸다.

 

ConnectableObservable 은 해당 옵저버블을 만든 소스 옵저버블을 바로 구독하는 것이 아니라 내부의 서브젝트를 구독하는 것이다.


2. multicast 연산자

multicast 연산자는 소스 옵저버블로부터 multicast 를 호출할 때 서브젝트 팩토리 함수를 사용해 커넥터블 옵저버블을 만들어 핫 옵저버블을 다룰 수 있다.

연산자 원형
multicast<T, R>(
       subjectOrSubjectFactory: Subject<T> | (() => Subejct<T>),
       selector?: (source: Observable<T>) => Observable<R>
): OperatorFunction<T, R>
  • 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

소스 옵저버블 대신 소스 옵저버블과 연결된 서브젝트를 사용하므로 소스 옵저버블을 한 번만 구독한다. 때문에, 소스 옵저버블을 한 번만 구독한 후 발행한 값을 선택자 함수에서 제공하는 스트림을 거쳐서 출력한다.

/* multicast 연산자의 구현 코드 일부 */
export class MulticastOperator {
  constructor(subjectFactory, selector) {
    this.subjectFactory = subjectFactory;
    this.selector = selector;
  }

  call(subscriber, source) {
    const { selector } = this;
    const subject = this.subjectFactory();
    const subscription = this.selector(subject).subscribe(subscriber);
    subscription.add(source.subscribe(subject));
    return subscription;
  }
}

[코드 12-7] multicast 연산자의 구현 코드 일부

 

구독할 때 순서

  1. 팩토리 함수를 호출해 서브젝트 리턴
  2. 선택자 함수에서 서브젝트를 사용한 후 호출했을 때 리턴되는 결과를 구독
  3. 서브젝트의 소스 옵저버블을 연결하여 구독 목록에 추가

3. publish 연산자

publish 연산자는 서브젝트나 서브젝트의 팩토리 함수를 사용할 필요가 없도록 추상화한 연산자다.

/* publish 연산자의 구현 코드 일부 */
import {multicast} from "rxjs/operators";
import {Subject} from "rxjs";

export function publish(selector) {
  return selector ? 
    multicast(() => new Subject(), selector) : 
    multicast(new Subject());
}

[코드 12-8] publish 연산자의 구현 코드 일부

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

연산자 원형
publish<T, R>(
       selector?: OperatorFunction<T, R>
): MonoTypeOperatorFunction<T> | OperatorFunction<T, R>
  • selector? - 선택자 함수다. 소스 옵저버블을 여러 번 구독하지 않고 서브젝트를 이용해 소스 옵버버블을 필요할 때마다 사용할 수 있다.

선택자 함수가 없는 기본 동작은 publish 연산자에서 생성한 서브젝트 인스턴스를 이용해서 멀티캐스팅할 수 있다. 선택자 함수가 있다면 당연히 connect 함수를 호출할 수 없다.

 

중요한 점

  1. 같은 서브젝트 객체를 공유하므로 소스 옵저버블 구독을 완료하면 내부에 생성한 서브젝트도 사용할 수 없다.
  2. connect 함수를 호출한 후 소스 옵저버블을 구독하다 오나료하면 다시 connect 함수를 호출해도 이후 구독하는 옵저버들이 값을 전달 받을수 없다 서브젝트를 사용할 수 없으므로 이를 구독하는 옵저버들은 값을 전달 받을수 없기 때문이다.
/* 서브젝트 객체를 재구독할 때 발생할 수 있는 문제 */
const { interval, Subject } = require('rxjs');
const { multicast, take, tap, publish } = require('rxjs/operators');

// const testSource$ = interval(500).pipe(
//   take(5),
//   tap(x => console.log(`tap ${x}`)),
//   multicast(() => new Subject())
// );

const testSource$ = interval(500).pipe(
  take(5),
  tap(x => console.log(`tap ${x}`)),
  publish()
);

const a = testSource$.subscribe(x => console.log(`a: ${x}`));
const b = testSource$.subscribe(x => console.log(`b: ${x}`));

testSource$.connect();

setTimeout(() => {
  console.log('timeout');
  a.unsubscribe();
  b.unsubscribe();
  testSource$.subscribe(x => console.log(`c: ${x}`));
  testSource$.connect();
}, 3000);

[코드 12-9] 서브젝트 객체를 재구독할 때 발생할 수 있는 문제

 

실행 결과

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 연산자만 호출하고 멀티캐스팅 되지 않는다.

 

실행 순서

  1. a, b 를 출력하는 두 옵저버블에서 publish 연산자를 호출해 만든 testSource$ 를 구독하고 connect 함수를 호출하여 500ms 마다 0~4까지 5개의 숫자를 멀티캐스팅한다.
  2. 5개 숫자를 모두 발행한 후 complete 함수를 호출한 3초 후에 재구독하도록 setTimeout 함수를 호출한다.
  3. 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>>
연산자 원형 - publishReplay
publishReplay<T, R>(
       bufferSize?: number,
       windowTime?: number,
       selectorOrScheduler?: SchedulerLike | OperatorFunction<T, R>,
       scheduler?: SchedulerLike
): UnaryFunction<Observable<T>, ConnectableObservable<R>>
연산자 원형 - 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 함수까지 자동으로 호출해준다.

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

연산자 원형
refCount<T>(): MonoTypeOperatorFunction<T>
/* 커넥터블 옵저저블에 refCount 연산자 추가 */
const { interval, Subject } = require('rxjs');
const { take, tap, multicast, publish, refCount } = require('rxjs/operators');

// const testSource$ = interval(500).pipe(
//   take(5),
//   tap(x => console.log(`tap ${x}`)),
//   multicast(new Subject()),
//   refCount()
// );

const testSource$ = interval(500).pipe(
  take(5),
  tap(x => console.log(`tap ${x}`)),
  publish(),
  refCount()
);

const a = testSource$.subscribe(x => console.log(`a: ${x}`));
const b = testSource$.subscribe(x => console.log(`b: ${x}`));

setTimeout(() => {
  console.log('timeout');
  testSource$.subscribe(x => console.log(`c: ${x}`));
}, 3000);

[코드 12-13] 커넥터블 옵저버블에 refCount 연산자 추가

 

실행 결과

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이 된 이후 재구독을 하여도 새로운 값을 전달 받을 수 있다.

/* share 연산자와 publish.refCount 를 사용했을 때의 차이 */
const {interval} = require('rxjs');
const { take, tap, publish, refCount, share } = require('rxjs/operators');

const testSource$ = interval(500).pipe(
  take(5),
  tap(x => console.log(`tap ${x}`)),
  share()
);

// const testSource$ = interval(500).pipe(
//   take(5),
//   tap(x => console.log(`tap ${x}`)),
//   publish(),
//   refCount()
// );

const a = testSource$.subscribe(x => console.log(`a: ${x}`));
const b = testSource$.subscribe(x => console.log(`b: ${x}`));

setTimeout(() => {
  console.log('timeout');
  testSource$.subscribe(x => console.log(`c: ${x}`));
}, 3000);

[코드 12-16] share 연산자와 publish.refCount 를 사용했을 때의 차이 - share 연산자 사용

 

실행 결과 - share 연산자 사용

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

share 연산자는 재구독을 할 수 있다.

/* share 연산자와 publish.refCount 를 사용했을 때의 차이 */
const {interval} = require('rxjs');
const { take, tap, publish, refCount, share } = require('rxjs/operators');

// const testSource$ = interval(500).pipe(
//   take(5),
//   tap(x => console.log(`tap ${x}`)),
//   share()
// );

const testSource$ = interval(500).pipe(
  take(5),
  tap(x => console.log(`tap ${x}`)),
  publish(),
  refCount()
);

const a = testSource$.subscribe(x => console.log(`a: ${x}`));
const b = testSource$.subscribe(x => console.log(`b: ${x}`));

setTimeout(() => {
  console.log('timeout');
  testSource$.subscribe(x => console.log(`c: ${x}`));
}, 3000);

[코드 12-16] share 연산자와 publish.refCount 를 사용했을 때의 차이 - publish.refCount 연산자 사용

 

실행 결과 - publish.refCount 연산자 사용

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

pipe(publish(), refCount()) 는 재구독을 할수 없다.


6. 마치며

커넥터블 옵저버블 - connect 함수를 제공하여 소스 옵저버블과 서브젝트를 연결시켜서 멀티캐스팅을 지원한다.

share 연산자 - 재구독 가능

publish 연산자 - 재구독 불가능 (서브젝트 자체를 사용하기 때문)

반응형

'RxJS' 카테고리의 다른 글

13장. 스케줄러 요약  (0) 2024.08.08
11장. 서브젝트 요약  (0) 2024.08.06
10장. 에러 처리 요약  (0) 2024.08.05
9장. 조건 연산자 요약  (0) 2024.08.05
8장. 유틸리티 연산자 요약  (0) 2024.08.02
반응형

1. 서브젝트의 특성

  • 콜드 옵저버블 - 멜티캐스팅을 지원하지 않는 옵저버블
  • 핫 옵저버블 - 멀티캐스팅을 지원하는 옵저버블

서브젝트란?

멀티캐스팅을 지원하기위해 옵저버이면서 옵저버블이라는 특성이 있다. 따라서, subscribe 함수를 호출했을 때 옵저버를 등록한다. 등록된 옵저버들은 서브젝트가 보내는 값이나 이벤트, 에러, 구독 완료 등의 정보를 받을 수 있다.

/* 옵저버블로 사용하는 서브젝트 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);

[코드 11-1] 옵저버블로 사용하는 서브젝트

 

서브젝트가 옵저버로서의 특성이 있으므로 각 옵저버로 값, 에러, 완료는 next, error, complete 함수를 호출해서 보낼 수 있다.

/* 옵저버로 사용하는 서브젝트 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);

subject.next(1);
subject.next(2);
subject.next(3);

[코드 11-2] 옵저버로 사용하는 서브젝트

 

실행 결과

observerA: 1
observerB: 1
observerC: 1
observerA: 2
observerB: 2
observerC: 2
observerA: 3
observerB: 3
observerC: 3

서브젝트의 next 함수를 호출할 때는 subscribe 함수로 등록한 옵저버들에게 값을 전파한다. 옵저버로서 error, complete 함수도 호출할 수 있다.

/* error 함수 호출 후 next 함수 호출 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);

subject.error(new Error('error!'));
subject.next(4);
subject.complete();

[코드 11-3] error 함수 호출 후 next 함수 호출

 

실행 결과

observerA: Error: error!
observerB: Error: error!
observerC: Error: error!
/* complete 함수 호출 후 next 함수 호출 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.subscribe(observerA);
subject.subscribe(observerB);
subject.subscribe(observerC);

subject.complete();
subject.next(4);
subject.error(new Error('error!'));

[코드 11-4] complete 함수 호출 후 next 함수 호출

 

실행 결과

observerA: complete
observerB: complete
observerC: complete

2. 서브젝트와 옵저버블의 연결

/* interval 생성 함수를 이용하는 콜드 옵저버블 동작 */
const { interval } = require('rxjs');
const { take } = require('rxjs/operators');

const intervalSource$ = interval(500).pipe(take(5));

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.log(`observerA: ${e}`),
  complete: () => console.log('observerA: complete'),
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.log(`observerB: ${e}`),
  complete: () => console.log('observerB: complete'),
};

intervalSource$.subscribe(observerA);
setTimeout(() => {
  intervalSource$.subscribe(observerB);
}, 2000);

[코드 11-5] interval 생성 함수를 이용하는 콜드 옵저버블 동작

 

실행 결과

observerA: 0
observerA: 1
observerB: 0
observerA: 2
observerB: 1
observerA: 3
observerB: 2
observerA: 4
observerA: complete
observerB: 3
observerB: 4
observerB: complete

콜드 옵저버블의 동작 원리

subscribe 함수를 호출하는 각 옵저버블 구독이 따로 동작하며 매번 새로 구독하는 구조다.

 

서브벡트를 이용해서 어떻게 기존 옵저버가 멀티캐스팅되는 구조로 연결할 수 있을까? 서브젝트를 선언하고 서브젝트를 구독하는 subject.subscribe 가 필요하다.

/* 서브젝트를 새로 선언해 대체 */
const { interval, Subject } = require('rxjs');
const { take } = require('rxjs/operators');

const subject = new Subject();
const intervalSource$ = interval(500).pipe(take(5));

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.log(`observerA: ${e}`),
  complete: () => console.log('observerA: complete'),
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.log(`observerB: ${e}`),
  complete: () => console.log('observerB: complete'),
};

subject.subscribe(observerA);
setTimeout(() => {
  subject.subscribe(observerB);
}, 2000);

[코드 11-6] 서브젝트를 새로 선언해 대체

 

[코드 11-5]를 기반으로 서브젝트를 새로 선언하고 intervalSource$ 대신 해당 자리를 subject 로 바꿨다.

다음으로 서브젝트에서 next, error, complete 함수로 값을 전달해야 한다.

/* intervalSource$ 를 구독해 서브젝트로 보내는 작업 추가 */
const { interval, Subject } = require('rxjs');
const { take } = require('rxjs/operators');

const subject = new Subject();
const intervalSource$ = interval(500).pipe(take(5));

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.log(`observerA: ${e}`),
  complete: () => console.log('observerA: complete'),
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.log(`observerB: ${e}`),
  complete: () => console.log('observerB: complete'),
};

subject.subscribe(observerA);

intervalSource$.subscribe({
  next: x => subject.next(x),
  error: e => subject.error(e),
  complete: () => subject.complete()
});

setTimeout(() => {
  subject.subscribe(observerB);
}, 2000);

[코드 11-7] intervalSource$를 구독해 서브젝트로 보내는 작업 추가

 

실행 결과

observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerB: 3
observerA: 4
observerB: 4
observerA: complete
observerB: complete

intervalSource$ 구독을 시작하고 2초 후에 observerA를 구독한다는 것을 보여주어야 하므로 setTimeout 직전에 코드를 위치 시켰다.

서브젝트의 동작 원리

서브젝트는 서브젝트 하나에서 스트림 하나가 여러 옵저버로 전파되는 구조이다.

/* 옵저버블에서 값을 바로 서브젝트로 전달 */
const { interval, Subject } = require('rxjs');
const { take } = require('rxjs/operators');

const subject = new Subject();
const intervalSource$ = interval(500).pipe(take(5));

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.log(`observerA: ${e}`),
  complete: () => console.log('observerA: complete'),
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.log(`observerB: ${e}`),
  complete: () => console.log('observerB: complete'),
};

subject.subscribe(observerA);

intervalSource$.subscribe(subject);

setTimeout(() => {
  subject.subscribe(observerB);
}, 2000);

[코드 11-8] 옵저버블에서 값을 바로 서브젝트로 전달

 

실행 결과

observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerB: 3
observerA: 4
observerB: 4
observerA: complete
observerB: complete

서브젝트는 옵저버이기도 하므로 옵저버블에서 값을 바로 서브젝트로 전달해줘서 함수 중복을 피할 수 있다. 즉, 서브젝트의 옵저버블 특성을 잘 활용하면 옵저버블과 연결하거나 직접 옵저버의 함수들을 호출해서 멀티캐스팅하려는 값, 이벤트, 에러, 완료에 관한 정보를 보낼 수 있다.


3. 서브젝트의 에러와 완료 처리

Subject 의 내부 함수

  • next 함수 - subscribe 함수 호출 전 전달한 값은 이후 구독하는 옵저버로 전달하지 않는다.
  • error 함수 - 호출 결과는 이후에 구독하는 옵저버에게도 전파한다.
  • complete 함수 - 호출 결과는 이후에 구독하는 옵저버에게도 전파한다.

즉, 이미 해당 서브젝트에서 에러가 발생했거나 서브젝트 구독을 완료했다는 것을 알려준다.

/* 서브젝트의 에러 발생 상황을 전파 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.error('error');
subject.subscribe(observerA);
subject.subscribe(observerB);

[코드 11-9] 서브젝트의 에러 발생 상황을 전파

 

실행 결과

observerA: error
observerB: error
/* 서브젝트의 구독 완료 상황을 전파 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.complete();
subject.subscribe(observerA);
subject.subscribe(observerB);

[코드 11-10] 서브젝트의 구독 완료 상황을 전파

 

실행 결과

observerA: complete
observerB: complete

 

서브젝트는 unsubscribe 함수도 제공한다. unsubscribe 함수를 호출하면 아무 일도 일어나지 않은 것 같지만 이후 모든 옵저버블 대상으로 멀티케스팅할 수 없다. next, error, complete 함수를 호출할 때도 에러가 발생하며 멀티캐스팅하려고 특정 옵저버가 구독을 시도해도 에러가 발생한다. 이는 등록된 옵저버가 있는 배열을 null 로 만들며, 더 사용할 수 없는 서브젝트로 취급해 closed 플래그를 true 로 인식해 에러가 발생하는 것이다.

/* 서브젝트의 unsubscribe 함수의 동작 */
const { Subject } = require('rxjs');

const subject = new Subject();

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

subject.subscribe(observerA);
subject.subscribe(observerB);

subject.unsubscribe();

// subject 구독 해제 후 다시 구독한다.
subject.subscribe(observerC);

// 하나씩만 주석 처리를 해제한 후 코드를 실행한다.
// subject.next(1);
// subject.error('error');
// subject.complete();

[코드 11-11] 서브젝트의 unsubscribe 함수의 동작

 

실행 결과

[Error [ObjectUnsubscribedError]: object unsubscribed]

4. 서브젝트의 종류

  • BehaviorSubject - 시간과 같은 연속인 값을 다루는 구조에 적합하다. 초기값이 있어 언제 구독해도 항상 값이 있다.
  • ReplaySubject - 서브젝트를 생성할 때 인자로 설정한 수만큼 최근 전달받은 아이템을 갖고 있다가 다음 구독할 때 해당 수만큼 이벤트를 전달한다.
  • AsyncSubject - 서브젝트 구독 완료 후 가장 마지막에 있는 아이템을 전달한다.

4-1. BehaviorSubject

BehaviorSubject 는 subscribe 함수를 호출하자마자 next 함수에서 전달 받을 수 있는 초기값이 있다. 생성할 때 초기값을 전달하고, 옵저버의 하수가 한 번도 호출되지 않아 아무 값도 전달받지 않는 다면 초기값을 그대로 사용한다. error 나 complete 함수가 호출되지 않았다면 subscribe 함수를 호출할 때마다 최근에 next 함수에서 전달받은 값을 준다. 즉, 초기값이 최근 값이든 subscribe 함수를 호출하자마자 전달받을 수 있는 값이 있는 서브젝트다.

/* BehaviorSubject 의 기본 동작 예 */
const { BehaviorSubject } = require('rxjs');

const behaviorSubject = new BehaviorSubject('초기값');

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const observerC = {
  next: x => console.log(`observerC: ${x}`),
  error: e => console.error(`observerC: ${e}`),
  complete: () => console.log('observerC: complete')
};

behaviorSubject.subscribe(observerA);
behaviorSubject.next('값1');

behaviorSubject.subscribe(observerB);
behaviorSubject.next('값2');

behaviorSubject.subscribe(observerC);
behaviorSubject.next('값3');
behaviorSubject.next('값4');
behaviorSubject.next('값5');

[코드 11-12] BehaviorSubject 의 기본 동작 예

 

실행 결과

observerA: 초기값
observerA: 값1
observerB: 값1
observerA: 값2
observerB: 값2
observerC: 값2
observerA: 값3
observerB: 값3
observerC: 값3
observerA: 값4
observerB: 값4
observerC: 값4
observerA: 값5
observerB: 값5
observerC: 값5

BehaviorSubject 는 어떤 옵저버든 subcribe 함수로 호출할 때마다 바로 전달받을 수 있는 값이 있다. 이후에는 멀티캐스팅으로 구독하는 모든 옵저버가 next 함수로 값을 전달 받을 수 있는 구조다.

현재 값을 가져올 수 있는 value 또는 getValue 함수

BehaviorSubject 는 서브젝트 중 유일하게 항상 값을 갖고 있기 때문에, 현재 값을 subcribe 함수의 호출 없이 바로 전달 받을 수 있는 게터함수 getValue 를 제공한다. behaviorSubject.value 로 바도 접근할 수 있다.

/* BehaviorSubject 룰 이용하여 구현한 숫자 동작 */
const { BehaviorSubject, interval } = require('rxjs');
const { take, map } = require('rxjs/operators');

const behaviorSubject = new BehaviorSubject(0);

const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};

const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

const incrementInterval$ = interval(1000).pipe(
  take(5),
  map(x => behaviorSubject.value + 1), // 최신 값에서 1 증가시킨 값으로 변환
  // map(x => behaviorSubject.getValue() + 1)
);

// incrementInterval$ 를 behaviorSubject 에 연결하여 구독 시작
incrementInterval$.subscribe(behaviorSubject);

// observerA 바로 구독
behaviorSubject.subscribe(observerA);

// observerB 는 3.2초 후 구독해 가장 최신 값 3이 바로 나오는지 확인
setTimeout(() => behaviorSubject.subscribe(observerB), 3200);

[코드 11-13] BehaviorSubject 를 이용하여 구현한 숫자 동작

 

실행 결과

observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerB: 3
observerA: 4
observerB: 4
observerA: 5
observerB: 5
observerA: complete
observerB: complete

4-2. ReplaySubject

ReplaySubject 는 next 함수로 지정한 개수만큼 연속해서 전달한 최신 값을 저장했다가 다음 구독 때 해당 개수만큼 옵저버로 발행한다. 그 후 멀티캐스팅되는 값을 발행하는 서브젝트다.

 

서브젝트를 생성할 때 연속해서 전달해야 하는 값 개수를 지정할 수 있다. 개수를 지정하지 않으면 메모리와 관련한 성능 문제가 발생한다.

/* ReplaySubject 의 기본 사용 예 */
const { ReplaySubject, interval } = require('rxjs');
const { take } = require('rxjs/operators');

const replaySubject = new ReplaySubject(3);
const interval$ = interval(500).pipe(take(8));
const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};
const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

console.log('try replaySubject.subscribe(observerA)');
replaySubject.subscribe(observerA);

console.log('try interval$.subscribe(replaySubject)');
interval$.subscribe(replaySubject);

setTimeout(() => {
  console.log('try replaySubject.subscribe(observerB) setTimeout 2600ms');
  replaySubject.subscribe(observerB);
}, 2600)

[코드 11-14] ReplaySubject 의 기본 사용 예

 

실행 결과

try replaySubject.subscribe(observerA)
try interval$.subscribe(replaySubject)
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerA: 4
try replaySubject.subscribe(observerB) setTimeout 2600ms
observerB: 2
observerB: 3
observerB: 4
observerA: 5
observerB: 5
observerA: 6
observerB: 6
observerA: 7
observerB: 7
observerA: complete
observerB: complete

연속해서 전달할 값 개수는 3개로 지정했고 observerB 를 구독하면 가장 최근 저장한 값 3개인 2, 3, 4를 전달해 발행한다.

제한 없이 값을 전달하는 ReplaySubject

ReplaySubject 생성 부분에서 인자를 사용하지 않으면 제한 없이 연속해서 값을 전달할 수 있다.

/* 제한 없이 값을 전달하는 ReplaySubject 의 사용 예 */
const { ReplaySubject, interval } = require('rxjs');
const { take } = require('rxjs/operators');

const replaySubject = new ReplaySubject();
const interval$ = interval(500).pipe(take(8));
const observerA = {
  next: x => console.log(`observerA: ${x}`),
  error: e => console.error(`observerA: ${e}`),
  complete: () => console.log('observerA: complete')
};
const observerB = {
  next: x => console.log(`observerB: ${x}`),
  error: e => console.error(`observerB: ${e}`),
  complete: () => console.log('observerB: complete')
};

console.log('try replaySubject.subscribe(observerA)');
replaySubject.subscribe(observerA);

console.log('try interval$.subscribe(replaySubject)');
interval$.subscribe(replaySubject);

setTimeout(() => {
  console.log('try replaySubject.subscribe(observerB), setTimeout 2600ms');
  replaySubject.subscribe(observerB);
}, 2600)

[코드 11-15] 제한 없이 값을 전달하는 ReplaySubject 의 사용 예

 

실행 결과

try replaySubject.subscribe(observerA)
try interval$.subscribe(replaySubject)
observerA: 0
observerA: 1
observerA: 2
observerA: 3
observerA: 4
try replaySubject.subscribe(observerB), setTimeout 2600ms
observerB: 0
observerB: 1
observerB: 2
observerB: 3
observerB: 4
observerA: 5
observerB: 5
observerA: 6
observerB: 6
observerA: 7
observerB: 7
observerA: complete
observerB: complete

4-3. AsyncSubject

AsyncSubject 는 비동기로 실행한 동작이 완료되면 마지막 결과를 받는 역할을 한다. 즉, AsyncSubject 로 비동기 동작을 실행한 후 complete 함수를 호출하기 직전에 비동기 연산의 마지막 결과를 전달해야 한다.

/* 피보나치 수열 옵저버블에 AsyncSubject 연산자 사용 */
const { interval, AsyncSubject } = require('rxjs');
const { take, scan, pluck, tap } = require('rxjs/operators');

const asyncSubject = new AsyncSubject();

const period = 500;
const lastN = 8;

const fibonacci = n => interval(period).pipe(
  take(n),
  scan((acc, index) => acc ? { a: acc.b, b: acc.a + acc.b } : { a: 0, b: 1 }, null),
  pluck('a'),
  tap(n => console.log(`tap log: emitting ${n}`))
);

fibonacci(lastN).subscribe(asyncSubject);

asyncSubject.subscribe(result => console.log(`1st subscribe: ${result}`));

setTimeout(() => {
  console.log('try 2nd subscribe');
  asyncSubject.subscribe(result => console.log(`2nd subscribe: ${result}`));
}, period * lastN + 1000);

[코드 11-16] 피보나치 수열 옵저버블에 AsyncSubject 연산자 사용

 

실행 결과

tap log: emitting 0
tap log: emitting 1
tap log: emitting 1
tap log: emitting 2
tap log: emitting 3
tap log: emitting 5
tap log: emitting 8
tap log: emitting 13
1st subscribe: 13
try 2nd subscribe
2nd subscribe: 13

finbonacci 함수는 n을 인자로 사용해 해당 개수만큼 period 에서 설정한 시간마다 피보나치 수열을 발행하는 옵저버블을 리턴한다. AsyncSubject 를 테스트하기 위해 첫번째 구독은 실제 구독 완료될 때까지 기다린 후 마지막 결과만 나오는지 확인하는 용도이며, 두번째 구독은 첫번째 구독이 완료된 이후 구독하면 바로 결과가 나오는지 확인하는 용도다.


5. 마치며

서브젝트는 옵저버의 특성과 옵저버블의 특성을 동시에 갖고 멀티캐스팅을 지원한다. 옵저버블과 연결해 멀티캐스팅을 지원하지 않는 옵저버블의 발행 값을 여러 옵저버로 전파할 수 있다.

반응형

'RxJS' 카테고리의 다른 글

13장. 스케줄러 요약  (0) 2024.08.08
12장. 멀티캐스팅 연산자 요약  (0) 2024.08.07
10장. 에러 처리 요약  (0) 2024.08.05
9장. 조건 연산자 요약  (0) 2024.08.05
8장. 유틸리티 연산자 요약  (0) 2024.08.02
반응형

1. catchError 연산자

catchError 연산자는 에러가 발생했을 때 인자로 사용하는 선택자 함수로 해당 에러를 전달하여 리턴하는 옵저버블을 대신 구독하는 연산자다.

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

연산자 원형
catchError<T, R>(
       selector: (err: any, caught: Observable<T>) => ObservableInput<R>
): OperatorFunction<T, T | R>

소스 옵저버블에서 에러가 발생했을 때는 값을 더 발행하지 않는다. 이 때 catchError 연산자는 선택자(selector) 함수에서 리턴하는 옵저버블을 구독해 값을 발행한다.

  • 에러 발생후에도 값을 발행해야 할 때
  • 에러 발생 없이 프로그램 실행을 끝낼 때
/* catchError 연산자의 사용 예 */
const { from, of } = require('rxjs');
const { map, tap, pluck, catchError } = require('rxjs/operators');

const integers = ['1', '2', '3', 'r', '5'];

from(integers).pipe(
  map((value, index) => ({ value, index })),
  tap(valueIndex => {
    const { value } = valueIndex;
    const { index } = valueIndex;
    if (!Number.isInteger(parseInt(value, 10))) {
      const error = new TypeError(`${value}은(는) 정수가 아닙니다`);
      error.index = index;
      error.integerCheckError = true;
      throw error;
    }
  }),
  pluck('value'),
  catchError(err => {
    if (err.name === 'TypeError' && err.integerCheckError) {
      const catchArray = [err.message];
      const restArray = integers
        .slice(err.index, integers.length)
        .map(x => `에러 후 나머지 값 ${x}`);
      return from([err.message].concat(restArray));
    }
    return of(err.message);
  })
).subscribe(x => console.log(x), err => console.error(err));

[코드 10-1] catchError 연산자의 사용 예

 

실행 결과

1
2
3
r은(는) 정수가 아닙니다
에러 후 나머지 값 r
에러 후 나머지 값 5

tap 연산자에서 해당 값을 parseInt 로 바꿔 정수인지 검사하는데 r은 정수값이 아니므로 TypeError 를 전달한다. 이 때 에러각 발생한 지점의 index 값을 error에 넣고 catchError 연산자의 선택자 함수는 err 객체를 전달 받는다. 해당 레러가 정수 확인 중 전달된 에러가 맞다면 에러 메세지를 출력한 후 해당 index 부터 나머지 값을 자례로 발행한다. 다른 에러면 에러 메세지만 발행하도록 옵저버블을 리턴한다.

 

subscribe 함수에 error 함수도 설정했지만, catchError 연산자가 에러를 적절히 처리해줘서 해당 함수를 호출하지 않는다. 하지만, catchError 연산자의 선택자 함수에서 리턴받아 구독한 옵저버블에서 에러가 발생하면 error 함수를 호출한다.

1-1. mergeMap 연산자를 사용한 catchError 연산자 응용

[코드 10-1]은 에러가 발생한 원래 스트림에서 에러가 발생한 값만 처리하고 나머지 값들을 그대로 처리할 수 없다. 따라서 mergeMap 연산자를 추가한후 소스 옵저버블에서 발행하는 값 각각을 옵저버블로 감싸서 여기에 catchError 연산자를 적용해야 한다. 이렇게 하면 에러가 발생해도 나머지 값들을 계속 발행할 수 있다.

/* mergeMap 연산자 안에서 catchError 연산자 사용 */
const { from, of } = require('rxjs');
const { mergeMap, tap, catchError } = require('rxjs/operators');

from(['1', '2', '3', 'r', '5', '6', 'u', '8']).pipe(
  mergeMap(x => {
    return of(x).pipe(
      tap(value => {
        if (!Number.isInteger(parseInt(value, 10))) {
          throw new TypeError(`${value}은(는) 정수가 아닙니다`);
        }
      }),
      catchError(err => of(err.message))
    );
  })
).subscribe(x => console.log(x), err => console.error(err));

[코드 10-2] mergeMap 연산자 안에서 catchError 연산자 사용

 

실행 결과

1
2
3
r은(는) 정수가 아닙니다
5
6
u은(는) 정수가 아닙니다
8

에러가 없는 값을 그대로 발행하면서 에러가 발생했을 때는 에러 메세지를 발행하도록 동작한다.


2. retry 연산자

retry 연산자는 에러가 발생했을 때 인자로 설정한 정수값만큼 소스 옵저버블 구독을 재시도하는 연산자다. 인자로 설정한 값만큼 재시도하다가 다시 에러가 발생하면 이를 처리한다. 에러가 발생했을 때만 구독을 재시도 한다.

  • 서버와 통신하면서 간헐적으로 에러가 발생했을 때 통신을 재시도하는 용도로 사용

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

연산자 원형
retry<T>(count: number = -1): MonoTypeOperatorFunction<T>

count 는 에러 발생 전 재시도 횟수를 설정한다. 기본값은 -1인데 이는 에러가 발생하지 않는다면 재시도를 하지 않는다는 것이다.

retry 연산자는 구독을 재시도 할 때 소스 옵저버블을 다시 구독해서 처음부터 값을 발행한다.

/* retry 연산자의 사용 예 */
const { interval, of } = require('rxjs');
const { take, mergeMap, tap, retry, catchError } = require('rxjs/operators');

interval(100).pipe(
  take(30),
  mergeMap(x => {
    return of(x).pipe(
      tap(value => {
        if (Math.random() <= 0.3) {
          throw new Error(`RANDOM ERROR ${value}`);
        }
      }),
      retry(10),  // 재시도 횟수와 유무에 따라 에러를 방지할 수 있음
      catchError(err => of(err.message))
    );
  })
).subscribe(x => console.log(x), err => console.error(err));

[코드 10-3] retry 연산자의 사용 예

 

실행 결과 - retry 연산자 사용 할 경우

0
1
2
3
... 중간 생략
27
28
29

 

실행 결과 - retry 연산자 주석 처리할 경우

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 에서 리턴하는 옵저버블의 값을 발행하면 이어소 소스 옵저버블 구독을 재시도한다. 에러가 발행하면 전체 스트림의 구독을 종료한다. 따라서, 구독 재시도가 필요하면 재시도하기 전 해줘야 할 일을 처리한 후 아무 값이든 발행해 해당 스트림을 종료시키지 않도록 해야한다.

/* retryWhen 연산자의 사용 예 */
const { interval, of } = require('rxjs');
const { take, mergeMap, tap, retryWhen, scan, catchError } = require('rxjs/operators');

interval(100).pipe(
  take(30),
  mergeMap(x => {
    return of(x).pipe(
      tap(value => {
        if (Math.random() <= 0.3) {
          throw new Error(`RANDOM ERROR ${value}`);
        }
      }),
      retryWhen(errors => {
        return errors.pipe(
          scan((acc, error) => {
            return {
              count: acc.count + 1,
              error
            };
          }, { count: 0 }),
          tap(errorInfo => {
            console.error(`retryCount: ${errorInfo.count}, error message: ${errorInfo.error.message}`);
          })
        );
      }),
      catchError(err => of(err.message))
    );
  })
).subscribe(x => console.log(x), err => console.error(err));

[코드 10-4] retryWhen 연산자의 사용 예

 

실행 결과

0
1
retryCount: 1, error message: RANDOM ERROR 2
retryCount: 2, error message: RANDOM ERROR 2
2
retryCount: 1, error message: RANDOM ERROR 3
3
4
retryCount: 1, error message: RANDOM ERROR 5
5
6
7
8
9
10
11
12
13
retryCount: 1, error message: RANDOM ERROR 14
retryCount: 2, error message: RANDOM ERROR 14
retryCount: 3, error message: RANDOM ERROR 14
14
15
16
17
18
retryCount: 1, error message: RANDOM ERROR 19
retryCount: 2, error message: RANDOM ERROR 19
19
20
21
22
23
24
25
retryCount: 1, error message: RANDOM ERROR 26
retryCount: 2, error message: RANDOM ERROR 26
26
27
28
29

errors 로 에러를 전달 받으면 scan 연산자를 통해 count 를 누적해 몇 번째 구독 재시도인지와 그에 따른 에러 메세지를 로그로 출력된다.

 

[코드 10-4]를 응용하면 가끔 에러가 발생하는 서버 환경에서 어떤 에러 때문에 몇 번을 구독 재시도해 성공이나 실패했는지 테스트 할 수 있다.

3-1. n회 재시도 후 에러 없이 complete 함수 호출

/* 재시도 후 에러 없이 complete 함수 호출 */
const { interval, of } = require('rxjs');
const { take, mergeMap, tap, retryWhen, scan, catchError } = require('rxjs/operators');

interval(100).pipe(
  take(30),
  mergeMap(x => {
    return of(x).pipe(
      tap(value => {
        if (Math.random() <= 0.5) {
          throw new Error(`RANDOM ERROR ${value}`);
        }
      }),
      retryWhen(errors => {
        return errors.pipe(
          take(2),
          scan((acc, error) => {
            return {
              count: acc.count + 1,
              error
            };
          }, { count: 0 }),
          tap(errorInfo => {
            console.error(`retryCount: ${errorInfo.count}, error message: ${errorInfo.error.message}`);
          })
        );
      })
    );
  }),
  catchError(err => of(err.message))
).subscribe(x => console.log(x), err => console.error(err));

[코드 10-5] 재시도 후 에러 없이 complete 함수 호출

 

실행 결과

0
retryCount: 1, error message: RANDOM ERROR 1
1
retryCount: 1, error message: RANDOM ERROR 2
2
retryCount: 1, error message: RANDOM ERROR 3
3
4
5
retryCount: 1, error message: RANDOM ERROR 6
6
retryCount: 1, error message: RANDOM ERROR 7
retryCount: 2, error message: RANDOM ERROR 7
8
9
retryCount: 1, error message: RANDOM ERROR 10
10
retryCount: 1, error message: RANDOM ERROR 11
11
12
13
14
retryCount: 1, error message: RANDOM ERROR 15
retryCount: 2, error message: RANDOM ERROR 15
retryCount: 1, error message: RANDOM ERROR 16
16
17
18
19
retryCount: 1, error message: RANDOM ERROR 20
retryCount: 2, error message: RANDOM ERROR 20
20
21
22
retryCount: 1, error message: RANDOM ERROR 23
retryCount: 2, error message: RANDOM ERROR 23
retryCount: 1, error message: RANDOM ERROR 24
24
retryCount: 1, error message: RANDOM ERROR 25
25
retryCount: 1, error message: RANDOM ERROR 26
26
27
retryCount: 1, error message: RANDOM ERROR 28
28
retryCount: 1, error message: RANDOM ERROR 29
retryCount: 2, error message: RANDOM ERROR 29
29

에러가 발생했을 때 원하는 횟수만큼만 구독을 재시도하다가 마지막 시도에서도 에러가 발생하면 재시도와 에러 처리 없이 complete 함수를 호출한다.

 

여기서 20은 두번의 재시도 끝에 성공했으므로 20이 발행된 것이다. 23은 두번의 재시도에서도 에러가 발생하여 23을 구독 하지 않고 완료한 것이다.

3-2. n회 재시도 후 에러 처리

/* 마지막 재시도 후 에러 처리 */
const { interval, of, throwError} = require('rxjs');
const { take, mergeMap, tap, retryWhen, scan, catchError } = require('rxjs/operators');

const n = 2;

interval(100).pipe(
  take(30),
  mergeMap(x => {
    return of(x).pipe(
      tap(value => {
        if (Math.random() <= 0.5) {
          throw new Error(`RANDOM ERROR ${value}`);
        }
      }),
      retryWhen(errors => {
        return errors.pipe(
          scan((acc, error) => {
            return {
              count: acc.count + 1,
              error
            };
          }, { count: 0 }),
          mergeMap(errorInfo => {
            if (errorInfo.count === n + 1) {
              return throwError(errorInfo.error);
            }
            return of(errorInfo);
          }),
          tap(errorInfo => {
            console.error(`retryCount: ${errorInfo.count}, error message: ${errorInfo.error.message}`);
          })
        );
      }),
      catchError(err => of(err.message))
    );
  })
).subscribe(x => console.log(x), err => console.error(err));

[코드 10-6] 마지막 재시도 후 에러 처리

 

실행 결과

retryCount: 1, error message: RANDOM ERROR 0
retryCount: 2, error message: RANDOM ERROR 0
0
retryCount: 1, error message: RANDOM ERROR 1
1
2
retryCount: 1, error message: RANDOM ERROR 3
3
4
retryCount: 1, error message: RANDOM ERROR 5
5
6
retryCount: 1, error message: RANDOM ERROR 7
7
8
retryCount: 1, error message: RANDOM ERROR 9
retryCount: 2, error message: RANDOM ERROR 9
RANDOM ERROR 9
retryCount: 1, error message: RANDOM ERROR 10
retryCount: 2, error message: RANDOM ERROR 10
10
retryCount: 1, error message: RANDOM ERROR 11
retryCount: 2, error message: RANDOM ERROR 11
11
retryCount: 1, error message: RANDOM ERROR 12
12
13
retryCount: 1, error message: RANDOM ERROR 14
retryCount: 2, error message: RANDOM ERROR 14
14
retryCount: 1, error message: RANDOM ERROR 15
retryCount: 2, error message: RANDOM ERROR 15
15
16
retryCount: 1, error message: RANDOM ERROR 17
17
retryCount: 1, error message: RANDOM ERROR 18
18
retryCount: 1, error message: RANDOM ERROR 19
19
retryCount: 1, error message: RANDOM ERROR 20
20
retryCount: 1, error message: RANDOM ERROR 21
21
retryCount: 1, error message: RANDOM ERROR 22
22
23
retryCount: 1, error message: RANDOM ERROR 24
retryCount: 2, error message: RANDOM ERROR 24
RANDOM ERROR 24
retryCount: 1, error message: RANDOM ERROR 25
retryCount: 2, error message: RANDOM ERROR 25
RANDOM ERROR 25
26
retryCount: 1, error message: RANDOM ERROR 27
27
28
retryCount: 1, error message: RANDOM ERROR 29
29

catchError 연산자가 없다면 n + 1 번째에서 에러가 발생했을 때 전체 스트림이 종료된다.

  • 최대 n 번 구독 재시도하다가 에러가 발생한다면 retry 연산자를 사용
  • 에러 처리 전, 에러 순서, 에러 특성에 따라 처리해야 할 것 이 있다면 retryWhen  연산자를 사용

next(재시도) 함수를 호출할지, error(에러를 냄) 함수를 호출할지, complete(에러 없이 완료) 함수를 호출할지 적절히 선택해서 사용해야 한다.

반응형

'RxJS' 카테고리의 다른 글

12장. 멀티캐스팅 연산자 요약  (0) 2024.08.07
11장. 서브젝트 요약  (0) 2024.08.06
9장. 조건 연산자 요약  (0) 2024.08.05
8장. 유틸리티 연산자 요약  (0) 2024.08.02
7장. 수학 및 결합 연산자 요약  (0) 2024.08.02
반응형

조건 연산자란?

특정 조건에 맞는지 알려주는 boolean 값을 리턴해 특정 조건에 해당할 때 정해진 값을 발행하는 연산자다.

 

조건 연산자와 필터링 연산자와의 차이점

  • 필터링 연산자 - 소스 옵저버블에서 발행하는 값을 확인하는 연산자
  • 조건 연산자 - 이미 발행한 값이 아닌 소스 옵저버블의 특성에 따라서 조건 자체를 분기해야 할 때 사용하는 연산자

1. defaultIfEmpty 연산자

defaultIfEmpty 연산자는 소스 옵저버블이 empty 함수로 생성한 옵저버블일 때 인자로 설정한 기본값을 발행하는 연산자다.

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

empty 함수 (생성함수) - 소스 옵저버블이 아무 값도 발행하지 않고 완료하는 empty 옵저버블을 생성한다.

defaultIfEmpty 연산자 - empty 옵저버블을 생성할 때 기본값을 설정할 필요가 있을 때 사용

연산자 원형
defaultIfEmpty<T, R>(defaultValue: R = null): OperatorFunction<T, T | R>
/* defaultIfEmpty 연산자의 구현 코드 일부 */
_next(value) {
  this.isEmpty = false;
  this.destination.next(value);
}

_complete() {
  if (this.isEmpty) {
    this.destination.next(this.defaultValue);
  }
  this.destination.complete();
}

[코드 9-1] defaultIfEmpty 연산자의 구현 코드 일부

 

소스 옵저버블이 empty 옵저버블일 때는 기본값을 출력하고, empty 옵저버블이 아니라면 소스 옵저버블의 원래 동작을 처리한다.

/* defaultIfEmpty 연산자를 이용하는 empty 옵저버블 사용 예 */
const { range } = require('rxjs');
const { defaultIfEmpty } = require('rxjs/operators');

const getRangeObservable = count => range(1, count);

function subscribeWithDefaultIfEmpty(count) {
  getRangeObservable(count)
    .pipe(defaultIfEmpty('EMPTY'))
    .subscribe(value => console.log(`개수(count): ${count}, 값(value): ${value}`));
}

subscribeWithDefaultIfEmpty(0);
subscribeWithDefaultIfEmpty(3);

[코드 9-2] defaultIfEmpty 연산자를 이용하는 Empty 옵저버블 사용 예

 

실행 결과

개수(count): 0, 값(value): EMPTY
개수(count): 3, 값(value): 1
개수(count): 3, 값(value): 2
개수(count): 3, 값(value): 3

소스 옵저버블이 empty 옵저버블이면 complete 함수를 호출했을 때 기본값을 발행하고, 그렇지 않으면 소스 옵저버블의 원래 동작을 실행한다.


2. isEmpty 연산자

isEmpty 연산자는 true/false 를 값으로 발행해 소스 옵저버블이 empty 옵저버블인지 아닌지를 알려주고 구독 완료하는 연산자다.

 

isEmpty 연산자와 defaultIfEmpty 연산자의 차이점

  • defaultIfEmpty 연산자 - 소스 옵저버블이 empty 옵저버블이 아닐 때 소스 옵저버블의 원래 동작을 실행
  • isEmpty 연산자 - 소스 옵저버블이 empty 옵저버블인지 확인 후 true/false 값을 발행하고 구독을 완료

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

연산자 원형
isEmpty<T>(): OperatorFunction<T, boolean>
/* isEmpty 연산자의 사용 예 */
const { range } = require('rxjs');
const { isEmpty } = require('rxjs/operators');

const getRangeObservable = count => range(1, count);

function subscribeWithIsEmpty(count) {
  getRangeObservable(count)
    .pipe(isEmpty())
    .subscribe(value => console.log(`개수(count): ${count}, 값(value): ${value}`));
}

subscribeWithIsEmpty(0);
subscribeWithIsEmpty(3);

[코드 9-3] isEmpty 연산자의 사용 예

 

실행 결과

개수(count): 0, 값(value): true
개수(count): 3, 값(value): false

소스 옵저버블은 어떠한 값도 발행하지 않는다.


3. find 연산자

find 연산자는 인자로 사용하는 predicate 함수로 소스 옵저버블에서 발행하는 값 중 처음으로 함수 조건을 만족했을 때 true 를 리턴하는 값을 발행하고 구독을 완료하는 연산자다. 구독을 완료 할 때까지 조건을 만족하는 값이 없었다면 undefined 라는 값을 발행한다.

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

연산자 원형
find<T>(
       predicate: (value: T, index: number, source: Observable<T>) => boolean,
       thisArgs?: any
): MonoTypeOperatorFunction<T>

predicate 함수를 호출해 동작하는 중 에러가 발행하면 error 함수를 호출해 에러를 전달 받는다. 소스 옵저버블에서 에러가 발생해도 error 함수로 에러를 전파한다.

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

const getRangeObservable = count => range(1, count);

function subscribeWithFindGreaterThan3(count) {
  getRangeObservable(count)
    .pipe(find(x => x > 3))
    .subscribe(value => console.log(`개수(count): ${count}, 값(value): ${value}`));
}

subscribeWithFindGreaterThan3(5);
subscribeWithFindGreaterThan3(1);

[코드 9-4] find 연산자의 사용 예

 

실행 결과

개수(count): 5, 값(value): 4
개수(count): 1, 값(value): undefined
반응형

'RxJS' 카테고리의 다른 글

11장. 서브젝트 요약  (0) 2024.08.06
10장. 에러 처리 요약  (0) 2024.08.05
8장. 유틸리티 연산자 요약  (0) 2024.08.02
7장. 수학 및 결합 연산자 요약  (0) 2024.08.02
6장. 조합 연산자 요약  (0) 2024.08.01
반응형

1. tap 연산자

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 연산자의 마블 다이어그램

연산자 원형
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 연산자의 마블 다이어그램

연산자 원형
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
반응형

수학 및 결합 연산자의 종류

  • count 연산자
  • max 연산자
  • min 연산자
  • reduce 연산자

특징

  1. 수학 개념을 기반으로 만들었다.
  2. 결합 속성이 있다
  3. 소스 옵저버블에서 complete 함수를 호출해야 결과를 next 함수로 발행할 수 있다.

1. reduce 연산자

reduce 연산자는 누적자 함수를 이용해 소스 옵저버블에서 발행한 값을 누적한다. 소스 옵저버블에서 complete 함수를 호출하면 지금까지 누적한 결과를 next 함수로 발행하고 완료한다.

 

scan 연산자와의 차이점 - 매번 누적한 결과를 리턴하는 것이 아니다. 값을 계속 누적만 하고 발행하지 않다가 complete 함수를 호출할 때 한 번만 누적 결과를 발행한다.

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

연산자 원형
reduce<T, R>(
       accumulator: (acc: R, value: T, index?: number) => R,
       seed?: R
): OperatorFunction<T, R>

1-1. 초기값이 없음

초기값이 없으면 처음 발행하는 값을 초기갑으로 삼아 누적자 함수를 적용하지 않는다. 두번째 발행하는 값을 처음 발행하는 값에 누적하도록 동작한다. 만약 소스 옵저버블에서 1개의 값만 발행하고 complete 함수를 호출하면 이 값만 발행하고 구독을 완료한다.

/* 초기값 없이 reduce 연산자를 사용하는 예 */
const { of } = require('rxjs');
const { reduce } = require('rxjs/operators');

of(0).pipe(reduce((acc, cur) => acc + cur))
  .subscribe(result => console.log(`result: ${result}`));

[코드 7-1] 초기값 없이 reduce 연산자를 사용하는 예

 

실행 결과

result: 0

2개 이상 값을 전달받지 못하므로 누적자 함수를 호출할 수 없다. 따라서 0 그대로 발행한다.

/* reduce 연산자와 range 함수로 4개 값을 발행하는 예 */
const { range } = require('rxjs');
const { reduce } = require('rxjs/operators');

range(1, 4).pipe(reduce((acc, cur) => acc + cur))
  .subscribe(result => console.log(`result: ${result}`));

[코드 7-2] reduce 연산자와 range 함수로 4개 값을 발행하는 예

 

실행 결과

result: 10

누적자 함수 첫번째 인자는 지금까지 누적한 값이고, 두번째 인자는 현재 값이다. 따라서, 1부터 4까지 더한 10을 발행한다.

 

[코드7-2]에서 누적자 함수를 호출해 누적하는 과정

  1. 소스 옵저버블 1 -> 누적값 = 1 저장
  2. 소스 옵저버블 2 -> 누적자 함수 호출 (acc = 1, curr = 2), 1 + 2 = 3 계산 후, 누적값 = 3
  3. 소스 옵저버블 3 -> 누적자 함수 호출 (acc = 3, curr = 3), 3 + 3 = 6 계산 후, 누적값 = 6
  4. 소스 옵저버블 4 -> 누적자 함수 호출 (acc = 6, curr = 4), 6 + 4 = 10 계산 후, 누적값 = 10
  5. complete -> 누적값 10 발행 후 구독 완료 (next(10), complete 함수를 차례로 호출)

1-2. 초기값이 있음

/* 초기값이 있는 reduce 연산자를 사용하는 예 */
const { of } = require('rxjs');
const { reduce } = require('rxjs/operators');

of(0).pipe(reduce((acc, cur) => acc + cur, 1))
  .subscribe(result => console.log(`result: ${result}`));

[코드 7-3] 초기값이 있는 reduce 연산자를 사용하는 예

 

실행 결과

result: 1

소스 옵저버블은 0이란 값 1개만 발행하지만 누적자 함수로 초기값 1을 함께 누적해 최종 1이라는 값을 발행한다.

 

[코드7-3]실행 과정

  1. 소스 옵저버블 0 -> 누적자 함수 호출 (초기값 acc = 1, curr = 0), 1 + 0 = 1 계산 후, 누적값 = 1
  2. 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. 소스 옵저버블 1 -> 누적자 함수 호출 (초기값 acc = 1, curr = 1), 1 + 1 = 2 계산 후, 누적값 = 2
  2. 소스 옵저버블 2 -> 누적자 함수 호출 (acc = 2, curr = 2), 2 + 2 = 4 계산 후, 누적값 = 4
  3. 소스 옵저버블 3 -> 누적자 함수 호출 (acc = 4, curr = 3), 4 + 3 = 7 계산 후, 누적값 = 7
  4. 소스 옵저버블 4 -> 누적자 함수 호출 (acc = 7, curr = 4), 7 + 4 = 11 계산 후, 누적값 = 11
  5. 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}`));

[코드 7-6] comparer 함수를 사용하지 않는 예

 

실행 결과

result: 10

소스 옵저버블이 발행하는 값이 숫자이므로 부동호 바교하여 가장 큰 값인 10을 발행한다.

2-1. max 연산자에서 comparer 함수 사용

/* comparer 함수를 사용하는 예 */
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 }
];

from(movies).pipe(max((x, y) => x.avg - y.avg))
  .subscribe(x => console.log(JSON.stringify(x)));

[코드 7-7] comparer 함수를 사용하는 예

 

실행 결과

{"title":"영화 2","avg":9.14}

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 연산자는 소스 옵저버블에서 값을 발행할 때마다 개수를 내부에서 카운트한다. 구독을 완료한 후 총 몇 개인지 발행한다.

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

predicate 함수는 조건을 만족하는지 아닌지를 검사해 리턴한다. 즉, 소스 옵저버블에서 발행하는 값을 predicate 함수에 전달해 조건을 만족할 때 리턴하는 값만 카운트해서 마지막에 그 개수를 리턴한다.

 

원하는 값만 필터링한 개수를 카운트 하는 용도로 사용 할 수 있다.

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

4-1. 기본 동작

/* 인자 없는 count 연산자의 사용 예 */
const { range } = require('rxjs');
const { count } = require('rxjs/operators');

range(1, 20).pipe(count())
  .subscribe(result => console.log(`result: ${result}`));

[코드 7-13] 인자 없는 count 연산자의 사용 예

 

실행 결과

result: 20

발행하는 값이 총 20개 이므로 20을 마지막에 발행한다.

4-2. predicate 함수 사용

/* predicate 함수를 사용해 짝수만 카운트 */
const { range } = require('rxjs');
const { count } = require('rxjs/operators');

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

[코드 7-14] predicate 함수를 사용해 짝수만 카운트

 

실행 결과

result: 3

소스 옵저버블에서 발행하는 값 중 짝수에 해당하는 값은 3개 이므로 3을 출력한다.

반응형

'RxJS' 카테고리의 다른 글

9장. 조건 연산자 요약  (0) 2024.08.05
8장. 유틸리티 연산자 요약  (0) 2024.08.02
6장. 조합 연산자 요약  (0) 2024.08.01
5장. 변환 연산자 요약  (0) 2024.07.31
4장. 필터링 연산자 요약  (0) 2024.07.29

+ Recent posts