skip to content

In RxJS, how do bufferTime, bufferCount and buffer(notifier) decide when to emit a batch, and what happens to leftover values?

level: middleimportance: should knowfreq 38%

answer

  1. time, size or another stream
  2. empty arrays still emit
  3. partial buffer flushed on complete
  4. discarded on error

basics

~20 s

bufferTime emits the collected array every time span, bufferCount after N values, and buffer each time its notifier emits. On completion the partial buffer is emitted; on error it is discarded. bufferTime and buffer emit empty arrays when nothing arrived.

solid answer

~40 s

All three collect source values into an array and emit the array as one value. `bufferTime(ms)` emits every `ms` - **including empty arrays** when nothing arrived - and its signature also takes a creation interval and a `maxBufferSize` that flushes early when full. `bufferCount(n)` emits after every `n` values; a second argument `startBufferEvery` makes overlapping or skipping windows. `buffer(notifier)` emits whatever has accumulated each time the notifier emits, again even if empty. When the source **completes**, the pending partial buffer is emitted before completion (`bufferCount` emits nothing if no buffer is open). When the source **errors**, the error passes through and buffered values are lost. For batching analytics events, `bufferTime(5000, null, 50)` plus a filter for non-empty batches is a common shape.

code

ts · 22 lines
ts
import { Subject, bufferCount, bufferTime, filter, of } from 'rxjs';

interface AnalyticsEvent {
  name: string;
  at: number;
}

const events$ = new Subject<AnalyticsEvent>();

events$
  .pipe(
    bufferTime(5000, null, 50), // every 5 s, or as soon as 50 events are collected
    filter((batch) => batch.length > 0), // drop the empty arrays a quiet period produces
  )
  .subscribe((batch) => console.log('send', batch.length, 'events'));

of(1, 2, 3, 4, 5, 6, 7)
  .pipe(bufferCount(3))
  .subscribe((b) => console.log(b));
// [1, 2, 3]
// [4, 5, 6]
// [7]  <- partial buffer flushed on completion

go deeper

for a junior

Recall that the buffer operators emit arrays, triggered by time, by count or by a notifier stream.

for a middle

Explain empty arrays from bufferTime and buffer, the maxBufferSize and startBufferEvery arguments, and that completion flushes the partial batch.

for a senior

Design batching that tolerates failure: filter empty batches, cap batch size, and handle errors before the buffer so a failure does not silently drop events.

for a principal

Trade batch size and flush interval against data loss on page unload and server load, and decide where batching belongs between client and collector.

## Buffering in RxJS The **buffer** operators turn a stream of single values into a stream of **arrays**. Each variant decides differently when to close the current array and emit it. | Operator | Emits a batch when | Empty batches | On source complete | On source error | |---|---|---|---|---| | `bufferTime(span)` | `span` ms have passed | yes | emits the open buffer(s) | error; buffer lost | | `bufferTime(span, null, max)` | `span` passes **or** `max` values collected | yes, on timer | emits the open buffer(s) | error; buffer lost | | `bufferCount(n)` | `n` values collected | no | emits the partial buffer if one is open | error; buffer lost | | `buffer(notifier$)` | the notifier emits | yes | emits the current buffer | error; buffer lost | ## bufferTime `bufferTime(bufferTimeSpan, bufferCreationInterval?, maxBufferSize?, scheduler?)` runs on RxJS's async scheduler by default. - With only a span, it keeps one buffer and emits it every `span` ms, then starts the next. - With a `maxBufferSize`, a buffer that fills up is emitted immediately and a new one begins with a fresh timer, so bursts do not create huge arrays. - With a `bufferCreationInterval`, a new buffer opens every interval and each lives for `span` ms, so buffers can overlap or leave gaps. The detail people miss: the timer fires whether or not anything arrived, so a quiet stream produces **a steady trickle of empty arrays**. Sending those to a server wastes requests; follow `bufferTime` with `filter(batch => batch.length > 0)`. ## bufferCount `bufferCount(bufferSize, startBufferEvery?)` is purely count-based and has no timer, so it never emits an empty array. - `bufferCount(3)` over `1..7` emits `[1, 2, 3]`, `[4, 5, 6]` and, on completion, `[7]`. - `bufferCount(3, 1)` starts a new buffer at every value, producing sliding windows: `[1, 2, 3]`, `[2, 3, 4]`, and so on. - `bufferCount(2, 3)` starts a buffer every third value and closes it after two, skipping one value in each cycle. On completion, every buffer still open is emitted, even if it is short. ## buffer(notifier) `buffer(closingNotifier)` collects until the **notifier** emits, then emits the collected array and starts a new one. The notifier can be any `ObservableInput`, such as a stream of clicks on a "send" button or a flush signal from elsewhere in the app. - Each notifier emission produces an array, **even an empty one**. - If the notifier **completes**, buffering carries on until the source completes, and the final buffer is emitted then. - If the notifier **errors**, the error goes downstream. ## Leftovers: complete versus error The difference between the two terminal paths is the important production detail: 1. **Completion** flushes: the partial batch is emitted, then the result completes. No values are lost. 2. **Error** does not flush: the error is forwarded and anything buffered is dropped. If losing a batch is unacceptable, errors must be handled before the buffer, so the source completes instead of failing - error recovery is a separate topic. 3. **Unsubscription** also drops the current buffer; the operators release their arrays in teardown. ## Choosing a buffer - **Fixed time slices**, such as a once-per-second chart update: `bufferTime(1000)`. - **Bounded requests**, such as posting at most 50 events at a time: `bufferTime(span, null, 50)` or `bufferCount(50)`. - **User- or app-driven flushes**, such as a "sync now" button: `buffer(syncClicks$)`. - **Everything at the end** of a finite stream: `toArray()` instead of any buffer operator. ## An analytics-batching example A front end records many small events and sends them in batches. `bufferTime(5000, null, 50)` produces at most one batch per five seconds and never more than 50 events per batch; `filter` drops empty batches; the subscriber posts each batch. When the page's stream is completed deliberately - on logout, say - the last partial batch is still emitted. ## Summary Choose by trigger: time (`bufferTime`), count (`bufferCount`) or an external signal (`buffer`). Expect empty arrays from the time- and notifier-based variants, a flush on completion, and data loss on error or unsubscription.

  • In RxJS, why does bufferTime(1000) keep emitting [] when the source is quiet?
    `bufferTime` closes its buffer on a timer, not on values, so every second it emits whatever it collected - an empty array when nothing arrived - and starts a new one. Add `filter(batch => batch.length > 0)` after it if empty batches should be ignored.
  • In RxJS, what does bufferCount(3, 1) emit over the values 1, 2, 3, 4?
    It starts a new buffer at every value and emits each once it holds three: `[1, 2, 3]` and then `[2, 3, 4]`. On completion the still-open shorter buffers, `[3, 4]` and `[4]`, are emitted too, because `bufferCount` flushes every open buffer when the source completes.

saying these in an interview costs you the question

  • bufferTime only emits when at least one value has arrived.
  • When the source errors, the buffer operators emit the partial batch first.
  • bufferCount drops the final partial buffer when the source completes.
  • buffer(notifier) stops emitting batches and completes as soon as the notifier completes.
  • bufferCount(3, 1) produces non-overlapping groups of three.