Using observables

The Observable API provides a mechanism for handling streams of values, including asynchronous events. This guide explains how to transform, subscribe to, and unsubscribe from existing observables, using browser event streams as examples.

Before proceeding, you may wish to read the Observable API overview to familiarize yourself with the core concepts.

Obtaining an observable

Observable objects (commonly called observables) represent a stream of values that can be observed and transformed. The code that supplies these values is the producer, and the code that subscribes to receive and use them is the consumer. There are three main ways to obtain observables:

On the web, EventTarget objects are perhaps the most common use case of observables. The basic idea is this: wherever you have been writing addEventListener(eventType, handler), you can now write when(eventType).subscribe(handler) to achieve the same effect. For example, here is how you can listen for click events on the document body using when():

js
document.body.when("click").subscribe((event) => {
  console.log("Clicked!", event);
});

This not merely a syntactic difference. Observables let you compose operations such as filtering events, transforming their data, and stopping a stream when another event occurs. Beyond events, you can use the same operations with promises, iterables, and custom streams, such as the timer and element-size notifications examples shown in Creating custom observables.

Transforming an observable

An observable is a stream of values that you can transform into a new stream. An Observable object has several transform methods that resemble what you may already know from Iterator or Array.

Observable method Iterator equivalent Description
Observable.drop() Iterator.drop() Skips the first n values from the source observable.
Observable.filter() Iterator.filter() Skips values that don't match a predicate function.
Observable.flatMap() Iterator.flatMap() Maps each value to an observable, then flattens the resulting observables into a single observable.
Observable.inspect() N/A Calls callbacks to inspect values and the subscription lifecycle, while allowing further chaining.
Observable.map() Iterator.map() Maps each value to a new value using a mapping function.
Observable.switchMap() N/A Maps each value to an inner observable and emits values from only the latest inner observable.
Observable.take() Iterator.take() Takes only the first n values from the source observable.
Observable.takeUntil() N/A Like take(), but stops when a second observable emits a value.

These methods can be chained together to apply multiple transformations. You then subscribe to the final observable to receive the transformed values.

In the following example, we print the mouse coordinates to the screen whenever the mouse is moved over a couple of <div> elements. We won't show the HTML because it just includes the <div> elements plus a single <p> element to display the data.

In the example CSS, we give the <div> elements a height, background-color, and margin-bottom:

css
div {
  height: 150px;
  background-color: purple;
  margin-bottom: 10px;
}

The JavaScript looks like this:

js
const outputElem = document.querySelector("p");

document.body
  .when("mousemove")
  .filter((e) => e.target.matches("div"))
  .map((e) => ({ x: e.clientX, y: e.clientY }))
  .subscribe((p) => {
    outputElem.textContent = `${p.x},${p.y}`;
  });

In this snippet, the page's <body> element is an EventTarget. We obtain a stream of mousemove events fired on it using the when() method.

We then specify a pipeline:

  • Observable.filter() filters the events passed through the pipeline to only events fired on the <div> element (tested using the Element.matches() method) and not other body descendants.
  • Observable.map() maps the fired mousemove event objects to new objects containing the coordinates of the mouse cursor when the event was fired.

Finally, Observable.subscribe() subscribes to the observable. We pass a callback function to the subscribe() method. This callback is called each time a mousemove event passes the filter.

The rendered output looks like this:

Try moving the mouse over the top of the example; the coordinates are printed to the <p> only when the <div> elements are moved over, not the areas outside the <div>s.

Note: Observables are "lazy" — events don't start being passed through them, nor do they queue any data, until they have at least one subscriber. In the previous example, if you remove the subscribe() method call and add logs inside the filter() and map() methods, you will see that they don't log anything. Once subscribe() is called at the end of the pipeline, all previous observables in the chain also become subscribed and start processing data.

Working with inner observables

Some operations produce another stream for each source value: a click might start an upload, or a change to a search field might start a request. These are called inner observables. Using map() alone would send the inner observable objects to your observer. Observable.flatMap() and Observable.switchMap() instead subscribe to them and forward their values.

  • flatMap() processes source values sequentially. It waits for the current inner observable to complete before calling the mapper for the next queued source value. Use this when every operation should finish in order. If an inner observable never completes, later source values remain queued.
  • switchMap() unsubscribes from the current inner observable when a new source value arrives, then calls the mapper and subscribes to its result. Use this when only the latest operation's results are relevant.

Both methods convert the mapper's result using Observable.from(), so the mapper can also return a promise, iterable, or async iterable. A promise contributes its fulfillment value and then completes; a rejection becomes an error in the inner observable. In contrast, map() forwards a returned promise as a value without awaiting it.

