skip to content

RxJS

2 roadmaps56 questionsupdated

RxJS is the observable library behind Angular's async APIs: creation functions, pipeable operators, subjects and schedulers. Interviewers probe operator choice, teardown and hot versus cold.

on this pageshow

guide

overview

~1 min

RxJS is a library for composing asynchronous and event-driven code as streams of values over time. Most frontend developers meet it through Angular, whose HTTP client, router and forms hand out observables, but it fits any JavaScript code that coordinates input, timers, sockets and requests. Interviewers use it to test whether you can reason about time: what starts work, what is still running after a view closes, and which of several near-identical operators encodes the behaviour you actually want. The hub follows the way code uses the library. [Observables and observers](/topics/fe-rxjs-observables) is the contract everything else builds on, and [stream creation functions](/topics/fe-rxjs-creation) turn values, promises, events and time into streams. [Pipeable operators](/topics/fe-rxjs-pipeable-operators) is the largest section: mapping, filtering, time-based throttling, flattening and joining several sources. [Subjects and multicasting](/topics/fe-rxjs-subjects) covers sharing one execution among many subscribers. [Catching errors and retrying](/topics/fe-rxjs-error-retry) and [unsubscribing and schedulers](/topics/fe-rxjs-schedulers) decide how a stream fails, when it stops and on which clock it runs. [Marble diagram tests](/topics/fe-rxjs-marble-testing) are how you prove any of that without waiting. Junior rounds stay close to the model: laziness, the comparison with promises, the difference between two similar operators. Middle and senior rounds become bug hunts: a leak that survives teardown, a stale search result, a shared stream that never releases its socket, a test that hangs. Learn the observable contract and subscription lifetime first, then the everyday operators, then flattening.

primer

### An observable is a recipe, not a result A plain observable does nothing until something subscribes, and every new subscription runs the recipe again from the start. That one fact explains the comparison with promises, the request that fires twice, and why hot versus cold, sharing and retrying behave as they do. ### Many values, then one ending An observer can receive any number of values, followed by at most one terminal signal: error or completion. After that the stream is finished for good. Many traps come from forgetting whether a source ever ends: an accumulator that stays silent, a join that waits forever, an error handler expected to resume. ### Every subscription owns resources Subscribing wires up listeners, timers, sockets or requests, and the subscription is the handle that releases them. Finite sources clean up on their own; endless ones hold on until something ends them. Leak questions come down to telling the two apart and choosing a declarative ending over manual bookkeeping. ### Operators are functions from stream to stream `pipe()` chains functions that each subscribe to what is upstream and return a new observable. Position in the chain is meaning, not style: the same operator moved two lines can fix a leak or cause one. ### Flattening is a concurrency policy When each value starts an inner stream — a request per keystroke, a save per click — the flattening operator decides what happens on overlap: cancel the old work, run both, queue the new, or ignore it. A good answer names the failure each choice brings under load. ### Time is an input Rate-limiting operators, timers and schedulers treat time as data. Whether a value leaves at the start of a window, at its end or after a quiet period separates operators that look alike, and the same abstraction lets tests swap in virtual time. ### Sharing is opt-in By default each subscriber gets its own execution. Subjects and the share operators let one execution feed many subscribers, and the questions then become what a late subscriber receives and when the shared source is allowed to stop.

Observable
A lazy description of a stream. Subscribing runs it and delivers values, then optionally an error or completion, to one observer.
Subscription
The handle returned by subscribe(). Calling unsubscribe() on it stops delivery and releases whatever the stream set up.
Teardown
The cleanup logic a stream registers when it starts, such as removing a listener or clearing a timer, run once when the subscription ends.
Cold observable
A stream whose producer is created per subscription, so each subscriber gets its own independent execution from the beginning.
Hot observable
A stream whose producer exists independently of any one subscriber, so subscribers share it and only see what happens after they join.
Pipeable operator
A function that takes an observable and returns a new one, chained inside pipe() to transform, filter, time or combine values.
Higher-order observable
An observable whose values are themselves observables, usually produced by mapping each value to a request or other inner stream.
Flattening
Subscribing to the inner streams of a higher-order observable and merging their values into one output, under a chosen overlap policy.
Subject
An object that is both observable and observer: values pushed into it are broadcast to every current subscriber.
BehaviorSubject
A subject that always holds a current value, starts with an initial one, and hands it to each new subscriber at once.
ReplaySubject
A subject that buffers recent values, limited by count and age, and replays them to each new subscriber before live values.
Scheduler
An object that decides when and in which execution context work runs: synchronously, as a microtask, on a timer or per animation frame.
Virtual time
A simulated clock used by the test scheduler, so time-based operators run instantly and deterministically in tests.

