Creating custom observables
The Observable API lets you create custom streams of values using the Observable() constructor. This guide explains how to produce values, complete a stream, and clean up resources when a subscription ends.
Before proceeding, read Using observables to familiarize yourself with subscribing to and consuming observable streams.
Creating an observable
Like promises, observables are created by passing a callback to the Observable() constructor. The callback's job is to do work and push data to the observable's subscribers. The callback isn't called immediately: it's called when the first observer subscribes, either via subscribe(), an aggregation method, or subscribing to a downstream observable created by a transformation method.
The Observable() callback receives a Subscriber object. You can call methods on this object to dispatch data to all observers subscribed to the observable. Additional observers share the same underlying subscription until it completes, errors, or all observers unsubscribe. After that, the callback is called again when the next observer subscribes.
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.
The Subscriber object has the following methods:
next(value): Sends a value to each observer'snextcallback. This can be called any number of times while the subscription is active.complete(): Ends the subscription successfully and calls each observer'scompletecallback without arguments.error(error): Ends the subscription with an error and passes the error to each observer'serrorcallback. If an observer has noerrorcallback, the error is reported as an uncaught error to the global object.addTeardown(callback): Registers a callback to clean up resources when the subscription ends. See Teardown for details.
With a subscription set up, the custom observable can send any number of values by calling subscriber.next(), optionally followed by a call to subscriber.complete() or subscriber.error() to signal that the stream of data is finished.
In this example, we print the numbers 1 to 10 to the page, then print a message to say that the count is complete. We won't show the HTML because it just includes a single <p> element to display the count and a <button> to start the count.
In the JavaScript, we use the Observable() constructor to create a new observable. Inside its callback function, we declare a variable i with a value of 1. We then use a Window.setInterval() call to check the value of i every timerInterval milliseconds. If the value has exceeded the specified number of iterations, we call the Subscriber.complete() method to complete the subscription. If not, we call Subscriber.next() to send the current value of i to the observers. At the end of the interval callback, i is incremented by 1.
function makeTimer(timerInterval, iterations = Infinity) {
return new Observable((subscriber) => {
let i = 1;
const interval = setInterval(() => {
if (i === iterations + 1) {
subscriber.complete();
clearInterval(interval);
} else {
subscriber.next(i);
}
i++;
}, timerInterval);
});
}
Note: This function is not currently production-ready! It allows multiple intervals to be created when the user clicks the button multiple times before the count finishes. We'll fix this problem in the Teardown section.
Next, we define an init() function in which we subscribe to the observable by calling Observable.subscribe(). The object passed to subscribe() defines the observer's callbacks: next() prints the value received from the producer to the <p> element, and complete() displays a completion message.
function init() {
makeTimer(500, 10).subscribe({
next(value) {
outputElem.textContent = value;
},
complete() {
outputElem.textContent = "Count complete; click to restart.";
},
});
}
Finally, the init() function is called in response to the click event on the <button> element, using when() and subscribe().
btn.when("click").subscribe(init);
The rendered output looks like this:
Click the button. Every 500 milliseconds, the value of i is printed to the page and then incremented by 1. On the next interval after printing 10, the subscription completes and the paragraph displays "Count complete; click to restart."
Note:
The producer calls methods on the Subscriber object to send notifications; the consumer defines the corresponding callbacks in the object passed to subscribe(). As you'll see in the next section, cleanup is also registered by the producer, using Subscriber.addTeardown() inside the constructor callback.
Teardown
The previous example is not production-ready because subscribing to each new makeTimer() observable creates a new interval. If the user clicks the <button> multiple times, they will create multiple intervals, all trying to update the same <p> element. To fix this, we need to unsubscribe from the previous observable and clear its interval before starting a new count.
First of all, we'll use an AbortController to unsubscribe when the user clicks the button again:
let controller;
function init() {
controller?.abort();
controller = new AbortController();
makeTimer(500, 10).subscribe(
{
next(value) {
outputElem.textContent = value;
},
complete() {
outputElem.textContent = "Count complete; click to restart.";
},
},
{ signal: controller.signal },
);
}
btn.when("click").subscribe(init);
The makeTimer() observable is stopped, but unfortunately, the interval created inside it is not cleared until the count reaches 11. This is fine for our example because the interval will eventually clear itself, but in a real-world scenario this could lead to memory leaks and unexpected behavior. We need to make sure that clearInterval is called deterministically when the observable becomes inactive, not just when it reaches the termination point. We do this by adding a teardown to the observable.
The teardown logic is passed as a callback to Subscriber.addTeardown():
function makeTimer(timerInterval, iterations = Infinity) {
return new Observable((subscriber) => {
let i = 1;
const interval = setInterval(() => {
if (i === iterations + 1) {
subscriber.complete();
} else {
subscriber.next(i);
}
i++;
}, timerInterval);
subscriber.addTeardown(() => {
clearInterval(interval);
});
});
}
The callback registered with addTeardown() runs when Subscriber.complete() or Subscriber.error() closes the subscription, before the observers' completion or error callbacks. It also runs when all observers unsubscribe.
In this case, the teardown callback clears the interval via Window.clearInterval(), so it stops running as soon as the subscription ends.
The example now renders like so:
Press the <button> while the count is running; the count will restart from 1.
Producing values synchronously
An observable does not have to wait for an event or asynchronous operation. The constructor callback runs synchronously when a subscription starts, and subscriber.next() invokes observers' callbacks synchronously. A subscription can therefore receive values and complete before subscribe() returns:
const numbers = new Observable((subscriber) => {
for (let value = 1; value <= 10; value++) {
if (!subscriber.active) {
return;
}
subscriber.next(value);
}
subscriber.complete();
});
console.log("Before subscribing");
numbers.take(3).subscribe({
next(value) {
console.log(value);
},
complete() {
console.log("Complete");
},
});
console.log("After subscribing");
// Before subscribing
// 1
// 2
// 3
// Complete
// After subscribing
After receiving three values, take(3) completes its output and unsubscribes from numbers. With no observers remaining, the producer's active property becomes false. Checking it before each iteration stops the producer from doing unnecessary work. This check also handles a subscription started with an already aborted signal.
Calling subscriber.complete() or subscriber.error() does not stop the producer's JavaScript execution. Use return, break, or a Subscriber.active check to stop producing values when appropriate. Functionality intended to interrupt synchronous production needs to be available before subscribe() is called, such as through an AbortSignal supplied in its options.
Canceling asynchronous work
When a consumer unsubscribes, the producer is responsible for stopping work that is no longer needed. For APIs that accept an AbortSignal, such as fetch(), you can pass subscriber.signal directly. This signal is aborted when the shared subscription ends, including when all observers unsubscribe.
The following function creates an observable that fetches JSON, emits the parsed data, and completes. The request starts when the observable is subscribed to:
function fetchJSON(url) {
return new Observable((subscriber) => {
fetch(url, { signal: subscriber.signal })
.then((response) => {
if (!response.ok) {
throw new Error(`Request failed: ${response.status}`);
}
return response.json();
})
.then((data) => {
subscriber.next(data);
subscriber.complete();
})
.catch((error) => {
if (subscriber.active) {
subscriber.error(error);
}
});
});
}
If the subscription ends while the request or response body is pending, aborting subscriber.signal cancels that work and causes the fetch or body-reading promise to reject. The rejection handler checks subscriber.active before forwarding the error: calling subscriber.error() after cancellation would report the error to the global object. While the subscription is active, request and JSON-parsing failures are forwarded to its observers.
The Observable() constructor callback is not asynchronous. Its return value is ignored, so returning a promise would not make the observable wait for it or automatically forward its rejection. Instead, this example explicitly calls next(), complete(), and error() from the promise handlers.
We can use fetchJSON() in a search pipeline like the one in Using observables. Suppose the page contains a search input and a results element:
const searchInput = document.querySelector("input[type='search']");
const results = document.querySelector("#results");
searchInput
.when("input")
.map(() => searchInput.value)
.switchMap((query) =>
fetchJSON(`/search?q=${encodeURIComponent(query)}`).catch((error) => {
results.textContent = error.message;
return [];
}),
)
.subscribe((data) => {
results.textContent = JSON.stringify(data);
});
Each new input event makes switchMap() unsubscribe from the previous inner observable. Because that request has no other observers, its subscriber.signal is aborted, canceling the pending request. This both prevents obsolete results from being displayed and stops the underlying work. The inner catch() handles request failures without ending the input subscription, so the user can try another search.
Example: observing element size
Custom observables can wrap APIs that deliver notifications through callbacks. In this example, we wrap a ResizeObserver to create a stream of an element's dimensions. Resizing an element does not fire a resize event on that element, so we cannot obtain this stream using when().
HTML and CSS
The markup contains a resizable panel and a paragraph to display its dimensions, as well as the Stop and Restart buttons. The resize property lets the user resize the panel by dragging its corner.
<div id="panel">Drag my corner to resize me.</div>
<p id="dimensions"></p>
<button>Stop</button>
<button id="restart" disabled>Restart</button>
#panel {
width: 200px;
height: 100px;
min-width: 100px;
max-width: 90%;
min-height: 50px;
max-height: 200px;
overflow: auto;
resize: both;
border: 1px solid;
}
JavaScript
Inside a custom observable's callback, we create a ResizeObserver and start observing the panel. Each notification passes the panel's content rectangle to Subscriber.next(). We also register a teardown callback to disconnect the ResizeObserver when the subscription ends.
In this example, clicking Stop completes the stream returned by takeUntil() and unsubscribes its only observer from sizes, triggering the teardown. Clicking Restart starts a new subscription and creates a new ResizeObserver.
const panel = document.getElementById("panel");
const dimensions = document.getElementById("dimensions");
const sizes = new Observable((subscriber) => {
const observer = new ResizeObserver(([entry]) => {
subscriber.next(entry.contentRect);
});
observer.observe(panel);
subscriber.addTeardown(() => observer.disconnect());
});
const stop = document.querySelector("button");
const restart = document.querySelector("#restart");
function start() {
restart.disabled = true;
sizes.takeUntil(stop.when("click")).subscribe({
next({ width, height }) {
dimensions.textContent = `Content size: ${Math.round(width)} × ${Math.round(height)} pixels`;
},
complete() {
restart.disabled = false;
},
});
}
restart.when("click").subscribe(start);
start();
Subscribing starts the ResizeObserver, which reports the initial size and subsequent size changes. The subscription callback displays these dimensions outside the panel so updating the output does not affect the observed element's size.