Thirty RxJS 7 operators (creation functions included) grouped by what they do, with the details that come up in Angular interviews.
Creation (6)
of(1, 2, 3): emits its arguments synchronously, then completes.from(input): turns an array, promise, iterable or other observable-like into an Observable:from(fetch(url)).fromEvent(el, 'click'): DOM events as a stream. It never completes on its own.interval(1000): emits0, 1, 2, ...every second; the first value arrives after one second.timer(500): emits0once after 500 ms, then completes.timer(0, 1000)emits after 0 ms, then every second.throwError(() => new Error('Nope')): errors immediately. Pass a factory; passing the error value directly is deprecated.
Combination (5)
combineLatest([a$, b$]): once every source has emitted, emits an array of the latest values whenever any source emits. An object works too:combineLatest({ user: user$, prefs: prefs$ }).forkJoin([a$, b$]): waits for every source to complete, then emits their last values once, likePromise.all. A source that completes without emitting makes it complete empty; any error errors it.merge(a$, b$): interleaves values from all sources as they arrive; completes when all of them complete.withLatestFrom(other$): on each source value, emits[value, latestOther].other$never triggers an emission, and source values are dropped until it has emitted once.startWith(initial): emitsinitialsynchronously before anything from the source.
Transformation (2)
map((x) => x * 2): transforms each value.scan((acc, x) => acc + x, 0): likereduce, but emits every intermediate result, which makes it a tiny state machine.
Flattening (4)
What happens when a new outer value arrives while an inner Observable is still running (deep dive):
switchMap: unsubscribes from the previous inner (cancels it). Typeahead, route params: latest wins.mergeMap: runs inners concurrently, optionally capped:mergeMap(fn, 3). Independent parallel work.concatMap: queues them and runs one inner at a time, in order. Sequential saves.exhaustMap: ignores new outer values until the current inner completes. Submit and login buttons.
Filtering & rate limiting (7)
filter((x) => x > 0): passes only matching values.take(3): emits the first 3 values, then completes and unsubscribes from the source.takeUntil(destroy$): completes whendestroy$emits a value (completing without one does not count). Keep it last in thepipe.first(): the first value, or first match of a predicate, then completes. Errors withEmptyErrorif the source completes first, unless you pass a default:first((x) => x > 0, 0).distinctUntilChanged(): drops a value that is===to the previous one. For objects, pass a comparator:distinctUntilChanged((a, b) => a.id === b.id).debounceTime(300): emits the latest value after 300 ms of silence. If the source completes during the wait, the pending value is emitted first.throttleTime(1000): emits a value, then ignores the rest for 1 s. For the trailing value too:throttleTime(1000, asyncScheduler, { leading: true, trailing: true }).
Error handling (2)
catchError((err) => of(fallback)): replaces the failed stream with the Observable you return, or rethrow withthrowError(() => err). The source is finished at that point, so to keep an outer stream alive, catch inside the innerpipeof aswitchMap.retry({ count: 3, delay: 1000 }): resubscribes on error.retry(3)retries immediately;delaycan also be a function returning a notifier, such as(err, n) => timer(n * 1000)for backoff.
Utility (3)
tap(fn): side effects (logging, analytics) without changing values. Also accepts{ next, error, complete, subscribe, unsubscribe, finalize }.finalize(() => spinner.hide()): runs once the stream completes, errors or is unsubscribed.delay(500): shifts each value later by 500 ms (or until a givenDate).
Multicasting (1)
shareReplay({ bufferSize: 1, refCount: true }): shares one subscription among all subscribers and replays the last value to late ones. The default isrefCount: false(which is whatshareReplay(1)gives you): the source stays subscribed even after everyone unsubscribes, a leak with never-ending sources.
Recipes
const results$ = query$.pipe(
debounceTime(300),
map((q) => q.trim()),
distinctUntilChanged(),
filter((q) => q.length >= 2),
switchMap((q) => api.search(q).pipe(catchError(() => of([])))),
);saveClicks$.pipe(exhaustMap(() => api.save(form.value))); // no double submits
config$ = http.get<Config>('/api/config').pipe(shareReplay(1)); // one request, cachedshareReplay(1) without refCount is fine for an HTTP call, because the source completes after one response.
Angular tips
- The
asyncpipe subscribes and unsubscribes for you. takeUntilDestroyed()from@angular/core/rxjs-interop(Angular’s, not RxJS’s) completes a stream when the component is destroyed. Call it in an injection context or pass aDestroyRef.toSignal(obs$)andtoObservable(sig)bridge Observables and signals.HttpClientrequests complete after one response, but unsubscribing still cancels one in flight, which is exactly whatswitchMaprelies on.- Operators are pure functions:
source$.pipe(...)returns a new Observable, and nothing runs until something subscribes.