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?
answer
- notifier must emit, not finish
- complete is ignored
- sync notifier at subscribe
- operators after takeUntil
basics
~20 stakeUntil 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 linesimport { 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
Recall that the notifier must emit a value; completing it without emitting does not stop the source.
Explain the notifier-first subscription, forwarded errors, and why a synchronous notifier means the source is never subscribed.
Diagnose a leak caused by takeUntil placed before switchMap and use the complete notification it produces to finish batches.
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.