Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

RxJS lets JavaScript developers represent events and asynchronous results as Observable sequences, transform them with operators, and connect a consumer with subscribe(). That makes flows such as clicks and search requests easier to compose—but RxJS is a practical reactive library, not a single definitive implementation of every formal meaning of functional reactive programming (FRP).

What reactive programming with RxJS means

In an ordinary function call, you typically ask for a value and receive it immediately. With an event stream, values may arrive later and repeatedly: a user clicks, a timer fires, or a request returns. RxJS represents such sequences with Observables and provides functions for composing them.

The official RxJS overview describes ReactiveX this way: “ReactiveX combines the Observer pattern with the Iterator pattern and functional programming with collections to fill the need for an ideal way of managing sequences of events.” In practice, RxJS brings together the pattern of producing notifications, the pattern of consuming a sequence, and functional transformations that can be chained.

A useful way to think about a flow is: create or adapt a source, compose transformations, then subscribe to observe the result. Defining that flow and executing its producer are conceptually distinct; for many Observables, the producer work begins when someone subscribes.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

The three parts of a basic RxJS flow

Observable: the source or sequence description

An Observable models values or events that may arrive over time. It can be created directly or adapted from an existing source. For example, fromEvent adapts DOM events into an Observable sequence:

import { fromEvent } from 'rxjs';
import { map } from 'rxjs/operators';

const clicks = fromEvent(document, 'click').pipe(
  map(event => ({ x: event.clientX, y: event.clientY }))
);

This creates a description of a click flow and a transformation; it does not yet log anything. Each click can be mapped to a simpler object containing its coordinates.

Operators and pipe(): transform or coordinate values

Operators are composable functions that transform or coordinate observable sequences. pipe() lays out the sequence of transformations, making the intended flow visible without mixing every step into a callback. Here, map turns each event into coordinates. Filtering, timing, and combining operators can build on the same pattern.

Observer and subscribe(): consume notifications

An Observer handles notifications from an Observable: ordinary values, an error, or completion. The consumer attaches by calling subscribe(). This is where observation is established and, for many sources, where producer work starts.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
const subscription = clicks.subscribe({
  next: point => console.log(point),
  error: error => console.error(error),
  complete: () => console.log('Click stream complete')
});

// When this consumer no longer needs the flow:
subscription.unsubscribe();

The returned Subscription represents that consumer’s connection and can be unsubscribed to clean up. For a DOM event source, cleanup prevents the consumer from continuing to receive clicks. Unsubscription is not the same as completion: it ends this subscription rather than sending the Observer a normal completion notification.

Subscription, cold sources, and multiple consumers

Learn RxJS describes Observables as cold and unicast by default: subscribers generally get their own execution rather than automatically sharing one producer run. A second subscription to a cold source can therefore start another execution. This matters for sources with side effects, such as making a request or registering event handlers.

If consumers should share work or receive earlier values, choose a multicasting approach deliberately. A Subject can act as a shared source that multicasts notifications; sharing operators can share an Observable execution. Those approaches differ in whether late subscribers receive previous values and how the shared connection’s lifetime relates to its subscribers. Do not assume that adding another subscriber is always just another view of the same work.

Build an asynchronous search flow

A typeahead search needs to avoid sending a request for every keystroke. A common RxJS composition waits for a pause in typing, ignores unchanged queries, and switches to a request for the latest query:

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
import { fromEvent } from 'rxjs';
import { debounceTime, distinctUntilChanged, map, switchMap } from 'rxjs/operators';

const results = fromEvent(searchInput, 'input').pipe(
  map(event => event.target.value.trim()),
  debounceTime(250),
  distinctUntilChanged(),
  switchMap(query => searchApi(query))
);

const subscription = results.subscribe({
  next: renderResults,
  error: showSearchError
});

debounceTime(250) waits for 250 milliseconds without a newer input before forwarding the latest value; distinctUntilChanged() suppresses a query equal to the previous forwarded query. The example’s 250 milliseconds is a chosen application setting, not an RxJS requirement. switchMap maps each query to an inner request Observable and switches the output to the newest one. For cancellable inner work, switching unsubscribes from the prior inner subscription; whether the underlying operation itself stops depends on the source and its cancellation behavior.

