RxJS 7 in Practice: Cold vs Hot, Flattening Operator Races, Error Placement and Subscription Leaks
Key takeaways
Most RxJS bugs come from a handful of behaviours rather than from the operator catalogue: cold observables re-run their work per subscriber, the four flattening operators resolve overlapping requests differently, one uncaught error ends a stream for good, and shareReplay or a misplaced takeUntil can keep subscriptions alive forever. This guide walks through each one with outputs verified against RxJS 7.8.
RxJS is a library for composing asynchronous events as streams. Its reputation for difficulty comes less from the size of the operator list than from a few behaviours that are easy to get wrong and hard to see: work that runs once per subscriber, operators that silently drop or cancel requests, errors that end a stream permanently, and subscriptions that are never released. This article focuses on those behaviours. The outputs shown were checked with RxJS 7.8 in Node.
If you are new to asynchronous JavaScript in general, the Promises and event loop article covers the foundation this builds on.
What an Observable adds over a Promise
A Promise represents one future value, starts work as soon as it is created, and cannot be cancelled. An Observable is a function that can deliver zero or more values over time, does nothing until someone subscribes, and can be torn down by unsubscribing.
import { Observable } from 'rxjs';
const ticks$ = new Observable<number>((subscriber) => {
let n = 0;
const id = setInterval(() => subscriber.next(n++), 1000);
// Teardown: runs on unsubscribe, error or complete
return () => clearInterval(id);
});
const sub = ticks$.subscribe((n) => console.log(n));
setTimeout(() => sub.unsubscribe(), 3500); // logs 0, 1, 2, then the interval is cleared
The teardown function is the part that matters. It is why switchMap can cancel an HTTP request and why an unsubscribed fromEvent removes its DOM listener. It is also why a subscription you forget to release keeps an interval, a socket or a listener alive.
For a single request whose result you just await, a Promise is simpler, and RxJS does not make it better. RxJS earns its complexity when you combine several event sources, need cancellation, or need to control what happens when events overlap.
Cold vs hot: why your request runs twice
A cold observable starts its producer separately for every subscriber. of, from, interval, timer, ajax and Angular’s HttpClient are all cold.
let calls = 0;
const cold$ = new Observable<number>((s) => {
calls++;
s.next(calls);
s.complete();
});
cold$.subscribe((v) => console.log('A', v)); // A 1
cold$.subscribe((v) => console.log('B', v)); // B 2 <- the producer ran again
Replace the body with an HTTP call and each subscription is a separate network request. In Angular, the typical way to hit this is to use the same user$ twice in a template with two async pipes, or to subscribe once in the component and again in a child.
A hot observable shares one producer between subscribers, and a subscriber only sees values emitted after it subscribes. Subject, fromEvent on a DOM element and WebSocket messages behave this way:
const s = new Subject<number>();
s.next(1); // nobody is listening; this value is gone
s.subscribe((v) => console.log(v));
s.next(2); // logs 2 only
To turn a cold observable into a shared one, use the multicasting operators:
share()shares one subscription to the source while at least one subscriber exists. Late subscribers do not get past values. When the subscriber count drops to zero, it unsubscribes from the source and the next subscriber starts a fresh execution.shareReplay({ bufferSize: 1, refCount: true })does the same but also replays the latest value to late subscribers. This is the usual choice for “load this once and let every component read it”.
The shareReplay refCount pitfall
shareReplay(1) without options looks equivalent, but it never unsubscribes from the source, even after every subscriber has gone:
const prices$ = interval(1000).pipe(
tap(() => console.log('tick')),
shareReplay(1),
);
const sub = prices$.subscribe();
sub.unsubscribe();
// 'tick' keeps logging: the interval is still subscribed
With shareReplay({ bufferSize: 1, refCount: true }) the interval stops when the last subscriber leaves. I verified both behaviours: after unsubscribing, the plain version kept ticking and the refCount version did not.
For an HTTP request, the plain version is usually harmless, because the source completes after one response and there is nothing left to keep alive. It even acts as a permanent cache, which may or may not be what you want. The problem appears with sources that never complete: polling, WebSockets, store selections, fromEvent. When I chase a “the page gets slower the longer the tab stays open” complaint in an RxJS codebase, this is one of the first things I grep for: a shareReplay(1) over a polling interval, created again on every visit to a page, with none of the copies ever stopping. The fix is one option object.
The four flattening operators, with an actual race
switchMap, mergeMap, concatMap and exhaustMap all map each outer value to an inner observable and flatten the results. They differ only in what happens when a new outer value arrives while an inner one is still running. The documentation says this in words; a concrete race makes it clearer.
Three “clicks” arrive at 0, 10 and 20 ms, and each triggers a “request” that takes 25 ms:
import { timer, take, map, mergeMap, switchMap, concatMap, exhaustMap } from 'rxjs';
const clicks$ = timer(0, 10).pipe(take(3)); // emits 0, 1, 2
const request = (i: number) => timer(25).pipe(map(() => i));
clicks$.pipe(mergeMap(request)).subscribe(console.log); // 0, 1, 2
clicks$.pipe(switchMap(request)).subscribe(console.log); // 2
clicks$.pipe(concatMap(request)).subscribe(console.log); // 0, 1, 2 (one after another)
clicks$.pipe(exhaustMap(request)).subscribe(console.log); // 0
| Operator | When a new value arrives mid-request | Good for | Wrong for |
|---|---|---|---|
switchMap | Unsubscribes from the old inner (cancels it) | Search-as-you-type, route params, “latest wins” | Writes: a cancelled POST may still have reached the server |
mergeMap | Runs both concurrently | Independent parallel work (with a concurrent limit) | Anything where result order matters; responses can arrive out of order |
concatMap | Queues it until the current inner completes | Ordered writes, saving edits in sequence | Fast event sources: the queue grows without bound |
exhaustMap | Ignores it | Submit, login and refresh buttons | Search: the user’s last keystroke is dropped |
Two details are easy to miss:
switchMapcancellation is client-side. Unsubscribing aborts theXMLHttpRequestorfetchin the browser, but if the request already reached the server, the server still processes it. That is whyswitchMapis the wrong choice for “save” actions: you can end up with a write applied on the server whose response the client discarded.mergeMapin a search box is a race. WithmergeMap, a slow response for “ca” can arrive after the fast response for “cat” and overwrite the correct results. That is precisely the bugswitchMapexists to prevent.
mergeMap accepts a second argument limiting concurrency, for example mergeMap(upload, 3) to run at most three uploads at once. concatMap is equivalent to mergeMap(fn, 1).
Errors terminate the stream: where catchError goes
The Observable contract is that after error or complete, a stream emits nothing more. That is the rule behind the most common RxJS production bug: a search box that works until the first failed request and then silently stops responding.
// Broken: the first HTTP error kills the whole search stream
searchTerms$.pipe(
debounceTime(300),
switchMap((q) => http.get<Result[]>(`/api/search?q=${q}`)),
catchError(() => of([])), // replaces the dead stream with [] and completes
).subscribe(render);
When the inner request errors, the error propagates through switchMap to the outer stream. catchError then replaces the failed stream with of([]), which emits once and completes. The subscription to searchTerms$ is gone; further typing does nothing, and there is no error in the console because the error was “handled”.
The fix is to catch inside the projection, so only that one request fails:
searchTerms$.pipe(
debounceTime(300),
distinctUntilChanged(),
switchMap((q) =>
http.get<Result[]>(`/api/search?q=${encodeURIComponent(q)}`).pipe(
catchError((err) => {
reportError(err);
return of([] as Result[]);
}),
),
),
).subscribe(render);
A minimal check of the difference, with x === 2 throwing:
of(1, 2, 3).pipe(
map((x) => { if (x === 2) throw new Error('bad'); return x; }),
catchError(() => of(-1)),
).subscribe(console.log); // 1, -1 (3 never arrives)
of(1, 2, 3).pipe(
mergeMap((x) => of(x).pipe(
map((y) => { if (y === 2) throw new Error('bad'); return y; }),
catchError(() => EMPTY),
)),
).subscribe(console.log); // 1, 3
The other half of error handling is what happens when nobody handles it. If you subscribe without an error callback and the stream errors, RxJS 7 rethrows the error asynchronously, so it surfaces as an uncaught exception (in the browser console, or uncaughtException in Node) rather than being thrown from subscribe(). A try/catch around subscribe() does not catch it.
Retrying with backoff
retryWhen is deprecated in RxJS 7. retry now takes a configuration object:
import { retry, timer } from 'rxjs';
http.get('/api/data').pipe(
retry({
count: 3,
delay: (_err, attempt) => timer(2 ** attempt * 250), // 500ms, 1s, 2s
}),
);
retry({ count: 2 }) means the source is subscribed three times in total: the original attempt plus two retries. Retrying re-subscribes to the source, which for a cold HTTP observable means a new request. For a non-idempotent POST, think about whether a retry could apply the write twice.
Unsubscribing and memory leaks
Observables that complete on their own (of, a single HTTP response, take(1), firstValueFrom) release their resources when they complete. Observables that never complete (interval, fromEvent, Subjects, store selections, router events) have to be unsubscribed, or they hold a reference to your component and its callbacks for as long as the source lives.
takeUntil must come last
The classic pattern uses a notifier Subject:
private destroy$ = new Subject<void>();
ngOnInit() {
this.route.params.pipe(
switchMap((p) => this.api.poll(p.id)), // never completes
takeUntil(this.destroy$), // last operator
).subscribe((data) => (this.data = data));
}
ngOnDestroy() {
this.destroy$.next();
this.destroy$.complete();
}
The position matters. If takeUntil is placed before switchMap, it completes the outer source when destroy$ emits, but switchMap keeps its current inner subscription alive until that inner completes. With a polling inner observable, that is never. I confirmed this: with takeUntil before switchMap, an inner interval kept ticking after the notifier fired. The no-unsafe-takeuntil lint rule (in eslint-plugin-rxjs and its maintained fork eslint-plugin-rxjs-x) exists to catch this.
The failure I run into most in code reviews is not a missing unsubscribe in simple cases but this one: takeUntil was added, the component “has cleanup”, and nobody notices the polling continuing in the Network tab after navigating away because the operator is two lines too early.
Angular: takeUntilDestroyed, async pipe and toSignal
In Angular 16 and later, takeUntilDestroyed from @angular/core/rxjs-interop replaces the manual Subject:
import { Component, DestroyRef, inject } from '@angular/core';
import { takeUntilDestroyed } from '@angular/core/rxjs-interop';
export class UsersComponent {
private destroyRef = inject(DestroyRef);
constructor() {
// In an injection context (constructor, field initializer), no argument is needed
this.store.users$.pipe(takeUntilDestroyed()).subscribe(/* ... */);
}
ngOnInit() {
// Outside an injection context, pass the DestroyRef explicitly
this.route.params.pipe(takeUntilDestroyed(this.destroyRef)).subscribe(/* ... */);
}
}
Calling takeUntilDestroyed() with no argument inside ngOnInit fails, because ngOnInit does not run in an injection context. Even better is not subscribing manually at all: the async pipe and toSignal() subscribe and unsubscribe with the component’s lifetime. The Angular article covers signals and how they relate to observables.
RxJS 7 changes that break older examples
Much RxJS code online was written for version 6. The main differences you will hit:
-
toPromise()is deprecated (it still exists in 7.x but is scheduled for removal). UsefirstValueFrom(obs$), which resolves with the first value and unsubscribes, orlastValueFrom(obs$), which waits for completion. UnliketoPromise(), which resolved withundefinedfor an empty stream, both reject:EmptyError: no elements in sequencePass a default to opt out:
firstValueFrom(obs$, { defaultValue: null }). Be careful withlastValueFromon a stream that never completes (a Subject,interval): the promise never settles. -
Operators are exported from
'rxjs'directly.import { map, switchMap } from 'rxjs'works;'rxjs/operators'still works for compatibility. -
subscribe(next, error, complete)with three callbacks is deprecated. Pass an observer object:subscribe({ next, error, complete }). A singlenextfunction is still fine. -
throwError(error)should bethrowError(() => error), a factory, so the error is created per subscription. -
retryWhenandrepeatWhenare deprecated in favour ofretry({ delay })andrepeat({ delay }). -
shareis configurable:share({ resetOnError, resetOnComplete, resetOnRefCountZero }), which covers many cases that used to needpublish,multicastandrefCount, all of which are deprecated.
Subjects: which one and when
| Type | Late subscriber receives | Typical use |
|---|---|---|
Subject | Only future values | Event bus, notifier for takeUntil |
BehaviorSubject(initial) | The current value, then future values | State with an initial value; .value for synchronous reads |
ReplaySubject(n) | The last n values | Late listeners that need recent history |
AsyncSubject | Only the final value, on completion | Rare; similar to a Promise |
A Subject is a hot, imperative entry point into the reactive world. Expose it as an Observable (subject.asObservable()) so that consumers cannot call next() on it, and keep the ability to push values inside the service that owns the state. When every component can call next() on a shared Subject, reasoning about who changed the state becomes as hard as with global variables.
In Angular apps written today, a BehaviorSubject used purely to hold component or service state is often better replaced by a signal. RxJS remains the better tool for event streams: debouncing, cancellation, combining several asynchronous sources.
A WebSocket with reconnection, done correctly
A common example pipes webSocket() through a retry operator and then calls next() on the result. That does not work: pipe() returns a plain Observable, not the WebSocketSubject, so there is no next() to call. Keep the subject and the message stream separate:
import { webSocket } from 'rxjs/webSocket';
import { retry, timer } from 'rxjs';
const socket = webSocket<ServerMessage>('wss://example.com/ws');
const messages$ = socket.pipe(
retry({ delay: (_err, attempt) => timer(Math.min(1000 * attempt, 10000)) }),
shareReplay({ bufferSize: 1, refCount: true }),
);
messages$.subscribe({
next: (msg) => handle(msg),
error: (err) => console.error('socket failed', err),
});
socket.next({ type: 'subscribe', channel: 'prices' }); // sends through the subject
retry without count retries forever, with the delay function providing backoff. Because retry re-subscribes, the socket reconnects, but any subscription messages you sent before the drop are not replayed: send them again when the connection opens, for example from the openObserver config option. The WebSocket article covers the protocol side of reconnection.
When not to use RxJS
RxJS pays off for event streams with timing, cancellation and composition. For a single request that you await, a Promise with async/await is shorter and easier for the rest of the team to read. Wrapping every fetch in an Observable, or storing all application state in Subjects, adds operator knowledge as a requirement for touching any code, with no benefit when nothing is ever cancelled or combined. In Angular, HttpClient returns Observables anyway, and the practical rule is to use them for streams and to convert to a signal or a Promise at the edge where a single value is all you need.
Related Articles
- Angular for Enterprise Apps: Components, Services, DI, Signals
- JavaScript Async Programming: Promises, async/await and the Event Loop
- Real-Time Apps with WebSocket and Socket.io
- TypeScript Generics: How Inference Picks T
External references: RxJS documentation, RxJS 7 deprecations