Follow one feature through the layers. A creation function turns something outside the library — an event, a timer, a promise — into an observable. Operators in `pipe()` reshape it, and where one value should start more work, a flattening operator subscribes to an inner stream and applies its overlap policy. Error handling sits where a failure should stop. A share operator decides whether each subscriber pays for its own execution. Finally something subscribes — a component, a template, a test — and something decides when that subscription ends; in Angular templates the async pipe does both. A polling price feed shows several of those decisions meeting in one chain: ```typescript const prices$ = timer(0, 30_000).pipe( exhaustMap(() => fetchPrices().pipe( retry({ count: 2, delay: 1_000 }), catchError(() => of(null)), // this poll fails, the feed lives on ), ), shareReplay({ bufferSize: 1, refCount: true }), ); ``` Each choice is one section of the hub. The timer comes from creation functions. The flattening operator keeps a slow response from overlapping the next tick. Retry and recovery live inside the inner stream, because an error that reached the timer would end the feed for everyone. The share at the end gives all subscribers one poll and one latest value, and the reference-counting option lets the timer stop once the last subscriber leaves. Marble tests verify all of it without waiting thirty seconds. Schedulers run underneath all of it. Most code never names one, but swapping the default — for animation frames, or virtual time in a test — changes when values arrive, not what the chain means.

  1. Observables and Observers →

    The contract every other section assumes: laziness, the three notifications, teardown and the comparison with promises.

  2. Pipeable Operators →

    How chaining works and the everyday operators for mapping, filtering and timing that most interview code uses.

  3. Higher-Order Flattening →

    The four overlap policies behind requests, saves and searches; the most frequent middle-level operator question.

  4. Unsubscribing & Schedulers →

    When a subscription ends, how leaks happen and how to stop them, then which clock work runs on.

  5. Subjects and Multicasting →

    Hot versus cold in practice: pushing values yourself, late subscribers and sharing one execution safely.

  6. Catching Errors & Retrying →

    Recovery after errors terminate a stream: replacement streams, retries with backoff, timeouts and cleanup.

  • Subscribing inside another subscribe callback instead of flattening: the inner subscriptions escape teardown and cancellation, and the overlap policy is never stated.

  • Reaching for mergeMap by default when user input triggers requests; name the overlap behaviour first — see Higher-Order Flattening.

  • Putting takeUntil early in the chain and assuming everything after it stops — see takeUntil placement.

  • Expecting catchError to resume the stream it caught: that source is finished, so a long-lived stream needs its recovery inside the inner observable.

  • Adding shareReplay to an endless source without deciding what happens when subscribers leave — see the refCount leak.

  • Treating from(promise) as lazy: the promise has already started, so repeat subscriptions share one result unless a factory runs per subscription.

  • Testing debounce or polling with real waits or ad-hoc fake timers, then chasing flaky tests — see Marble Diagram Tests.

This guide assumes RxJS 7. A lot of production code, tutorials and interviewers still carry habits from RxJS 6, so several questions turn on what 7 changed: - **Promise conversion.** `toPromise()` is deprecated in favour of `firstValueFrom` and `lastValueFrom`, which state which value they wait for. - **Multicasting.** The `multicast`, `publish` and `refCount` family is deprecated; `connectable`, `connect`, `share` with a configuration object and `shareReplay` cover the same ground. - **Operator names.** Pipeable joins that shared a name with a creation function gave way to `*With` variants such as `mergeWith`, and during the 7.x line operators became importable from the main `rxjs` entry point, not only `rxjs/operators`. - **Retry.** `retry` gained a configuration object, and `retryWhen` is deprecated in favour of it. Further back, RxJS 6 made pipeable operators the standard, replacing operators patched onto the observable prototype; chained `.map().filter()` code is from that older era. When an answer depends on version, say which one you mean.

RxJS is the JavaScript member of the ReactiveX family, so its vocabulary transfers to RxJava and its relatives. In the browser it competes with simpler tools, and interviewers expect you to say when it is worth its weight. - **Promises and async/await** fit a single future value; observables earn their cost with many values over time, cancellation and composition across time. - **Signals**, including Angular's, hold synchronous current state and track dependencies automatically; observables model events and asynchronous flows. Angular ships interop helpers between the two, and a common modern answer uses signals for view state and RxJS for event streams, requests and timing. - **State libraries** in the Angular world lean on it too: NgRx's classic store exposes observables and its side effects are written as streams, so operator fluency carries over. A defensible choice names the workload. Debounced search, polling, cancellable requests and joins of several live sources favour RxJS; a value read once or a flag the template displays often does not need it.

explore

report an issue with this guide →

questions

page 2 of 2

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

In RxJS, how do asyncScheduler, asapScheduler and queueScheduler differ in when they run a task scheduled with zero delay?

level: middleimportance: should knowfreq 35%

basics

~20 s

With zero delay, RxJS queueScheduler runs the task synchronously (queuing nested tasks until the current one ends), asapScheduler runs it as a microtask after the current code, and asyncScheduler runs it on a timer, after pending microtasks.

