skip to content

Pipeable Operators

The operators chained inside pipe() to transform, filter, rate-limit, flatten and join streams. Choosing between switchMap, mergeMap, concatMap and exhaustMap is a near-certain interview question.

part ofRxJSoverview, primer and where to startread it →
on this pageshow

explore

questions

25

In RxJS, what does distinctUntilChanged() drop from a sensor stream, and why can the same reading still appear twice?

level: juniorimportance: must knowfreq 64%

answer

  1. memory of one
  2. compares with the previous emission
  3. === by default
  4. distinct() remembers everything

basics

~20 s

distinctUntilChanged() drops a value only when it equals the last value it emitted, using === by default. It remembers one key, so 21, 21, 22, 21 becomes 21, 22, 21; distinct() is the operator that suppresses every earlier repeat.

solid answer

~40 s

`distinctUntilChanged()` keeps exactly one piece of state: the key of the last value it let through. Each new value is compared with that key (`===` unless you pass a comparator); if they are equal the value is dropped, otherwise it is emitted and becomes the new reference. The first value is always emitted. Because the memory is only one value deep, a reading that returns after a different one passes again: `21, 21, 22, 21` emits `21, 22, 21`. That is exactly what you want for a sensor feed or a form value, where only a *change* matters. If you need to suppress any value ever seen, that is `distinct()`, which keeps a `Set` of every key and grows without bound on an endless stream unless you give it a `flushes` notifier.

code

ts · 10 lines
ts
import { interval, map, distinctUntilChanged } from 'rxjs';

const readSensor = () => Math.round(20 + Math.random());

const temperature$ = interval(1000).pipe(
  map(() => readSensor()),
  distinctUntilChanged()
);

temperature$.subscribe(c => console.log('temperature changed to', c));

go deeper

for a junior

Recall that only consecutive duplicates are dropped, that the first value always passes, and give the 21, 21, 22, 21 example with its output.

for a middle

Explain the single previousKey, the default === comparison, why objects defeat it, and when distinct() with its growing Set is the right tool instead.

for a senior

Show where it belongs in a real pipeline: after reducing to a primitive, before costly work, and why distinct() on an endless stream is a memory leak without flushes.

for a principal

Frame it as a contract decision: which layer guarantees change-only emission, and whether producers should emit primitives or stable references so consumers stay cheap.

