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, если источник ещё не завершился.