The following live example searches a small list of fruit names. Type into the search field to start a search; the simulated response takes one second, so you can type again while a previous search is pending. An async mapper lets us treat each search as one operation:

html
<label>Search fruit: <input type="search" /></label>
<p id="results">Type a fruit name</p>
js
const searchInput = document.querySelector("input[type='search']");
const results = document.querySelector("#results");

async function search(query) {
  const fruits = ["Apple", "Apricot", "Banana", "Cherry", "Pear"];
  await new Promise((resolve) => setTimeout(resolve, 1000));
  return fruits.filter((fruit) =>
    fruit.toLowerCase().includes(query.toLowerCase()),
  );
}

const queries = searchInput.when("input").map(() => searchInput.value);

queries.switchMap(search).subscribe({
  next(data) {
    results.textContent = JSON.stringify(data);
  },
  error(error) {
    results.textContent = error.message;
  },
});

In a real search app, search() could fetch and parse a JSON response instead:

js
async function search(query) {
  const response = await fetch(`/search?q=${encodeURIComponent(query)}`);
  if (!response.ok) {
    throw new Error(`Search failed: ${response.status}`);
  }
  return response.json();
}

If another input event arrives before the previous search finishes, the previous result is no longer forwarded. However, unsubscribing from an observable created from a promise does not cancel the work behind that promise: the previous request can still finish. To cancel the request itself, return a custom observable that passes subscriber.signal to fetch(), as shown in Canceling asynchronous work.

Inspecting a pipeline

Observable.inspect() lets you run a side effect, such as logging, while forwarding values unchanged. Unlike subscribe(), it returns an observable and does not start the pipeline by itself:

js
document.body
  .when("click")
  .inspect((event) => console.log("Click:", event.target))
  .map((event) => ({ x: event.clientX, y: event.clientY }))
  .subscribe((point) => console.log("Coordinates:", point));

You can also pass an object with subscribe, next, error, complete, and abort callbacks to inspect the subscription's lifecycle. An inspect() callback can affect the pipeline if it throws; for example, an exception in its next callback becomes an error in the returned observable.

See the inspect() reference for details.

Aggregating values

The previously introduced group of methods return another observable, allowing you to chain multiple transformations together. You can then subscribe to the final observable to receive the transformed values.

But you don't always want to process each value individually. Sometimes you are just interested in getting a single aggregated value from the entire stream. For this purpose, the Observable API provides the following methods:

Observable method Iterator equivalent Description
Observable.every() Iterator.every() Returns false when the predicate returns false for any value; true otherwise.
Observable.find() Iterator.find() Returns the first value that matches a predicate.
Observable.first() N/A Returns the first value from the observable.
Observable.forEach() Iterator.forEach() Calls a function for each value.
Observable.last() N/A Returns the last value from the observable.
Observable.reduce() Iterator.reduce() Aggregates the values to a single value.
Observable.some() Iterator.some() Returns true when the predicate returns true for any value; false otherwise.
Observable.toArray() Iterator.toArray() Collects all values into an array.

All these methods return promises that fulfill with the described results. Depending on the method, the promise fulfills as soon as the result is determined or when the observable completes. It can also reject, for example, if the observable errors or the subscription is aborted.

Unlike the transformation methods, these aggregation methods implicitly subscribe: the pipeline starts receiving values as soon as one of these methods is called.

This example searches for the first mouse position beyond 200 on both axes within a single panel. Until a match is found, it displays the current position (also demonstrating how inspect() works). Click Restart after finding a match to try again.

html
<div id="target">Move the pointer beyond 200,200.</div>
<p></p>
<button disabled>Restart</button>
css
#target {
  height: 300px;
  background-color: lavender;
}
js
const target = document.querySelector("#target");
const outputElem = document.querySelector("p");
const restart = document.querySelector("button");

function start() {
  restart.disabled = true;
  outputElem.textContent = "Move the mouse around...";
  target
    .when("mousemove")
    .map((event) => ({ x: event.offsetX, y: event.offsetY }))
    .inspect(({ x, y }) => {
      outputElem.textContent = `Target of 200,200 not yet reached (current ${x},${y})`;
    })
    .find(({ x, y }) => x > 200 && y > 200)
    .then(({ x, y }) => {
      outputElem.textContent = `Target coordinates found: ${x},${y}`;
      restart.disabled = false;
    });
}

restart.when("click").subscribe(start);
start();

The rendered output looks like this:

Subscribing to an observable

We already showed basic subscribe() usage in the previous sections, but let's look at it in a bit more detail.

