skip to content

Multi-Source Joins

combineLatest, forkJoin, zip, withLatestFrom, merge, concat and race join several streams with different waiting rules. Interviewers ask why forkJoin never emits if one source never completes.

part ofRxJSoverview, primer and where to startread it →
on this pageshow

explore

questions

5

In RxJS, how do combineLatest and forkJoin differ, and why can a forkJoin never emit when one source never completes?

level: middleimportance: must knowfreq 76%

answer

  1. waiting rule versus finishing rule
  2. first value from every source
  3. last values, only after completion
  4. one emission, then complete
  5. empty source ends forkJoin silently

basics

~20 s

combineLatest emits the latest value of every source once each has emitted, then again on every new value; forkJoin waits for every source to complete and emits their last values once, so a never-completing source blocks it forever.

solid answer

~40 s

`combineLatest([a$, b$])` subscribes to all sources, stays silent until **each** has emitted at least once, then emits an array of the latest values every time any source emits; it completes only when all sources have completed. `forkJoin([a$, b$])` also subscribes to all sources at once, but it emits **exactly once** — the last value of each — and only after **every** source has completed. So a `forkJoin` over an `interval`, a `Subject` or a long-lived state stream never emits, because that source never completes. And if any source completes **without** emitting, `forkJoin` completes immediately with no value at all. Use `forkJoin` for a batch of one-shot requests that each emit once and complete; use `combineLatest` for long-lived streams whose latest values you want to keep recombining.

code

ts · 14 lines
ts
import { forkJoin, of, defaultIfEmpty, EMPTY } from 'rxjs';

const profile$ = of({ id: 'u1', name: 'Ada' });
const avatar$ = EMPTY; // e.g. a lookup that found nothing

forkJoin({ profile: profile$, avatar: avatar$ }).subscribe({
  next: (v) => console.log('next', v),
  complete: () => console.log('complete'),
});
// complete   (no next: avatar$ completed without a value)

forkJoin({ profile: profile$, avatar: avatar$.pipe(defaultIfEmpty(null)) })
  .subscribe((v) => console.log(v));
// { profile: { id: 'u1', name: 'Ada' }, avatar: null }

go deeper

for a junior

Recall the two triggers: combineLatest needs a first value from every source, forkJoin needs every source to complete and then emits once.

for a middle

Explain the slot mechanics: latest values kept per source, emission only when all slots are full, forkJoin completing silently when a source completes empty, and the array and object call forms.

for a senior

Diagnose a join that never fires by checking which source never emits or never completes, and pick forkJoin for one-shot request batches and combineLatest for long-lived state.

for a principal

Argue where joins belong in a data layer: one-shot aggregation versus live recomposition, and how default values for empty sources become a contract every caller relies on.