Choose flattening behavior from the work’s requirements

Flattening operators define how an outer sequence of inputs relates to inner sequences such as requests. Select based on cancellation, concurrency, ordering, and whether queued work is acceptable—not on a universal ranking.

Requirement Operator family to consider Behavior to account for
Only the newest inner result should remain relevant; prior inner subscriptions can be cancelled. switchMap Switches to the latest inner Observable. Cancellation of the underlying operation depends on the source.
Every inner operation should be allowed to finish, potentially at the same time as others. mergeMap Allows concurrent inner subscriptions; completion order may differ from input order.
Inner operations should run one at a time in input order, with later work waiting. concatMap Queues subsequent inner work until the current inner Observable completes.
Ignore new inputs while one inner operation is active. exhaustMap Does not queue or switch to new inputs during the active inner subscription.

For search suggestions, keeping the latest query is often the desired behavior, making switchMap a natural choice. For writes that must all be processed, switching away from an earlier operation may be wrong; sequential or concurrent processing may better fit the requirement.

Handle errors at the intended scope

Errors and completion are terminal notification paths, not ordinary values. An unhandled error terminates the affected Observable execution. Recovery operators such as catchError let a flow replace a failed sequence with a fallback or another Observable, but their placement determines what continues.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Recover a request without ending the input flow

For a typeahead, handle failure inside the inner request sequence if later keystrokes should still trigger searches. The following sketch assumes searchApi returns an Observable and of creates a fallback Observable:

import { of } from 'rxjs';
import { catchError, switchMap } from 'rxjs/operators';

const results = queries.pipe(
  switchMap(query =>
    searchApi(query).pipe(
      catchError(error => {
        showSearchError(error);
        return of([]);
      })
    )
  )
);

The failed inner request becomes an empty result, while the outer query sequence remains available for subsequent inputs. The fallback here is a design choice; an application may instead return a cached result or another recovery Observable.

Recovering outside the operation changes the scope

Placing error handling outside switchMap handles a failure of the composed flow. If it replaces that flow with a fallback sequence, the original input sequence may no longer be observed. Use that placement when ending or replacing the whole flow is intended. Retrying is another distinct strategy: it resubscribes according to the retry policy, so consider whether repeating the operation is safe and useful.

Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Reason about timing and test observable behavior

Time-based flows can be difficult to test reliably using real delays. RxJS marble diagrams use a compact timeline: - represents virtual time, letters represent emitted values, | marks completion, and # marks an error. Subscription diagrams use ^ and ! for subscription and unsubscription points.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

The official marble testing guide uses TestScheduler to run many timing-sensitive Observable tests on virtual time. A minimal test can assert both the output timeline and the subscription window:

import { TestScheduler } from 'rxjs/testing';
import { debounceTime } from 'rxjs/operators';

test('debounces a burst of values', () => {
  const scheduler = new TestScheduler((actual, expected) => {
    expect(actual).toEqual(expected);
  });

  scheduler.run(({ cold, expectObservable }) => {
    const source = cold('a--b----|');
    const result = source.pipe(debounceTime(2));
    expectObservable(result).toBe('-----b--|');
  });
});

In this timeline, a is followed by b before its debounce window elapses, so only b is emitted after the pause; the source then completes. Marble testing can also make cancellation and subscription intervals explicit, which is valuable when the timing of unsubscription is part of the behavior being tested.

There is an important boundary: Promise scheduling is not virtualized by TestScheduler. RxJS code that consumes Promises cannot be tested directly and reliably with virtual time alone. Test that part with the ordinary asynchronous testing facilities of your test framework.

Where to continue learning

The official RxJS overview introduces the library’s concepts, and the Observer guide explains notifications and consumers. The glossary and semantics guide clarifies RxJS terminology; consult version-specific API and migration material when applying details to a particular project.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

For more worked examples, see the Learn RxJS primer and its resource directory. For test structure and virtual-time details, use the RxJS marble testing guide.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.