Trust is earned, not given

A different perspective

2018-09-11 · Projects

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.