skip to content

In RxJS, a chat room's messages$ must show late joiners recent history; how do ReplaySubject's bufferSize and windowTime decide what they get?

level: middleimportance: should knowfreq 50%

answer

  1. count limit and age limit
  2. both default to unlimited
  3. replayed in order, synchronously
  4. replays even after completion
  5. AsyncSubject keeps only the last

basics

~10 s

ReplaySubject(bufferSize, windowTime) keeps at most bufferSize values that are younger than windowTime milliseconds and replays them in order to each new subscriber before live values; both limits default to infinity.

solid answer

~40 s

`new ReplaySubject<Message>(50, 10 * 60_000)` stores every pushed message together with a timestamp and trims the buffer to at most 50 entries and to entries younger than ten minutes. A new subscriber synchronously receives whatever survives both limits, oldest first, and then live messages. Both parameters default to `Infinity`, so `new ReplaySubject()` keeps **every** message forever — a memory leak in a long-lived room. A `bufferSize` below one is raised to one, so `ReplaySubject(0)` still replays the latest value. Unlike a `BehaviorSubject`, it needs no seed and keeps replaying its buffer even after `complete()` or `error()`, followed by that notification. An `AsyncSubject` would be wrong here: it emits only its final value, and only when it completes.

code

ts · 10 lines
ts
import { AsyncSubject } from 'rxjs';

const result = new AsyncSubject<number>();
result.subscribe((v) => console.log('early', v));

result.next(1);
result.next(2);      // nothing logged yet
result.complete();   // early 2

result.subscribe((v) => console.log('late', v)); // late 2

go deeper

for a junior

Recall that a ReplaySubject replays past values to new subscribers, limited by a count and optionally by age.

for a middle

Explain bufferSize and windowTime together, the Infinity defaults, synchronous oldest-first replay, behaviour after completion, and how AsyncSubject differs.

for a senior

Size replay buffers for memory and for the burst a late subscriber receives, and notice that age is measured from next(), not from the data's own timestamps.

for a principal

Decide whether history belongs in an in-memory replay buffer at all or should come from a paged server query, weighing memory, consistency and reconnect behaviour.

## The late-joiner problem A chat room pushes each incoming message into a subject. With a plain `Subject`, someone who opens the room after the conversation started sees an empty window until the next message arrives. A `BehaviorSubject` would show exactly one message and needs an invented initial value. What the room actually wants is **recent history**: the last N messages, and not ones from last week. That is `ReplaySubject`. ## How ReplaySubject works `ReplaySubject<T>` is a `Subject` with a **buffer**: 1. Every `next(v)` appends `v` to the buffer (with a timestamp when a time window is set) and forwards it to current subscribers. 2. The buffer is **trimmed** to the configured limits. 3. When someone subscribes, the subject first emits every buffered value **synchronously, oldest first**, then forwards live values. 4. If the subject has completed or errored, the new subscriber gets the buffered values **and then** the completion or error. The constructor takes up to three arguments: - **`bufferSize`** — the maximum number of values kept. Default `Infinity`. Values below one are raised to one. - **`windowTime`** — how long, in milliseconds, a value stays eligible after it was pushed. Default `Infinity`. - **`timestampProvider`** — an object with a `now()` method used to measure age; by default it reads the system clock. Schedulers satisfy this, which is how virtual time works in tests. When both limits are set, **both** apply: a value is replayed only if it is among the last `bufferSize` values **and** younger than `windowTime`. ```ts import { ReplaySubject } from 'rxjs'; interface Message { from: string; text: string } // last 50 messages, and none older than 10 minutes const messages$ = new ReplaySubject<Message>(50, 10 * 60_000); messages$.next({ from: 'ana', text: 'hi' }); messages$.next({ from: 'bo', text: 'hello' }); // someone opens the room later messages$.subscribe((m) => console.log(`${m.from}: ${m.text}`)); // ana: hi // bo: hello (replayed synchronously), then live messages ``` ## The four subjects compared | Subject | Seed needed | Late subscriber receives | After completion | |---|---|---|---| | `Subject` | no | only future values | completion only | | `BehaviorSubject` | yes | current value, then future ones | completion only | | `ReplaySubject(n, t)` | no | buffered values within `n` and `t`, then future ones | buffered values, then completion | | `AsyncSubject` | no | nothing until completion | the last value, then completion | `AsyncSubject` is the odd one out: it forwards **nothing** while active, and on `complete()` emits only its final value (if it had one) followed by completion — to current and later subscribers alike. It fits "one result computed once", not a chat feed. ## Traps worth naming - **Unbounded by default.** `new ReplaySubject()` keeps every value for the life of the subject. A room open all day accumulates every message in memory, and every new joiner receives all of them at once. - **Burst on subscribe.** Replay is synchronous: a joiner with a 500-message buffer receives 500 emissions before its subscribe call returns. Size the buffer for what the view can render. - **Age is measured from `next()`, not from the message's own timestamp.** Messages loaded from a server with old timestamps are "fresh" to the buffer the moment they are pushed. - **Replay after completion.** When the room closes and you call `complete()`, later subscribers still receive the buffered history, then the completion. That is often desired (read-only history) but surprises people who expect a completed subject to be empty. - **`ReplaySubject(0)` is not "no replay".** The buffer size is clamped to at least one, so it still replays the most recent value; use a plain `Subject` for no replay. ## Where this sits among the sharing tools A `ReplaySubject` is also the connector inside `shareReplay`, which puts the same buffer between a source and its subscribers. Choose a hand-held `ReplaySubject` when your own code pushes the values (as the chat socket handler does); choose `shareReplay` when an existing observable should be run once and replayed.

  • What does a subscriber get from a ReplaySubject that has already completed?
    Every value still in the buffer, synchronously and oldest first, followed by the completion notification. The same applies after `error()`: buffered values first, then the error. This differs from `BehaviorSubject`, which sends only the terminal notification once it has stopped.
  • Why is new ReplaySubject(0) not a way to disable replay?
    The constructor raises any `bufferSize` below one to one, so `ReplaySubject(0)` still keeps and replays the most recent value. If late subscribers should receive nothing from the past, use a plain `Subject`.

saying these in an interview costs you the question

  • new ReplaySubject() with no arguments replays only the latest value
  • windowTime delays delivery to late subscribers instead of limiting the age of replayed values
  • A completed ReplaySubject sends new subscribers only the completion
  • ReplaySubject needs an initial value, just like BehaviorSubject
  • AsyncSubject emits each value as soon as it is pushed