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.
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.
| call | zeroed? | pooled? | use |
|---|---|---|---|
Buffer.alloc(n) | yes | no | default; safe |
Buffer.allocUnsafe(n) | no (may contain old data) | yes if n < 4 KiB | only when you overwrite every byte immediately |
Buffer.allocUnsafeSlow(n) | no | no | long-lived buffers that should not pin a pool slab |
Buffer.from(arrayBuffer, off, len) | n/a | shares memory | zero-copy views |
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.
Stream types, state machines and backpressure
| type | implement | internal 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 |
| Duplex | both, independent sides | two separate buffers (a TCP socket) |
| Transform | _transform(chunk, enc, cb), _flush(cb) | output coupled to input (gzip, parsers) |
// 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); } },
});
};pipeline, Web Streams, async iterators and performance
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 lever | effect |
|---|---|
| chunk size / highWaterMark | bigger chunks mean fewer callbacks; 64 KiB to 1 MiB for bulk file work |
| object mode | every 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 Node | convenient for fetch interop, currently slower than Node streams for heavy throughput |