## What the operator remembers `distinctUntilChanged` is a pipeable operator from the `rxjs` package (importable from `'rxjs'` since 7.2; the `'rxjs/operators'` path still works but is deprecated). It answers one question for every incoming value: **is this the same as the value I emitted last?** In RxJS 7.8 the implementation keeps two variables per subscription: - a `first` flag, so the very first value is **always emitted** — there is nothing to compare it with; - `previousKey`, the key of the **last emitted** value (by default the value itself). For each later value it calls the comparator — `(a, b) => a === b` unless you supply one — with `previousKey` and the new key. If the comparator returns `true` the values are considered equal and the new one is **dropped**; if it returns `false` the value is **emitted** and becomes the new `previousKey`. Errors and completion pass straight through, and nothing is emitted at completion. Two implementation details are worth knowing: - the state lives **per subscription**, so two subscribers to the same cold pipeline each get their own memory; - the key is updated **before** the value is emitted, so re-entrant code that pushes a value back into the source during emission is compared against the right reference (fixed in the RxJS 7.0 beta cycle). ## Walking a sensor stream Imagine a thermometer that reports every second, even when nothing changed. Downstream, a chart or an alerting rule only cares when the reading moves. | incoming | previousKey | equal? | emitted | |---|---|---|---| | 21 | (none) | first value | 21 | | 21 | 21 | yes | — | | 21 | 21 | yes | — | | 22 | 21 | no | 22 | | 22 | 22 | yes | — | | 21 | 22 | no | 21 | The final `21` is emitted even though `21` was seen before: the operator only compares **adjacent** emissions. That is the correct behaviour for a sensor — the temperature really did change back. ## distinctUntilChanged versus distinct | | `distinctUntilChanged()` | `distinct()` | |---|---|---| | compares against | the last emitted key only | every key ever emitted | | state | one value | a `Set` that keeps growing | | `21, 22, 21` | `21, 22, 21` | `21, 22` | | safe on an endless stream | yes | only with a `flushes` notifier that clears the set | | typical use | sensors, form values, store selections | de-duplicating ids in a finite batch | RxJS documents that `distinct` keeps its keys in a `Set` and offers an optional `flushes` observable to clear it; without one, a long-lived stream of unique values leaks memory. ## Where it fits in a pipeline 1. **Reduce to a primitive first** when you can. `map(r => r.celsius)` followed by `distinctUntilChanged()` compares numbers, which `===` handles correctly. 2. **Place it before the expensive work** — a re-render, an HTTP call, a recalculation — so repeated values never reach it. 3. **Keep it after operators that produce new values**, not before: filtering duplicates and then mapping to a rounded value can re-introduce adjacent duplicates. For objects, `===` compares **references**, so a producer that builds a fresh object per reading defeats the default comparison; the fix is a comparator, a key selector or `distinctUntilKeyChanged`, which is a separate concern. One more placement rule: when several subscribers share a sensor feed, de-duplicate **once**, upstream of the sharing point, rather than repeating the operator in every consumer. Each copy of the operator keeps its own memory, so duplicated operators cost memory and CPU without changing the result. ## Common misreadings - "It removes all duplicates" — it removes consecutive ones only. - "It swallows the first value" — the first value is always emitted. - "It waits for the stream to settle" — it has no notion of time; it decides synchronously per value. Waiting for a quiet period is `debounceTime`, a different operator. - "It emits the last value on completion" — completion is forwarded as-is. ```ts import { of, distinctUntilChanged, distinct } from 'rxjs'; const celsius$ = of(21, 21, 21, 22, 22, 21); celsius$.pipe(distinctUntilChanged()).subscribe(v => console.log('changed', v)); // changed 21, changed 22, changed 21 celsius$.pipe(distinct()).subscribe(v => console.log('unique', v)); // unique 21, unique 22 ``` In an Angular app the same operator commonly sits on a form control's `valueChanges` or on a selection from a service's state stream, so that a subscriber only reacts when the value actually differs from the last one it saw.

  • What does distinctUntilChanged() do when the source errors or completes?
    It forwards both notifications unchanged. It holds no buffered value, so nothing extra is emitted at completion; the subscriber simply receives `complete` or `error` right after the last value that passed.
  • Why prefer distinctUntilChanged() over distinct() on a long-lived sensor stream?
    `distinct()` stores every key it has emitted in a `Set`, so an endless stream of changing readings grows memory without bound unless a `flushes` observable clears it. It would also hide a legitimate return to an earlier value. `distinctUntilChanged()` keeps one key and reports every real change.

A departures board that only repaints a row when the new status differs from what it currently shows: "Delayed, Delayed, Boarding, Delayed" repaints three times, because it compares only with what is on the board right now.

saying these in an interview costs you the question

  • distinctUntilChanged() removes every duplicate the stream has ever produced.
  • The first value is dropped because there is nothing to compare it with.
  • It compares objects by their contents out of the box.
  • It waits for the stream to go quiet before emitting, like debounceTime.
  • distinct() is safe on endless streams because it forgets old keys by itself.
open as a page

In RxJS, why should a typeahead search use switchMap rather than mergeMap to call the search API for each query?

level: juniorimportance: must knowfreq 76%

basics

~20 s

Search responses can return out of order. mergeMap forwards every response, so a slow reply for an old query can overwrite the newest results. switchMap unsubscribes the previous request when a new query arrives, so only the latest results arrive.

open as a page

In RxJS, how do debounceTime and throttleTime differ, and which fits autosave while typing versus a scroll-position tracker?

level: juniorimportance: must knowfreq 72%

basics

~20 s

debounceTime emits the latest value once the stream has been silent for the given time, which suits autosave after typing pauses. throttleTime emits a value, then ignores the source for a fixed window, which suits steady scroll updates.

open as a page

In RxJS, what is the difference between the map and tap operators, and why do side effects belong in tap?

