Node.js, part 3: streams, backpressure, and why pipeline() replaced .pipe()
Part 3from the Node.js series · 24 parts in all
Streams are Node's answer to data larger than memory, and the source of its most misdiagnosed bug: a process whose memory climbs while it copies a file. The cause is almost always missing backpressure — a fast producer writing into a slow consumer with nothing telling it to slow down. Part 3 is the mental model and the one function that gets it right.
The model: a pipe with a buffer on each side
Every stream carries a highWaterMark — how much it is willing to buffer.
.write() returns false when that buffer is full, which is the
consumer saying "stop". .pipe() honours that signal; hand-rolled
on('data') handlers do not:
// The classic leak: the source keeps reading at full speed because
// nothing ever tells it the socket is backed up.
source.on('data', (chunk) => {
const ok = destination.write(chunk);
if (!ok) source.pause(); // ... and now you must resume() too
});
Doing it correctly by hand is possible — pause on false, resume on
'drain' — which is exactly why you should not.
pipeline(): correct teardown, not just flow control
const { pipeline } = require('stream');
const { createReadStream, createWriteStream } = require('fs');
const { createGzip } = require('zlib');
// Errors from ANY stage destroy every stream in the chain, so a failure
// half-way cannot leave three file handles open.
pipeline(
createReadStream('access.log'),
createGzip(),
createWriteStream('access.log.gz'),
(err) => { if (err) console.error('copy failed', err); }
);
Quiet failure is the other thing pipeline fixes: .pipe() does not
forward errors, so a read error in the middle of a chain is simply an unhandled
'error' event — a crash, or worse, a silent truncation. Later Node versions
expose the promise form, which composes with async functions:
const { pipeline } = require('stream/promises');
await pipeline(
createReadStream('access.log'),
createGzip(),
createWriteStream('access.log.gz')
);
Object mode and the transform shape
A stream need not carry bytes. With objectMode: true each chunk is a value,
which is how a parser becomes a stream stage: read lines in, emit records out, and the whole
file never exists in memory at once. Next: worker threads, for the CPU-bound work that no
amount of async plumbing can help with.