Part 3 · 3 chapters · ~18 min

Streams, Buffers and Backpressure

Buffer internals and pooling, typed arrays and zero-copy, the four stream types and their state machines, backpressure through highWaterMark, pipeline versus pipe, Web Streams and their interop, async iteration, writing correct Duplex and Transform streams, and stream performance.

11

Buffers, typed arrays and zero-copy

A Buffer is a Uint8Array subclass over memory outside the V8 heap (counted in process.memoryUsage().arrayBuffers, not heapUsed). Small allocations come from a shared 8 KiB pool: Buffer.allocUnsafe(100) returns a slice of the pool without zeroing it.

callzeroed?pooled?use
Buffer.alloc(n)yesnodefault; safe
Buffer.allocUnsafe(n)no (may contain old data)yes if n < 4 KiBonly when you overwrite every byte immediately
Buffer.allocUnsafeSlow(n)nonolong-lived buffers that should not pin a pool slab
Buffer.from(arrayBuffer, off, len)n/ashares memoryzero-copy views
code
const big = Buffer.alloc(1024);
const view = big.subarray(0, 16);      // zero-copy: same memory
view[0] = 1; console.log(big[0]);      // 1
// a small slice of a pooled buffer keeps the whole 8 KiB slab alive: copy if you keep it long-term
const keep = Buffer.from(view);        // copy

Zero-copy patterns: pass ArrayBuffers to workers in the transfer list instead of cloning; use subarray rather than slicing into new buffers; use fs.read into a reused buffer for parsers. SharedArrayBuffer is covered with Atomics in part 6.

12

Stream types, state machines and backpressure

typeimplementinternal state that matters
Readable_read(size), call push()flowing vs paused mode, buffer length vs highWaterMark, ended
Writable_write(chunk, enc, cb)buffered length, needDrain, corked, finished
Duplexboth, independent sidestwo separate buffers (a TCP socket)
Transform_transform(chunk, enc, cb), _flush(cb)output coupled to input (gzip, parsers)
code
// a correct Transform: NDJSON lines → objects, with backpressure handled by the base class
import { Transform } from 'node:stream';
export const ndjson = () => {
  let tail = '';
  return new Transform({
    readableObjectMode: true,
    transform(chunk, _enc, cb) {
      const lines = (tail + chunk).split('\n'); tail = lines.pop();
      try { for (const l of lines) if (l) this.push(JSON.parse(l)); cb(); } catch (e) { cb(e); }
    },
    flush(cb) { try { if (tail) this.push(JSON.parse(tail)); cb(); } catch (e) { cb(e); } },
  });
};
BACKPRESSURE
a fast reader, a slow writer, and the highWaterMark that keeps memory bounded
write()flushReadablefast diskwritable bufferhighWaterMark 16 KiBWritableslow socket
swipe the figure sideways, or tap expand for full screen
1/5
without it
Read a 2 GB file and write it to a slow client with on("data") and write() ignoring the return value: data arrives far faster than the socket drains, so it piles up in the writable buffer and memory grows toward 2 GB.
ignore write() return value → memory grows unboundedthe classic streaming bug
13

pipeline, Web Streams, async iterators and performance

code
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';

// pipeline: propagates errors, destroys every stream on failure, and resolves when all finish
await pipeline(createReadStream('big.log'), createGzip(), createWriteStream('big.log.gz'));

// pipe() does NOT forward errors or destroy upstream on failure: a leak waiting to happen
// src.pipe(gzip).pipe(dst)   // an error in dst leaves src open forever

// async iteration respects backpressure: the loop body is the consumer
for await (const obj of createReadStream('events.ndjson').pipe(ndjson())) {
  await save(obj);                                // the stream waits while you await
}

// Web Streams interop (fetch bodies are Web ReadableStreams)
import { Readable } from 'node:stream';
const res = await fetch(url);
await pipeline(Readable.fromWeb(res.body), createWriteStream('out.bin'));
performance levereffect
chunk size / highWaterMarkbigger chunks mean fewer callbacks; 64 KiB to 1 MiB for bulk file work
object modeevery object is a separate callback and allocation; batch objects in hot paths
writev and cork()/uncork()coalesce many small writes into one syscall
Web Streams in Nodeconvenient for fetch interop, currently slower than Node streams for heavy throughput