level: juniorimportance: must knowfreq 70%

basics

~20 s

map replaces each value with its projection's result; tap runs a callback for each notification and passes the original value through, ignoring the callback's return. Side effects in tap keep map pure and visible in the pipe.

open as a page

In RxJS, how do combineLatest and forkJoin differ, and why can a forkJoin never emit when one source never completes?

level: middleimportance: must knowfreq 76%

basics

~20 s

combineLatest emits the latest value of every source once each has emitted, then again on every new value; forkJoin waits for every source to complete and emits their last values once, so a never-completing source blocks it forever.

open as a page

In RxJS, why does first() error with EmptyError when the source completes without a value, while take(1) simply completes?

level: middleimportance: must knowfreq 60%

basics

~20 s

first() promises one value: filter plus take(1) plus a check that errors with EmptyError if the source completes empty. take(1) promises at most one and just completes. first(pred, fallback) emits the fallback instead of erroring.

open as a page

In RxJS, how do switchMap, mergeMap, concatMap and exhaustMap differ when a new value arrives while an inner observable is still running?

level: middleimportance: must knowfreq 82%

basics

~20 s

switchMap unsubscribes the running inner observable and switches to the new one; mergeMap runs both at once; concatMap queues the new value until the current inner completes; exhaustMap ignores the new value until the current inner completes.

open as a page

In RxJS, what is the difference between scan and reduce, and which would you use for a live cart total?

level: middleimportance: must knowfreq 62%

basics

~20 s

scan emits the updated accumulator after every value; reduce emits only the final accumulator when the source completes. A live cart total needs scan, because the cart stream never completes and reduce would never emit.

open as a page

In RxJS, what is the difference between merge() and concat() when you combine two observables into one stream?

level: juniorimportance: should knowfreq 58%

basics

~10 s

merge subscribes to every source at once and forwards values as they arrive, interleaved; concat subscribes to the next source only after the previous one completes, so values come out source by source.

open as a page

In RxJS, given of(1, 2, 3, 4, 5, 1), what do take(3), takeWhile, skip(2) and skipWhile each emit with the predicate x < 3?

level: juniorimportance: should knowfreq 50%

basics

~20 s

take(3) emits 1, 2, 3 and completes; takeWhile(x => x < 3) emits 1, 2 and completes at the 3; skip(2) and skipWhile(x => x < 3) both emit 3, 4, 5, 1, since skipWhile stops checking once its predicate fails.

open as a page

In RxJS, when would you use withLatestFrom instead of combineLatest, and why can withLatestFrom silently drop source values?

level: middleimportance: should knowfreq 48%

basics

~10 s

Use withLatestFrom when only the source stream should trigger output and other streams just supply their latest value; source values that arrive before every other stream has emitted once are discarded, not queued.

open as a page

In RxJS, why does distinctUntilChanged() emit every sensor-reading object, and how do a comparator or distinctUntilKeyChanged() fix it?

level: middleimportance: should knowfreq 47%

basics

~20 s

Its default === compares object references, and each reading is a new object, so none look equal. Pass a comparator returning true for equal readings, a key selector as the second argument, or use distinctUntilKeyChanged('celsius') for one top-level property.

open as a page

In RxJS, why does exhaustMap suit a login button and a slow polling loop, and what exactly happens to the values it ignores?

level: middleimportance: should knowfreq 55%

basics

~20 s

exhaustMap starts an inner observable only when none is active and discards source values that arrive meanwhile, without buffering them. A login button then sends one request however often it is clicked, and a poll never overlaps a slow previous request.

open as a page

In RxJS, how do auditTime and sampleTime differ from throttleTime, and which one reports a scroll tracker's latest position?

level: middleimportance: should knowfreq 36%

basics

~20 s

auditTime starts a window on a value and emits the latest value when it ends; sampleTime emits the latest value on a fixed clock; default throttleTime emits a window's first value. auditTime usually suits a scroll tracker's latest position.

open as a page

In RxJS, why does a scroll-position tracker using throttleTime(100) miss the final position, and what do its leading and trailing options change?

level: middleimportance: should knowfreq 44%

basics