Just like a promise can send notifications either as "fulfilled" or "rejected", an observable can also send multiple types of notifications to its subscribers, each one corresponding to a different method you can pass into subscribe().

  • next(value): Called whenever a new value is available in the observable stream. In the examples above, we passed a single function into subscribe(), which is shorthand for passing an object with just a next() method.
  • error(err): Called when the observable signals an error. Exceptions thrown by the observer's own callbacks are reported as uncaught errors to the global object instead of passed to this error() callback.
  • complete(): Called when the observable has finished sending values. We'll discuss this more in the Creating custom observables guide. Event streams are infinite, but can be made finite by calling take() or takeUntil().

For example, this observable completes after three clicks and logs a completion message:

js
document.body
  .when("click")
  .take(3)
  .subscribe({
    next(event) {
      console.log("Clicked at", event.clientX, event.clientY);
    },
    complete() {
      console.log("Observable complete");
    },
  });

Internally, each active observable subscription has a list of observers — objects containing zero or more of these three callbacks. You can call subscribe() multiple times on the same observable to register multiple observers. For example:

js
const clickObservable = document.body.when("click");

clickObservable.subscribe((e) => {
  console.log("Observer 1: Clicked at", e.clientX, e.clientY);
});
clickObservable.subscribe((e) => {
  console.log("Observer 2: Clicked at", e.clientX, e.clientY);
});

// For every click, both observers will be called

Concurrent observers share the subscription, and each receives values emitted while it is subscribed; previously emitted values are not replayed to new observers. This differs from sharing an iterator, where each consumer's next() call advances the same iterator rather than broadcasting a value to all consumers.

Note: This shared-subscription behavior may change. A proposal to give each observer its own Subscriber would make each subscription start a separate execution instead of reusing an active subscription.

For this particular example, the behavior would be the same, but each subscribe() call will register a new click event listener instead of reusing the same event listener.

Unsubscribing from an observable

An observer can also be unsubscribed from the observable, which means its callbacks will no longer be called. If an observable has no more observers, its shared subscription becomes inactive, and it runs callbacks known as teardown callbacks. These callbacks release resources, such as the event listener registered by when(). Custom observables must implement this cleanup themselves, as described in Creating custom observables.

The canonical way to unsubscribe from an observable is to use an AbortController. With this method, you can unsubscribe mid-way through observable data processing, at any point you like. To do this, you create an AbortController and pass its signal when you call subscribe(). You can then call AbortController.abort() on the controller, which unsubscribes all observers associated with that signal.

For example, we can modify our earlier Basic when() example to unsubscribe when the user clicks anywhere on the page. This means the output will stop updating. Click Restart to start a new subscription.

js
const outputElem = document.querySelector("p");
const restart = document.querySelector("button");

function start() {
  restart.disabled = true;
  outputElem.textContent = "Move the mouse over a purple area";
  // Create controller
  const controller = new AbortController();

  document.body
    .when("mousemove")
    .filter((e) => e.target.matches("div"))
    .map((e) => ({ x: e.clientX, y: e.clientY }))
    .subscribe(
      (p) => {
        outputElem.textContent = `${p.x},${p.y}`;
      },
      // Register observer with signal
      { signal: controller.signal },
    );

  document.body
    .when("click")
    .filter((event) => event.target !== restart)
    .take(1)
    .subscribe(() => {
      // Unsubscribe on click
      controller.abort();
      outputElem.textContent += " — Stopped. Click Restart to try again.";
      restart.disabled = false;
    });
}

restart.when("click").subscribe(start);
start();

The take(1) call completes the click stream after the first click, removing its event listener. This is similar to using { once: true } with addEventListener().

In the next example, another observable emits a value to trigger the abort condition: document.body.when("click"). The takeUntil() method is a transformation method and therefore returns an observable. This means you can insert it in the pipeline to specify a condition under which the unsubscribe action occurs. The following code achieves the same effect as the previous example:

js
const outputElem = document.querySelector("p");

document.body
  .when("mousemove")
  .filter((e) => e.target.matches("div"))
  .map((e) => ({ x: e.clientX, y: e.clientY }))
  // When the click event fires, unsubscribe
  .takeUntil(document.body.when("click"))
  .subscribe((p) => {
    outputElem.textContent = `${p.x},${p.y}`;
  });

Note: The takeUntil() method converts its input to an observable. You can pass a promise to stop when it fulfills, or sync/async iterables for when they produce the first value, if any. Check Observable.from() for how it does the conversion.

An AbortController lets you unsubscribe at any point in your code. Separate controllers let you unsubscribe observers independently. Aborting does not call the observer's complete callback. In contrast, takeUntil() completes the observable it returns and notifies that observable's observers through their complete callbacks. Other observers subscribed directly to the source observable remain subscribed.