## Two different waiting rules RxJS has several **static join functions** — functions imported from `'rxjs'` that take several source observables and return one output observable. `combineLatest` and `forkJoin` look alike because both take an array (or an object) of sources and both emit arrays (or objects) of values. They differ in **what they wait for**: - **`combineLatest`** waits until every source has produced a **first value**, then emits on every later value from any source. - **`forkJoin`** waits until every source has **completed**, then emits once and completes. Everything else — the famous "forkJoin never fires" bug included — follows from those two rules. ## combineLatest step by step 1. On subscription it subscribes to every source, in array order. 2. Each source's most recent value is stored in its own slot. 3. While any slot is still empty, nothing is emitted — a slot is never filled with `undefined`. 4. From the moment the last empty slot is filled, **every** new value from **any** source produces an emission containing a copy of all current slots. 5. A source that completes after emitting keeps contributing its last value; the output completes only when **all** sources have completed. 6. An error from any source is forwarded immediately and the other sources are unsubscribed. ## forkJoin step by step 1. On subscription it subscribes to every source at once. 2. It records only whether each source has emitted and what its **latest** value was. 3. When every source has completed, it emits one array of those last values and completes. 4. If any source completes **without ever emitting**, `forkJoin` completes at that moment **without emitting**, even if other sources already have values — the array could never be full. 5. If any source errors, `forkJoin` errors and unsubscribes from the rest. This is why a source that never completes — `interval(1000)`, a `Subject`, a `BehaviorSubject` holding application state, a `fromEvent` stream — makes `forkJoin` wait forever. No error is raised; the subscriber simply never hears anything. A request observable that emits one response and completes (the usual shape of an HTTP call) is the natural input for `forkJoin`. ## Side by side, with zip for contrast | | `combineLatest` | `forkJoin` | `zip` | |---|---|---|---| | First emission | every source has emitted once | every source has completed | every source has an unconsumed value | | How often | on every value, once primed | exactly once | once per complete set | | Which values | latest of each | last of each | the n-th of each, in order | | Completes when | all sources complete | all complete, or one completes empty | a completed source has no buffered values | | Fits | long-lived state, filters, form values | parallel one-shot requests | pairing items by position | `zip` is the third waiting rule: it buffers values per source and emits only when every buffer has at least one value, pairing them by position. A fast source against a slow one makes `zip`'s buffer grow without bound, which is why it is rarer in application code. ## The calling conventions in RxJS 7 - Pass an **array**: `forkJoin([a$, b$])` emits a tuple typed from the inputs. - Or pass an **object**: `forkJoin({ user: user$, prefs: prefs$ })` emits an object with the same keys, which reads better when there are more than two sources. `combineLatest` accepts both shapes too. - Passing the sources as separate arguments (`forkJoin(a$, b$)`) has been deprecated since RxJS 6.5 and is slated for removal in v8. - An empty array makes both functions complete immediately without emitting. ```ts import { combineLatest, forkJoin, interval, of, take } from 'rxjs'; const a$ = of(1, 2); // emits 1, 2 synchronously, then completes const tick$ = interval(1000); // 0, 1, 2, ... and never completes forkJoin([a$, tick$.pipe(take(3))]).subscribe((v) => console.log('forkJoin', v)); // after ~3 s: forkJoin [2, 2], then completes forkJoin([a$, tick$]).subscribe((v) => console.log('stuck', v)); // never logs: tick$ never completes combineLatest([a$, tick$.pipe(take(3))]).subscribe((v) => console.log('latest', v)); // latest [2, 0] latest [2, 1] latest [2, 2] ``` ## Choosing between them - Several independent requests whose results you need **together, once** (load a record, its lookup table and its translations before showing a form): `forkJoin`. - Several long-lived streams whose **current combination** you want to recompute whenever any of them changes (filter controls, a selected tab, a sort order): `combineLatest`. - If you reach for `forkJoin` over a state stream, either you wanted `combineLatest`, or you meant to take one value from that stream first (for example with `take(1)`) so that it completes. - If a source may legitimately complete empty, give it a default value before the join so that `forkJoin` still emits.

  • What does forkJoin emit if one of its sources completes without emitting a value?
    Nothing. The moment a source completes empty, `forkJoin` completes without emitting, even if the other sources already produced values, because that slot could never be filled. `EMPTY`, a filtered-out result or an early `take(0)` all cause it. If an empty result is legitimate, pipe that source through `defaultIfEmpty(fallback)` so it contributes a value and the join still emits.
  • How do you pass sources to combineLatest and forkJoin in RxJS 7?
    As an array, which emits a typed tuple, or as an object, which emits an object with the same keys — `forkJoin({ user: user$, prefs: prefs$ })`. Passing the observables as separate arguments has been deprecated since RxJS 6.5 and is slated for removal in v8. To reshape the result, pipe it through `map` rather than relying on a trailing result-selector function.
  • When does combineLatest complete, and does a completed source stop contributing?
    `combineLatest` completes only after every source has completed. A source that completes after emitting keeps its last value in its slot, so later emissions from the other sources still include it. That is why combining a one-shot request with a long-lived stream works: the request's value stays pinned while the other stream keeps driving emissions.

saying these in an interview costs you the question

  • forkJoin emits again every time one of its sources emits a new value
  • combineLatest emits as soon as the first of its sources emits
  • forkJoin is fine over a BehaviorSubject because it already holds a value
  • combineLatest completes as soon as any one of its sources completes
  • A source that completes empty just leaves an undefined slot in forkJoin's result
open as a page

In RxJS, what is the difference between merge() and concat() when you combine two observables into one stream?

level: juniorimportance: should knowfreq 58%

basics

~10 s

merge subscribes to every source at once and forwards values as they arrive, interleaved; concat subscribes to the next source only after the previous one completes, so values come out source by source.

open as a page

In RxJS, when would you use withLatestFrom instead of combineLatest, and why can withLatestFrom silently drop source values?

level: middleimportance: should knowfreq 48%

basics

~10 s

Use withLatestFrom when only the source stream should trigger output and other streams just supply their latest value; source values that arrive before every other stream has emitted once are discarded, not queued.

open as a page

An RxJS dashboard stream built as combineLatest([user$, settings$, permissions$]) never produces a value; how do you find the cause and fix it?

level: seniorimportance: should knowfreq 44%

basics

~20 s

combineLatest stays silent until every source has emitted once, so one of the three never emitted, completed empty or errored; trace each source with tap, then seed or default the culprit or fix its producer.

open as a page

In RxJS 7, why were the pipeable merge, concat, zip, race and combineLatest operators deprecated, and what replaces them?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

Each shared its name with a static creation function, which clashed once all operators moved into the 'rxjs' entry point; they were renamed mergeWith, concatWith, zipWith, raceWith and combineLatestWith, and the old names are deprecated.

open as a page