On this page

Concepts

History

All data in this API is represented as Uint8Array bytes. Strings are automatically UTF-8 encoded when passed to from(), push(), or pipeTo(). This removes ambiguity around encodings and enables zero-copy transfers between streams and native code.

Each iteration yields a batch -- an Array of Uint8Array chunks (Uint8Array[]). Batching amortizes the cost of await and Promise creation across multiple chunks. A consumer that processes one chunk at a time can simply iterate the inner array:

for await (const batch of source) {
  for (const chunk of batch) {
    handle(chunk);
  }
}
async function run() {
  for await (const batch of source) {
    for (const chunk of batch) {
      handle(chunk);
    }
  }
}

Transforms come in two forms:

  • Stateless -- a function (chunks, options) => result called once per batch. Receives Uint8Array[] (or null as the flush signal) and an options object. Returns Uint8Array[] | null | Iterable.

  • Stateful -- an object { transform(source, options) } where transform is a generator (sync or async) that receives the entire upstream iterable and an options object, and yields output. This form is used for compression, encryption, and any transform that needs to buffer across batches.

Both forms receive an options parameter with the following property:

Attributes
options.signal:AbortSignal
An AbortSignal that fires when the pipeline is cancelled, encounters an error, or the consumer stops reading. Transforms can check signal.aborted or listen for the 'abort' event to perform early cleanup.

The flush signal (null) is sent after the source ends, giving transforms a chance to emit trailing data (e.g., compression footers).

// Stateless: uppercase transform
const upper = (chunks) => {
  if (chunks === null) return null; // flush
  return chunks.map((c) => new TextEncoder().encode(
    new TextDecoder().decode(c).toUpperCase(),
  ));
};

// Stateful: line splitter
const lines = {
  transform: async function*(source) {
    let partial = '';
    for await (const chunks of source) {
      if (chunks === null) {
        if (partial) yield [new TextEncoder().encode(partial)];
        continue;
      }
      for (const chunk of chunks) {
        const str = partial + new TextDecoder().decode(chunk);
        const parts = str.split('\n');
        partial = parts.pop();
        for (const line of parts) {
          yield [new TextEncoder().encode(`${line}\n`)];
        }
      }
    }
  },
};

The API supports two models:

  • Pull -- data flows on demand. pull() and pullSync() create lazy pipelines that only read from the source when the consumer iterates.

  • Push -- data is written explicitly. push() creates a writer/readable pair with backpressure. The writer pushes data in; the readable is consumed as an async iterable.

Pull streams have natural backpressure -- the consumer drives the pace, so the source is never read faster than the consumer can process. Push streams need explicit backpressure because the producer and consumer run independently. The budget and backpressure options on push(), broadcast(), and share() control how this works.

Push streams use a two-part buffering system. Think of it like a bucket (buffer) being filled through a hose (pending writes), with a float valve that closes when the bucket is full:

                          budget (e.g., 16384)
                                 |
    Producer                     v
       |                    +---------+
       v                    |         |
  [ write() ] ----+    +--->| buffer  |---> Consumer pulls
  [ write() ]     |    |    | (bucket)|     for await (...)
  [ write() ]     v    |    +---------+
              +--------+         ^
              | pending|         |
              | writes |    float valve
              | (hose) |    (backpressure)
              +--------+
                   ^
                   |
          'strict' mode limits this too!
  • Buffer (the bucket) -- data ready for the consumer, capped at budget bytes. When the consumer pulls, it drains all buffered data at once into a single batch.

  • Pending writes (the hose) -- writes waiting for buffer space. After the consumer drains, pending writes are promoted into the now-empty buffer and their promises settle.

How each policy uses these buffers:

PolicyBuffer limitPending writes limit
'strict'budget1
'unbounded'budgetUnbounded
'drop-oldest'budgetN/A (never waits)
'drop-newest'budgetN/A (never waits)

Strict mode catches "fire-and-forget" patterns where the producer calls write() without awaiting, which would cause unbounded memory growth. It limits the buffer to budget bytes and the pending writes queue to a single entry.

If you properly await each write, you can only ever have one pending write at a time (yours), so you never hit the pending writes limit. Unawaited writes accumulate in the pending queue and throw once it overflows:

import { push, text } from 'node:stream/iter';

const { writer, readable } = push({ budget: 16384 });

// Consumer must run concurrently -- without it, the first write
// that fills the buffer blocks the producer forever.
const consuming = text(readable);

// GOOD: awaited writes. The producer waits for the consumer to
// make room when the buffer is full.
for (const item of dataset) {
  await writer.write(item);
}
await writer.end();
console.log(await consuming);
const { push, text } = require('node:stream/iter');

async function run() {
  const { writer, readable } = push({ budget: 16384 });

  // Consumer must run concurrently -- without it, the first write
  // that fills the buffer blocks the producer forever.
  const consuming = text(readable);

  // GOOD: awaited writes. The producer waits for the consumer to
  // make room when the buffer is full.
  for (const item of dataset) {
    await writer.write(item);
  }
  await writer.end();
  console.log(await consuming);
}