~20 s

throttleTime defaults to { leading: true, trailing: false }: it emits each window's first value and drops the rest, so a scroll ending mid-window loses its final position. { leading: true, trailing: true }, passed third, also emits each window's latest value.

open as a page

In RxJS, how do bufferTime, bufferCount and buffer(notifier) decide when to emit a batch, and what happens to leftover values?

level: middleimportance: should knowfreq 38%

basics

~20 s

bufferTime emits the collected array every time span, bufferCount after N values, and buffer each time its notifier emits. On completion the partial buffer is emitted; on error it is discarded. bufferTime and buffer emit empty arrays when nothing arrived.

open as a page

An RxJS dashboard stream built as combineLatest([user$, settings$, permissions$]) never produces a value; how do you find the cause and fix it?

level: seniorimportance: should knowfreq 44%

basics

~20 s

combineLatest stays silent until every source has emitted once, so one of the three never emitted, completed empty or errored; trace each source with tap, then seed or default the culprit or fix its producer.

open as a page

In RxJS, a sensor feed piped through takeUntil(stop$) keeps emitting after stop$.complete() is called — why, and what other notifiers or placements make takeUntil misfire?

level: seniorimportance: should knowfreq 42%

basics

~20 s

takeUntil completes only when its notifier emits a value; a notifier that completes without emitting is ignored. It also misfires when the notifier errors, emits synchronously at subscribe time, or when a flattening operator after it keeps an inner stream alive.

open as a page

In RxJS, what risk does concatMap carry when saves arrive faster than the server answers, and how does mergeMap's concurrent argument change it?

level: seniorimportance: should knowfreq 45%

basics

~20 s

concatMap runs one save at a time and buffers the rest in an unbounded queue, so a slow server grows the backlog, memory and latency. mergeMap(save, n) runs up to n saves at once, draining faster but losing ordered completion.

open as a page

In RxJS, why can an autosave built on debounceTime(1000) lose the user's last edit when the editor closes, and how do you flush it?

level: seniorimportance: should knowfreq 34%

basics

~10 s

debounceTime holds the latest edit until its quiet period ends; unsubscribing first discards it. Completing the source instead makes debounceTime emit the held value immediately, so end the stream upstream rather than unsubscribing.

open as a page

In RxJS 7, how do you write a reusable custom operator, both by composing existing operators with pipe() and from scratch?

level: seniorimportance: should knowfreq 35%

basics

~20 s

A custom operator is a function from a source Observable to a new one. Compose existing operators with the standalone pipe(); only when that cannot work, return a new Observable that forwards next, error and complete and unsubscribes the source in teardown.

open as a page

In RxJS 7, why were the pipeable merge, concat, zip, race and combineLatest operators deprecated, and what replaces them?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

Each shared its name with a static creation function, which clashed once all operators moved into the 'rxjs' entry point; they were renamed mergeWith, concatWith, zipWith, raceWith and combineLatestWith, and the old names are deprecated.

open as a page

In RxJS, what do mergeAll, concatAll, switchAll and exhaustAll do, and how do they relate to the matching *Map operators?

level: middleimportance: nice to knowfreq 26%

basics

~20 s

They flatten an observable that already emits observables, using the same strategies as the Map operators: run concurrently, queue, switch to the newest, or ignore while busy. map(project) followed by mergeAll() behaves like mergeMap(project); concatAll() is mergeAll(1).

open as a page

In RxJS, why does pairwise() emit nothing for the first value, and how does startWith() change what it emits?

level: middleimportance: nice to knowfreq 28%

basics

~10 s

pairwise emits [previous, current] tuples, so it needs two values before its first emission. Putting startWith(x) before it supplies a synthetic first value, so the first real value produces [x, first] immediately.

open as a page

In RxJS, how does debounce() with a duration selector differ from debounceTime, and why does returning EMPTY from the selector hold the value?

level: seniorimportance: nice to knowfreq 24%

basics

~20 s

debounce() gets a duration observable per value from its selector and emits the value when that observable first emits. Since RxJS 7 completing does not end the duration, so EMPTY holds the value until a newer one or source completion.

open as a page