반응형

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

+ Recent posts