src.pipe(dest) forwards data and nothing else. If src errors, dest is left open; if dest errors, src keeps reading into a queue nobody drains — in a server, a leaked file descriptor per request. pipeline propagates errors in both directions and destroys every stage when one fails or the caller aborts; node:stream/promises returns a promise instead of taking a callback. Any stage may be a stream, an async generator, or — last — an async function consuming the iterable.
import { pipeline } from 'node:stream/promises';
let read = 0, kept = 0;
async function* onlyApac(src) {
for await (const l of src) { read++; if (l.includes(',apac,')) { kept++; yield l + '\n'; } }
}
await pipeline(
createReadStream('orders.csv'),
new LineSplitter(), // from the previous subsection
onlyApac,
createGzip({ level: 6 }),
createWriteStream('apac.csv.gz'),
);
const mb = (f) => (statSync(f).size / 1048576).toFixed(1);
console.log(`${read} lines read, ${kept} kept, ${mb('apac.csv.gz')} MB written, ` +
`peak rss ${(process.memoryUsage().rss / 1048576).toFixed(0)} MB`);6000001 lines read, 1500000 kept, 4.5 MB written, peak rss 157 MB
Five stages, 108 MB in, 157 MB resident — dominated by V8 86,723 's heap for six million short strings, not by stream buffers. Speed depends on the work per item: this per-line version took 35 s, while an equivalent pipeline whose transform scanned whole 64 KB chunks and never materialized a line object finished in 1.6 s. Pay for object mode where the records matter, not on the hot path of a byte pump.