open as a page

In RxJS, what is the difference between observeOn and subscribeOn, and which one changes when values reach the observer?

level: middleimportance: should knowfreq 28%

basics

~20 s

RxJS subscribeOn schedules the moment the source is subscribed, so a synchronous source still emits in one burst later. observeOn reschedules every next, error and complete notification on the scheduler, so it changes when values reach downstream observers.

open as a page

In RxJS, how do you end several subscriptions at once with a parent Subscription, and how does Subscription.add() behave?

level: middleimportance: should knowfreq 45%

basics

~20 s

Create one RxJS Subscription, add() each child subscription or teardown function to it, and call unsubscribe() once. add() returns void in RxJS 7, runs a teardown immediately if the parent is already closed, and closed children drop out automatically.

open as a page

In RxJS, a chat room's messages$ must show late joiners recent history; how do ReplaySubject's bufferSize and windowTime decide what they get?

level: middleimportance: should knowfreq 50%

basics

~10 s

ReplaySubject(bufferSize, windowTime) keeps at most bufferSize values that are younger than windowTime milliseconds and replays them in order to each new subscriber before live values; both limits default to infinity.

open as a page

In RxJS, which creation functions emit synchronously inside subscribe(), and what bugs can that synchronous emission cause in real code?

level: seniorimportance: should knowfreq 38%

basics

~20 s

of, from over arrays or iterables, range and EMPTY deliver all values and completion before subscribe() returns; from(promise), interval and timer deliver later. Synchronous delivery breaks code that uses the subscription variable in its callback or that assumes asynchrony.

open as a page

In RxJS 7, how would you retry a flaky price endpoint with exponential backoff using retry's config, and what do count, delay and resetOnSuccess control?

level: seniorimportance: should knowfreq 62%

basics

~10 s

Use retry({ count, delay }) where delay is a function of (error, retryCount) returning timer(base * 2 ** (retryCount - 1)); count caps retries, and resetOnSuccess restarts the count after the source emits again.

open as a page

In an RxJS marble test, how do you use expectSubscriptions and subscription marbles to prove that switchMap unsubscribes a stale inner stream?

level: seniorimportance: should knowfreq 25%

basics

~20 s

In RxJS marble tests, hot() and cold() fixtures log every subscription. Pass fixture.subscriptions to expectSubscriptions(...).toBe(...) with subscription marbles, where ^ marks subscribe and ! unsubscribe; an array covers several subscriptions, proving the stale inner ended when the next began.

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, why can shareReplay(50) keep a chat room's socket stream open after every subscriber has left, and how do you fix it?

level: seniorimportance: should knowfreq 60%

basics

~20 s

shareReplay's refCount defaults to false, so when the subscriber count drops to zero its inner ReplaySubject stays subscribed to the source; for a never-completing source use shareReplay({ bufferSize, refCount: true }) or share with explicit reset options.

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 7, how do timeout's each, first and with options behave, and how do you combine timeout with retry for a price request that hangs?

level: seniorimportance: nice to knowfreq 34%

basics

~20 s

timeout errors with a TimeoutError when a value is late: each limits every gap, first the wait for the first value, and with switches to a replacement instead; put timeout before retry so each attempt gets its own timer.

open as a page

Why can an RxJS TestScheduler.run marble test hang forever or see no values at all, and how do you fix each case?

level: seniorimportance: nice to knowfreq 20%

basics

~20 s

RxJS TestScheduler.run removes the frame limit, so an endless source such as interval keeps virtual time flushing forever; add an unsubscription marble. Promises and direct timers are not virtualised, so their values never arrive; test those parts asynchronously or stub them.

open as a page

An RxJS subscribe() call passes only a next callback and the source errors - what happens, and why does a surrounding try/catch not catch it?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

RxJS 7 treats the error as unhandled: the subscription closes and runs its teardown, and the error is rethrown asynchronously in a setTimeout, or passed to config.onUnhandledError. A try/catch around subscribe() has already exited, so it never sees it.

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

In RxJS, how would you drive a smooth countdown animation with animationFrameScheduler, and why does interval(16, animationFrameScheduler) not stay in step with frames?

level: seniorimportance: nice to knowfreq 18%

basics

~20 s

Use interval(0, animationFrameScheduler) or animationFrames() to emit once per frame, compute remaining time from a clock, and end with takeWhile. A positive delay makes animationFrameScheduler fall back to a timer, so interval(16, ...) is not frame-aligned.

open as a page

In RxJS 7, what replaced the deprecated multicast, publish and refCount operators, and when would you use connectable() instead of share()?

level: seniorimportance: nice to knowfreq 24%

basics

~10 s

RxJS 7 reduced multicasting to connectable, connect, share and shareReplay; use connectable() when you must attach every subscriber before the source starts, because share() starts it on the first subscription.

open as a page

showing 31–56 of 56