Handling errors

An error ends the affected subscription. Errors from a source propagate through the pipeline, and exceptions thrown by transformation callbacks, such as a map() mapper or filter() predicate, become errors in the returned observable. An observer's error callback reports or handles the failure, but does not resume the subscription. If the observer has no error callback, the error is reported as an uncaught error to the global object.

Observable.catch() lets a pipeline recover by subscribing to a replacement stream. Its callback receives the error and returns an observable, or any value convertible by Observable.from(). For example, returning [] completes the replacement without emitting a value. It does not retry the failed source.

Placement matters for inner observables. In the search example, a failure in the current request ends the search subscription, so subsequent input events no longer start searches. We can replace that subscription code with the following to recover from failures while the inner subscription is active:

js
queries
  .switchMap((query) =>
    Observable.from(search(query)).catch((error) => {
      results.textContent = error.message;
      return [];
    }),
  )
  .subscribe((data) => {
    results.textContent = JSON.stringify(data);
  });

Here, catch() handles only the inner request's failure. Its empty replacement completes, while the outer subscription continues listening for input. Placing catch() after switchMap() would instead replace the entire search pipeline: returning [] there would complete it and stop listening for input.

js
queries
  .switchMap(search)
  .catch((error) => {
    results.textContent = error.message;
    return [];
  })
  .subscribe({
    next(data) {
      results.textContent = JSON.stringify(data);
    },
    complete() {
      console.log("Search pipeline ended");
    },
  });

The same distinction applies to flatMap().

This does not handle failures from requests that switchMap() has already unsubscribed from. If such a request's promise later rejects, Observable.from() reports the error to the global object because its subscriber is inactive; the catch() callback is no longer subscribed. To cancel obsolete requests and avoid reporting their cancellation as an error, use the custom fetchJSON() producer in Canceling asynchronous work, which checks subscriber.active before forwarding a rejection.

Exceptions thrown by callbacks passed to subscribe() are different: they are reported to the global object, rather than becoming errors that a pipeline's catch() can recover from. Similarly, an async next callback's returned promise is not awaited; handle its rejections yourself, or use flatMap() or switchMap() to incorporate the asynchronous work into the pipeline.

For example, a pipeline's catch() does not handle an exception thrown by its observer:

js
Observable.from([1, 2, 3])
  .catch(() => [0])
  .subscribe(() => {
    throw new Error("Reported as an uncaught error on window");
  });

If an observer has foreseeable error conditions, handle its failure explicitly:

js
queries.subscribe(async (query) => {
  try {
    results.textContent = JSON.stringify(await search(query));
  } catch (error) {
    results.textContent = error.message;
  }
});

Unlike the switchMap() version, this version handles failures, but does not discard outdated results.

Running cleanup

Observable.finally() returns an observable that forwards the source's values and notifications, and runs a callback when its subscription ends through completion, error, or unsubscribing. This makes it useful for cleanup that a complete callback alone would miss:

js
document.body
  .when("click")
  .take(3)
  .finally(() => console.log("Stopped observing clicks"))
  .subscribe((event) => console.log(event.target));

In the previous snippet, the finally() callback runs after the subscription ends (after three clicks, as specified by take(3)). It will also run if the subscription is aborted early. It does not await a returned promise. Use finally() to attach cleanup when composing a pipeline; use addTeardown() to implement the producer's own resource cleanup, as described in Creating custom observables.

Example: canvas drawing

In this example we create a basic <canvas>-based drawing app, which brings together the APIs we have seen so far to demonstrate how observables help you declaratively implement complex event handling logic.

HTML

The markup includes a <canvas> element to draw onto, and a <form> containing two <input> controls to allow the user to choose a new pen size and color (a range slider and a color picker, respectively). We also include an <output> element to display the current range value.

html
<canvas></canvas>
<form>
  <div>
    <label for="size">Choose pen size:</label>
    <input id="size" type="range" min="1" max="40" value="10" />
    <output for="size">10</output>
  </div>
  <div>
    <label for="color">Choose pen color:</label>
    <input id="color" type="color" />
  </div>
</form>

CSS

In the CSS, we make the <body> element span the full width and height of the page. We also make the <form> sit on top of the <canvas>, and use inset properties to make it stick to the top, left, and right of the <body>.

The rest of the CSS isn't important to the understanding of the overall example, so we won't explain it here, but we've included it all below so you can see it.

css
* {
  box-sizing: border-box;
}

html {
  font-family: Arial, Helvetica, sans-serif;
  height: 100%;
}

body {
  margin: 0;
  height: inherit;
  overflow: hidden;
}

