RxJS: Cold/Hot и мультикастинг

Cold vs Hot Observable

Cold:

  • Источник данных создаётся внутри конструктора Observable.
  • У каждого подписчика свой источник — значит, и значения свои.

Hot:

  • Источник данных создаётся снаружи конструктора, поэтому данные могут меняться, даже если нет подписчиков.
  • Поток делит значения из источника со всеми подписчиками.

Преобразование Cold в Hot

const observable = coldData();
const subject = new Subject<string>();
observable.subscribe(subject);
const s1 = subject.subscribe(console.log);

Тот же приём используют операторы мультикастинга (share, shareReplay):

const hot$ = cold$.pipe(shareReplay(1));

share / publish / refCount / multicast

⚠️ Deprecation (RxJS 7+): операторы multicast, refCount, publish, publishReplay, publishLast устарели и удалены в RxJS 8. Современная замена — share({ connector, resetOnRefCountZero, ... }) и connectable. Эквивалентности ниже приведены для понимания, как это работало.

share() позволяет нескольким подписчикам использовать один Observable без пересоздания источника — превращает cold в hot.

pokemon$ = http.get(/* ... */).pipe(share());
// теперь pokemon$ hot — не будет нескольких HTTP-запросов

Осторожно: hot Observable не реплицирует источник — при поздней подписке ранее эмитированные значения недоступны. Решение — shareReplay().

// концептуально (устаревший API):
share() === multicast(() => new Subject()).refCount();
  • share() — пока есть хотя бы один подписчик, Subject эмитит значения. Когда подписчиков нет — отключается от источника. Новый подписчик пересоздаёт подписку на источник через новый Subject.
  • publish() + refCount() — аналогично, но: если источник завершился, новые подписчики получают «completed»; если нет — Subject переподписывается на источник.
  • publishLast() === multicast(new AsyncSubject()) — не эмитит значения, пока источник не завершится, затем отдаёт последнее значение всем подписчикам.

shareReplay / publishReplay

Оба используют ReplaySubject, но shareReplay() не реализован через multicast(). Дают поздним подписчикам доступ к ранее эмитированным значениям.

const hotObs = basicColdObs.pipe(shareReplay(1));
hotObs.subscribe((d) => console.log("1st " + d));
hotObs.subscribe((d) => console.log("2nd " + d)); // получит последнее значение
  • shareReplay({ refCount: true, bufferSize: n }) — пока есть подписчики, ReplaySubject эмитит значения; при их отсутствии отключается от источника.
  • publishReplay(n) + refCount() — новые подписчики получают последние N значений и переподписываются на источник через тот же ReplaySubject, если источник ещё не завершился.