skip to content

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%

answer

  1. notifier must emit, not finish
  2. complete is ignored
  3. sync notifier at subscribe
  4. operators after takeUntil

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.

solid answer

~30 s

`takeUntil(notifier)` subscribes to the notifier first and completes the output on the notifier's first `next`. Its completion handler is a no-op, so `stop$.complete()` alone changes nothing: call `stop$.next()` first. A notifier `error` is forwarded, so the output errors instead of completing. If the notifier emits synchronously during subscription — `of(true)`, a `BehaviorSubject` — the source is never subscribed and the output completes at once. Placement matters too: in `source$.pipe(takeUntil(stop$), switchMap(id => poll(id)))`, takeUntil completes only the outer stream, and `switchMap` completes when its active inner stream also completes, so polling continues. Put `takeUntil` after operators that open inner subscriptions.

code

ts · 16 lines
ts
import { Subject, interval, map, switchMap, takeUntil, toArray } from 'rxjs';

const stop$ = new Subject<void>();

const batch$ = interval(500).pipe(
  switchMap(() => interval(100).pipe(map(() => Math.random()))),
  takeUntil(stop$), // after switchMap: the inner interval stops too
  toArray()
);

batch$.subscribe(values => console.log('recorded', values.length));

setTimeout(() => {
  stop$.next(); // complete() alone would not stop the feed
  stop$.complete();
}, 2000);

go deeper

for a junior

Recall that the notifier must emit a value; completing it without emitting does not stop the source.

for a middle

Explain the notifier-first subscription, forwarded errors, and why a synchronous notifier means the source is never subscribed.

for a senior

Diagnose a leak caused by takeUntil placed before switchMap and use the complete notification it produces to finish batches.

for a principal

Define team rules for stream termination so that notifiers, placement and lint checks make these misfires impossible to merge.

## How takeUntil actually works `takeUntil(notifier)` is a pipeable operator that mirrors its source **until the notifier emits a value**. In RxJS 7.8 the implementation does three things on subscription: 1. It subscribes to the **notifier first**, with a `next` handler that completes the output and a `complete` handler that is a **no-op**. 2. If the output is still open, it subscribes to the **source**. 3. From then on, source values pass through; the first notifier value completes the output, which also unsubscribes from both source and notifier. The notifier can be any `ObservableInput`: an observable, a promise, an array. The documentation states it plainly: if the notifier completes without emitting, `takeUntil` passes all values. ## Why stop$.complete() does not stop the feed ```ts const stop$ = new Subject<void>(); const feed$ = sensor$.pipe(takeUntil(stop$)); stop$.complete(); // feed$ keeps emitting ``` Completion of the notifier reaches the no-op handler. Nothing emitted, so nothing stops. The fix is to emit, and optionally complete afterwards: ```ts stop$.next(); stop$.complete(); ``` ## The other ways it misfires | situation | what happens | why | |---|---|---| | notifier completes without a value | output keeps mirroring the source | completion handler is a no-op | | notifier errors | output **errors** | notifier errors are forwarded downstream | | notifier emits synchronously on subscribe | source never subscribed, output completes at once | the source subscription is skipped when the output is already closed | | flattening operator placed after `takeUntil` | inner work continues | the flattening operator waits for its active inner stream | ### Synchronous notifiers `takeUntil(of(true))` or `takeUntil(new BehaviorSubject(false))` emit during subscription. The output completes before the source is even subscribed — RxJS changed this deliberately so a notifier that has already fired never starts the source. A `BehaviorSubject` is a poor notifier for exactly this reason: it replays its current value to each new subscriber. A plain `Subject` or a `ReplaySubject` used on purpose makes the intent explicit. ### Placement before a flattening operator `switchMap`, `mergeMap`, `concatMap` and `exhaustMap` complete only when **the outer source has completed and no inner stream is active**. If `takeUntil` sits before them, it completes the outer stream, but a running inner stream — a polling loop, a long request — keeps going, and its values still reach the subscriber: ```ts sensorIds$.pipe( takeUntil(stop$), switchMap(id => pollReading(id)) // keeps polling after stop$ ); ``` Move it after the flattening operator so the whole chain completes on the notifier: ```ts sensorIds$.pipe( switchMap(id => pollReading(id)), takeUntil(stop$) ); ``` ## Completion, not just unsubscription `takeUntil` produces a real **`complete` notification** downstream. That differs from calling `unsubscribe()`, which sends nothing to the observer: - `complete` callbacks run, and operators such as `last`, `toArray` or `reduce` can emit their result; - `finalize` runs in both cases. That makes `takeUntil` useful beyond teardown: *record readings until the operator presses stop, then emit the collected batch* is `takeUntil(stop$)` followed by `toArray()`. ## Testing the behaviour These rules are easy to pin down in a unit test with synchronous sources: - `of(1, 2, 3).pipe(takeUntil(EMPTY))` emits `1, 2, 3` — an empty notifier only completes, so it never stops the source; - `of(1, 2, 3).pipe(takeUntil(of(true)))` emits nothing and completes — the synchronous notifier wins before the source is subscribed; - `of(1, 2, 3).pipe(takeUntil(throwError(() => new Error('stop'))))` errors — the notifier's error is forwarded. Time-based cases, such as a notifier firing between two readings, are best expressed as marble tests with virtual time. ## A review checklist 1. Does the notifier **emit** (not only complete) when stopping is intended? 2. Can the notifier **error**, and is an error on the output acceptable? 3. Could it emit **synchronously** at subscribe time and suppress the source entirely? 4. Is `takeUntil` placed **after** every operator that opens inner subscriptions? For component teardown in Angular, the framework's own `takeUntilDestroyed` covers the lifecycle case; the semantics above are the plain RxJS operator.

  • Why is a BehaviorSubject a risky takeUntil notifier?
    It replays its current value to every new subscriber. `takeUntil` subscribes to the notifier before the source, receives that value synchronously and completes the output without ever subscribing to the source, whatever the stored value is.
  • How does completion from takeUntil differ from calling unsubscribe()?
    `takeUntil` sends a `complete` notification, so observers' complete callbacks run and operators such as `toArray` or `last` emit their result. `unsubscribe()` sends nothing downstream; only teardown and `finalize` run.

saying these in an interview costs you the question

  • Completing the notifier Subject is enough to stop takeUntil.
  • A notifier error is swallowed and the output simply completes.
  • takeUntil can go anywhere in the pipe with the same effect.
  • A BehaviorSubject is a good notifier because it always has a value.
  • takeUntil only unsubscribes and never sends a complete notification.