run().catch(console.error);

Forgetting to await will eventually throw:

// BAD: fire-and-forget. Strict mode throws once both buffers fill.
for (const item of dataset) {
  writer.write(item); // Not awaited -- queues without bound
}
// --> throws "Backpressure violation: too many pending writes"

Unbounded mode caps buffered bytes at budget but places no limit on the pending writes queue. Awaited writes block until the consumer makes room, just like strict mode. The difference is that unawaited writes silently queue forever instead of throwing -- a potential memory leak if the producer forgets to await.

This is the mode that existing Node.js classic streams and Web Streams default to. Use it when you control the producer and know it awaits properly, or when migrating code from those APIs.

import { push, text } from 'node:stream/iter';

const { writer, readable } = push({
  budget: 16384,
  backpressure: 'unbounded',
});

const consuming = text(readable);

// Safe -- awaited writes block until the consumer reads.
for (const item of dataset) {
  await writer.write(item);
}
await writer.end();
console.log(await consuming);
const { push, text } = require('node:stream/iter');

async function run() {
  const { writer, readable } = push({
    budget: 16384,
    backpressure: 'unbounded',
  });

  const consuming = text(readable);

  // Safe -- awaited writes block until the consumer reads.
  for (const item of dataset) {
    await writer.write(item);
  }
  await writer.end();
  console.log(await consuming);
}

run().catch(console.error);

Writes never wait. When the slots buffer is full, the oldest buffered chunk is evicted to make room for the incoming write. The consumer always sees the most recent data. Useful for live feeds, telemetry, or any scenario where stale data is less valuable than current data.

import { push } from 'node:stream/iter';

// Keep only the most recent ~16 KB of readings
const { writer, readable } = push({
  budget: 16384,
  backpressure: 'drop-oldest',
});
const { push } = require('node:stream/iter');

// Keep only the most recent ~16 KB of readings
const { writer, readable } = push({
  budget: 16384,
  backpressure: 'drop-oldest',
});

Writes never wait. When the slots buffer is full, the incoming write is silently discarded. The consumer processes what is already buffered without being overwhelmed by new data. Useful for rate-limiting or shedding load under pressure.

import { push } from 'node:stream/iter';

// Accept up to 16 KB of buffered data; discard anything beyond that
const { writer, readable } = push({
  budget: 16384,
  backpressure: 'drop-newest',
});
const { push } = require('node:stream/iter');

// Accept up to 16 KB of buffered data; discard anything beyond that
const { writer, readable } = push({
  budget: 16384,
  backpressure: 'drop-newest',
});

A writer is any object conforming to the Writer interface. Only write() is required; all other methods are optional.

Each async method has a synchronous *Sync counterpart designed for a try-fallback pattern: attempt the fast synchronous path first, and fall back to the async version only when the synchronous call indicates it could not complete:

if (!writer.writeSync(chunk)) await writer.write(chunk);
if (!writer.writevSync(chunks)) await writer.writev(chunks);
if (writer.endSync() < 0) await writer.end();
writer.fail(err);  // Always synchronous, no fallback needed
Attributes

Returns true if the next write is likely to be accepted (buffered data is below capacity), false if backpressure is active, or null if the writer is closed or the consumer has disconnected.

This is a hint, not a guarantee: the state can change between the check and the write. Use ondrain() to wait for capacity rather than polling.

writer.end(options?): Promise
Attributes
options:Object
Cancel just this operation. The signal cancels only the pending end() call; it does not fail the writer itself.
Returns:Promise
Fulfills with the total number of bytes written.

Signals that no more data will be written and waits for buffered data to drain.

writer.endSync(): number
Returns:number
Total bytes written, or -1 if ending cannot complete synchronously.

Synchronous variant of writer.end(). A return value of -1 means closing has started but requires asynchronous draining. Use the try-fallback pattern to await completion:

const result = writer.endSync();
if (result < 0) {
  writer.end();
}
writer.fail(reason): void
Attributes
reason:any

Put the writer into a terminal error state. If the writer is already closed or errored, this is a no-op. Unlike write() and end(), fail() is unconditionally synchronous because failing a writer is a pure state transition with no async work to perform.

write(chunk, options?): Promise
Attributes
options:Object
Cancel just this write operation. The signal cancels only the pending write() call; it does not fail the writer itself.
Returns:Promise
Fulfills with undefined when buffer space is available.

Write a chunk.

writer.writeSync(chunk): boolean
Attributes
Returns:boolean
true if the write was accepted, false if the buffer is full.

Synchronous write. Does not block; returns false if backpressure is active.

writer.writev(chunks, options?): Promise
Attributes
chunks:Uint8Array[] | string[]
options:Object
Cancel just this write operation. The signal cancels only the pending writev() call; it does not fail the writer itself.
Returns:Promise

Write multiple chunks as a single batch.

writer.writevSync(chunks): boolean
Attributes
chunks:Uint8Array[] | string[]
Returns:boolean
true if the write was accepted, false if the buffer is full.

Synchronous batch write.