form {
  position: absolute;
  top: 0;
  right: 0;
  left: 0;
  padding: 10px;
  background: rgb(218 112 214 / 0.85);
  box-shadow: 0 1px 3px black;
}

form > div {
  display: flex;
  align-items: center;
}

form > div:first-child {
  margin-bottom: 10px;
}

form label {
  width: 140px;
}

form input {
  width: 100px;
}

JavaScript

In our script, we first grab references to our <canvas>, <form>, <input>, and <output> elements:

js
const canvas = document.querySelector("canvas");
const form = document.querySelector("form");
const sizeInput = document.querySelector("[type='range']");
const sizeOutput = document.querySelector("output");
const colorInput = document.querySelector("[type='color']");

Next, we synchronize the canvas's width and height to the clientWidth/clientHeight of the <body>. This is implemented in the sizeCanvas() function. It is called when the app starts and whenever the window resizes, using when("resize").subscribe(sizeCanvas) (which has the same effect as addEventListener("resize", sizeCanvas)). Setting the canvas dimensions also clears the drawing.

js
function sizeCanvas() {
  canvas.width = document.body.clientWidth;
  canvas.height = document.body.clientHeight;
}

sizeCanvas();

window.when("resize").subscribe(sizeCanvas);

Next, we define the variables and functions we need to draw on our <canvas>. First, we grab a reference to the <canvas> 2D rendering context, and store initial values for the pen size and color in penSize and penColor, respectively. The updatePenSize() function sets penSize to the range slider's valueAsNumber and displays its value in the <output> element. The updatePenColor() function sets penColor to the color picker's value. These functions handle the input event on the range slider and the change event on the color picker, respectively.

js
const ctx = canvas.getContext("2d");
let penSize = 10;
let penColor = "black";

function updatePenSize() {
  penSize = sizeInput.valueAsNumber;
  sizeOutput.textContent = sizeInput.value;
}

function updatePenColor() {
  penColor = colorInput.value;
}

sizeInput.when("input").subscribe(updatePenSize);
colorInput.when("change").subscribe(updatePenColor);

Now on to our main draw() function. Here, we hide the <form> so we can see the whole <canvas> when we start to draw. We then set the canvas context's fillStyle to the penColor, start drawing a path using beginPath(), draw a single circle of the specified penSize at the event object's x and y coordinates (more on those later) using arc(), and render the drawing on the canvas using fill(). This (along with the observable code you'll see later) draws a circle at the current mouse coordinates every time the mouse moves.

js
function draw(e) {
  form.style.display = "none";

  ctx.fillStyle = penColor;
  ctx.beginPath();
  ctx.arc(
    e.x - canvas.offsetLeft,
    e.y - canvas.offsetTop,
    penSize,
    0,
    2 * Math.PI,
  );
  ctx.fill();
}

The last function we define is finishDraw(), which shows the <form> again when the drawing is finished.

js
function finishDraw() {
  form.style.display = "block";
}

Finally, we create an observable for mousedown events on the <canvas>. For each press of the primary mouse button, Observable.flatMap() subscribes to a stream of mousemove events that ends when the button is released. We listen for mouseup on the document object so drawing also stops if the mouse is released outside the canvas. We then use Observable.map() to extract the mouse coordinates and subscribe() to pass them to draw().

js
canvas
  .when("mousedown")
  .filter((e) => e.button === 0)
  .flatMap(() => {
    const mouseUp = document
      .when("mouseup")
      .filter((e) => e.button === 0)
      .finally(finishDraw);
    return canvas.when("mousemove").takeUntil(mouseUp);
  })
  .map((e) => ({ x: e.clientX, y: e.clientY }))
  .subscribe(draw);

The pipeline processes a mousedown → mousemove… → mouseup sequence:

  1. Each primary-button mousedown event passes the filter and triggers the flatMap() callback.
  2. The callback returns canvas.when("mousemove").takeUntil(mouseUp). Subscribing to this inner observable starts listening for mouse movements and for the button release.
  3. Each mouse movement on the canvas passes through flatMap() to map(), which extracts its coordinates, and then to draw().
  4. When the primary button is released, Observable.takeUntil() completes the inner observable and unsubscribes from both event streams. The Observable.finally() callback runs during cleanup and calls finishDraw() to show the form again.
  5. The outer mousedown subscription remains active, so the next press starts a new drawing sequence.

The final effect is that we react to mousemove events (just like our first example on this page!) but we only start listening when mousedown fires and stop listening when mouseup fires. With the help of Observable, we have successfully composed three parallel event streams into a single, coherent sequence.

Result

The example renders like this:

See also