Node.js 20 🟢 Piping Streams, Backpressure Handling, and stream/promises
Piping is the mechanism that connects a readable stream to a writable stream and moves data between them. The pipe method has been part of Node.js since the earliest versions, and it handles the mechanics that make streaming practical: reading chunks from the source, writing them to the destination, pausing when the destination is full, and resuming when it drains. Backpressure is what those pauses represent: the destination is signaling that its buffer is full and the source should slow down. The stream/promises module, introduced in Node.js 15, wraps the pipeline machinery in promises so that error handling and cleanup work with async/await.
The distinction between pipe and pipeline is the most important practical detail. pipe connects two streams and handles backpressure, but it does not propagate errors from the source to the destination, does not clean up the streams when one fails, and does not return a promise. pipeline does all three. The Node.js documentation recommends pipeline over pipe for new code, and the promise version makes it a one-line await in an async function. This chapter covers the pipe method, the backpressure mechanism, the pipeline function from both node:stream and node:stream/promises, the finished utility, and the patterns that make stream processing correct under failure.
Key point: pipe connects streams and handles backpressure but does not propagate errors or clean up on failure. pipeline does both and returns a promise when imported from node:stream/promises. The promise-based pipeline is the recommended way to connect streams in async code. Backpressure is the mechanism by which a writable signals a readable to pause; pipe and pipeline handle it automatically.
Why piping and backpressure matter
The memory problem. A program that reads a file into memory and then writes it to another file uses memory proportional to the file size. A program that pipes the readable to the writable uses memory proportional to the chunk size. The difference is the difference between a program that works on a 100 MB file and one that works on a 100 GB file.
The backpressure problem. A fast readable — a file read from a local SSD — can produce data much faster than a slow writable — a network socket — can consume it. Without backpressure, the writable’s internal buffer grows without bound, and the process runs out of memory. Backpressure is the mechanism that pauses the readable when the writable’s buffer is full and resumes it when the buffer drains.
The error-propagation problem. When a stream in the middle of a chain fails, the other streams need to be told. With pipe, the error is not propagated, and the streams may leak. With pipeline, the error is propagated, and all the streams are destroyed.
The cleanup problem. A file descriptor that is not closed leaks. A socket that is not destroyed stays open. pipeline destroys every stream in the chain when the pipeline succeeds or fails, so no resources are leaked.
The composition problem. A pipeline can include any number of stages: read, decompress, transform, compress, write. Each stage is a stream, and the whole chain processes data with constant memory. The pipeline function makes the composition a single call.
a. The pipe method
The pipe method connects a readable to a writable and returns the writable so that chains can be formed.
const fs = require('node:fs');
const input = fs.createReadStream('input.txt');
const output = fs.createWriteStream('output.txt');
input.pipe(output);
The data flows from input to output. The pipe method handles the backpressure: when output.write returns false, pipe pauses input until output emits drain.
A chain of pipes:
const zlib = require('node:zlib');
fs.createReadStream('input.txt')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('input.txt.gz'));
Each pipe call returns the destination, so the next pipe can be chained. The chain reads input.txt, compresses the data, and writes the compressed bytes to input.txt.gz.
The pipe method emits events on the readable: pipe when a destination is connected, and unpipe when one is removed. A readable can have multiple destinations, and the data is copied to each.
input.pipe(output1);
input.pipe(output2);
Both outputs receive the same data. If one output is slower than the other, the readable is paused until the slower output drains, which slows the faster output. This is a limitation of pipe with multiple destinations.
The pipe method does not forward errors. If input emits an error, the error is not sent to output, and output is not closed. The program must attach error listeners to both streams and handle the cleanup manually.
b. Backpressure mechanics
Backpressure is the flow control that prevents a fast producer from overwhelming a slow consumer. The mechanism is the return value of write and the drain event.
When a chunk is written to a writable, the chunk goes into an internal buffer if it cannot be sent immediately. The buffer has a high-water mark, which is a threshold in bytes (or objects, in object mode). When the buffer size exceeds the high-water mark, write returns false.
const ok = stream.write(chunk);
// ok === true → buffer below high-water mark
// ok === false → buffer above high-water mark
When write returns false, the producer should stop writing. When the buffer drains back below the high-water mark, the drain event fires, and the producer can resume.
With pipe, this is handled automatically. The pipe implementation listens for the drain event and calls resume on the readable. The producer does not need to do anything.
When writing manually, the producer must respect the return value:
const { once } = require('node:events');
async function writeAll(stream, chunks) {
for (const chunk of chunks) {
if (!stream.write(chunk)) {
await once(stream, 'drain');
}
}
stream.end();
}
The once function returns a promise that resolves when the event fires. The loop writes each chunk, waits for the drain if the buffer is full, and ends the stream when done.
The default high-water mark is 16 KB for byte streams and 16 objects for object-mode streams. The values can be changed per stream with the highWaterMark option.
const stream = fs.createReadStream('file.txt', { highWaterMark: 64 * 1024 });
A larger high-water mark means fewer pauses but more memory. A smaller one means more pauses but less memory. The default is a reasonable balance for most workloads.
c. The pipeline function
The pipeline function connects any number of streams and handles error propagation and cleanup. It is available in two forms: callback-based from node:stream, and promise-based from node:stream/promises.
The promise-based form is the recommended one for async code:
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
const zlib = require('node:zlib');
async function compress() {
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('input.txt.gz')
);
console.log('compressed');
}
The await resolves when the pipeline finishes. If any stream in the chain emits an error, the promise rejects with that error, and all the other streams are destroyed.
The callback-based form:
const { pipeline } = require('node:stream');
pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('input.txt.gz'),
(err) => {
if (err) {
console.error('pipeline failed:', err);
} else {
console.log('compressed');
}
}
);
The callback receives the error from any stream in the chain. The promise form is preferred in async functions because it composes with try/catch.
The pipeline function accepts the following stream types at any position: Readable, Writable, Duplex, and Transform. The last argument must be a Writable or a function. The intermediate arguments must be Duplex or Transform streams, because they need both a readable and a writable side.
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGunzip(), // Transform
zlib.createGzip(), // Transform
fs.createWriteStream('output.txt.gz')
);
The chain reads a gzipped file, decompresses it, recompresses it, and writes the result. All the stages run with constant memory.
The pipeline function does not require the streams to be in a specific order beyond the first (Readable) and last (Writable). Any number of transforms can be inserted between them.
d. The finished utility
The finished function waits for a stream to finish or fail. It is useful when you need to know when a single stream has completed, not the whole pipeline.
The promise-based form from node:stream/promises:
const { finished } = require('node:stream/promises');
const stream = fs.createReadStream('input.txt');
stream.resume(); // start flowing
await finished(stream);
console.log('stream finished');
The await resolves when the stream emits end or finish and rejects if the stream emits an error. It does not consume the data; the caller must attach a data listener or call resume to start the flow.
The callback-based form from node:stream:
const { finished } = require('node:stream');
finished(stream, (err) => {
if (err) {
console.error('stream failed:', err);
} else {
console.log('stream finished');
}
});
The finished function is the tool for waiting on a single stream. For multiple streams, use pipeline.
e. Error handling in pipelines
Error handling is the reason pipeline exists. Without it, a failure in the middle of a pipe chain leaves the other streams in an undefined state.
With pipeline:
try {
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('output.txt.gz')
);
console.log('done');
} catch (err) {
console.error('pipeline failed:', err);
}
When any stream in the chain emits an error, pipeline destroys all the streams in the chain and rejects the promise with the error. The catch block handles it.
Common errors in a pipeline:
| Error | Cause |
|---|---|
ENOENT | Input file not found |
EACCES | Permission denied |
EISDIR | Input is a directory |
ERR_STREAM_PREMATURE_CLOSE | A stream closed before finishing |
ERR_STREAM_DESTROYED | A stream was destroyed |
ERR_INVALID_ARG_TYPE | Wrong argument type |
Each stream in the chain should also have an error listener when not using pipeline. With pipeline, the listeners are managed automatically.
The pipeline function emits the error to the callback or rejects the promise, and the streams are destroyed. The caller does not need to manually destroy or close anything.
Complete Example Session
// ============================================
// PART 1: BASIC PIPE
// ============================================
const fs = require('node:fs');
const input = fs.createReadStream('input.txt');
const output = fs.createWriteStream('output.txt');
input.pipe(output);
// ============================================
// PART 2: CHAINED PIPE
// ============================================
const zlib = require('node:zlib');
fs.createReadStream('input.txt')
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream('input.txt.gz'));
// ============================================
// PART 3: MULTIPLE DESTINATIONS
// ============================================
input.pipe(output1);
input.pipe(output2);
// ============================================
// PART 4: BACKPRESSURE WITH WRITE
// ============================================
const ok = stream.write(chunk);
if (!ok) {
stream.once('drain', () => {
// resume writing
});
}
// ============================================
// PART 5: MANUAL BACKPRESSURE LOOP
// ============================================
const { once } = require('node:events');
async function writeAll(stream, chunks) {
for (const chunk of chunks) {
if (!stream.write(chunk)) {
await once(stream, 'drain');
}
}
stream.end();
}
// ============================================
// PART 6: PIPELINE WITH PROMISES
// ============================================
const { pipeline } = require('node:stream/promises');
async function compress() {
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('input.txt.gz')
);
console.log('compressed');
}
// ============================================
// PART 7: PIPELINE WITH ERROR HANDLING
// ============================================
try {
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('output.txt.gz')
);
} catch (err) {
console.error('pipeline failed:', err);
}
// ============================================
// PART 8: PIPELINE WITH MULTIPLE TRANSFORMS
// ============================================
await pipeline(
fs.createReadStream('input.txt.gz'),
zlib.createGunzip(),
zlib.createGzip(),
fs.createWriteStream('output.txt.gz')
);
// ============================================
// PART 9: FINISHED UTILITY
// ============================================
const { finished } = require('node:stream/promises');
const stream = fs.createReadStream('input.txt');
stream.resume();
await finished(stream);
console.log('stream finished');
// ============================================
// PART 10: HIGH-WATER MARK
// ============================================
const stream = fs.createReadStream('file.txt', {
highWaterMark: 64 * 1024,
});
These ten parts cover a basic pipe, a chained pipe, multiple destinations, backpressure with write, a manual backpressure loop, pipeline with promises, pipeline with error handling, pipeline with multiple transforms, the finished utility, and the high-water mark option.
Quick Reference
pipe vs pipeline
| Aspect | pipe | pipeline |
|---|---|---|
| Error propagation | No | Yes |
| Cleanup on error | No | Yes |
| Return value | Destination | Promise (promise version) |
| Backpressure | Yes | Yes |
| Multiple streams | Chained | Single call |
| Recommended | Legacy | Yes |
pipeline Imports
| Form | Import | Usage |
|---|---|---|
| Promise | node:stream/promises | await pipeline(...) |
| Callback | node:stream | pipeline(..., cb) |
Backpressure Signals
| Signal | Meaning |
|---|---|
write returns false | Buffer above high-water mark |
drain event | Buffer below high-water mark |
pipe pauses source | Automatic backpressure |
pipe resumes source | Automatic after drain |
High-Water Mark Defaults
| Stream | Default |
|---|---|
| Byte streams | 16 KB |
| Object mode | 16 objects |
Common Errors
| Error | Cause |
|---|---|
ENOENT | File not found |
EACCES | Permission denied |
EISDIR | Input is a directory |
ERR_STREAM_PREMATURE_CLOSE | Closed before finishing |
ERR_STREAM_DESTROYED | Stream destroyed |
finished Utility
| Form | Import |
|---|---|
| Promise | node:stream/promises |
| Callback | node:stream |
Best Practices
✅ Do This:
// Use pipeline with promises
await pipeline(source, transform, destination);
// Handle errors with try/catch
try { await pipeline(...); } catch (err) { }
// Respect backpressure
if (!stream.write(chunk)) await once(stream, 'drain');
// Use a larger high-water mark for high-throughput pipelines
fs.createReadStream('file.txt', { highWaterMark: 64 * 1024 });
// Use finished for a single stream
await finished(stream);
❌ Don’t Do This:
// Use pipe without error handling
source.pipe(destination); // ❌ errors not propagated
// Ignore the write return value
stream.write(chunk); // ❌ buffer may overflow
// Forget to consume the stream
const stream = fs.createReadStream('file.txt');
await finished(stream); // ❌ never fires without resume or data listener
// Use pipeline without the last argument
await pipeline(source); // ❌ needs at least a source and a destination
// Use a non-stream in the middle
await pipeline(source, {}, destination); // ❌ invalid
Common Pitfalls
| Pitfall | Why It Happens | Fix |
|---|---|---|
| Process crashes on error | Used pipe without listeners | Use pipeline |
| Memory grows unbounded | Backpressure ignored | Respect write return |
finished never resolves | Stream not flowing | Call resume() or attach a data listener |
ERR_STREAM_PREMATURE_CLOSE | Source closed before end | Handle in the catch |
| Slow pipeline | High-water mark too small | Increase it |
write after end | Wrote after end() | Check the stream state |
| Errors lost in the middle | Used pipe | Use pipeline |
Real-World Examples
1. Copy a File
await pipeline(
fs.createReadStream('input.txt'),
fs.createWriteStream('output.txt')
);
2. Compress a File
await pipeline(
fs.createReadStream('input.txt'),
zlib.createGzip(),
fs.createWriteStream('input.txt.gz')
);
3. Decompress and Recompress
await pipeline(
fs.createReadStream('input.txt.gz'),
zlib.createGunzip(),
zlib.createGzip(),
fs.createWriteStream('output.txt.gz')
);
4. Stream to HTTP Response
await pipeline(
fs.createReadStream('file.txt'),
res
);
5. Stream from HTTP Request
await pipeline(
req,
fs.createWriteStream('upload.bin')
);
6. Custom Transform in Pipeline
await pipeline(
fs.createReadStream('input.txt'),
new Transform({ transform(chunk, enc, cb) {
this.push(chunk.toString().toUpperCase());
cb();
}}),
fs.createWriteStream('output.txt')
);
7. Manual Backpressure
for (const chunk of chunks) {
if (!res.write(chunk)) {
await once(res, 'drain');
}
}
res.end();
8. Finished for a Single Stream
const stream = fs.createReadStream('file.txt');
stream.resume();
await finished(stream);
9. Larger High-Water Mark
const stream = fs.createReadStream('file.txt', {
highWaterMark: 64 * 1024,
});
10. Pipeline with Error Handling
try {
await pipeline(source, transform, destination);
} catch (err) {
console.error('pipeline failed:', err.message);
}
Visual
pipe vs pipeline
┌──────────────────────────────────────────────────────────────┐
│ PIPE: │
│ source ──pipe──▶ transform ──pipe──▶ destination │
│ └── Backpressure handled │
│ └── Errors NOT propagated │
│ └── No cleanup on failure │
│ │
│ PIPELINE: │
│ pipeline(source, transform, destination) │
│ └── Backpressure handled │
│ └── Errors propagated │
│ └── All streams destroyed on failure │
└──────────────────────────────────────────────────────────────┘
Backpressure Flow
┌──────────────────────────────────────────────────────────────┐
│ source ──chunk──▶ writable │
│ │ │
│ ├── buffer below HWM → write returns true│
│ │ │
│ └── buffer above HWM → write returns false│
│ │ │
│ ▼ │
│ source pauses until the drain event fires │
│ │
│ drain event ──▶ source resumes │
└──────────────────────────────────────────────────────────────┘
Pipeline Stages
┌──────────────────────────────────────────────────────────────┐
│ await pipeline( │
│ fs.createReadStream('input.txt'), ← Readable │
│ zlib.createGunzip(), ← Transform │
│ zlib.createGzip(), ← Transform │
│ fs.createWriteStream('output.txt') ← Writable │
│ ); │
│ │
│ Any number of transforms between the Readable and Writable. │
│ Constant memory regardless of file size. │
└──────────────────────────────────────────────────────────────┘
Error Handling Comparison
┌──────────────────────────────────────────────────────────────┐
│ WITH pipe: │
│ source.pipe(dest); │
│ source.on('error', handler); │
│ dest.on('error', handler); │
│ └── Manual listener on every stream │
│ └── Manual cleanup required │
│ │
│ WITH pipeline: │
│ try { await pipeline(source, dest); } │
│ catch (err) { } │
│ └── Single catch │
│ └── All streams destroyed automatically │
└──────────────────────────────────────────────────────────────┘
Summary
| Item | Value |
|---|---|
pipe | Connects streams, handles backpressure |
pipe limitation | No error propagation, no cleanup |
pipeline | Connects streams, propagates errors, cleans up |
| Promise import | node:stream/promises |
| Callback import | node:stream |
| Backpressure signal | write returns false |
| Resume signal | drain event |
| High-water mark | 16 KB (byte streams), 16 (object mode) |
finished | Waits for a single stream |
| Last argument | Must be Writable or a function |
Key takeaways:
pipeconnects a readable to a writable and handles backpressure. It returns the destination so chains can be formed. It does not propagate errors or clean up on failure.pipelineis the recommended way to connect streams. It propagates errors from any stage, destroys all the streams when one fails, and returns a promise when imported fromnode:stream/promises.- Backpressure is the flow control between a fast producer and a slow consumer. The writable signals with
writereturningfalse, and the readable pauses until thedrainevent fires.pipeandpipelinehandle this automatically. - The high-water mark determines when backpressure kicks in. The default is 16 KB for byte streams. A larger value means fewer pauses and more memory; a smaller value means more pauses and less memory.
- The
finishedutility waits for a single stream. It resolves when the stream emitsendorfinishand rejects on error. The stream must be flowing for the promise to resolve. - Use
try/catchwithawait pipeline. The promise rejects with the error from any stage, and all the streams are destroyed automatically. - The last argument to
pipelinemust be a Writable or a function. The intermediate arguments must be Duplex or Transform streams. A Readable alone is not a valid pipeline.
Remember: Piping is the mechanism that moves data from a readable to a writable, and backpressure is the mechanism that keeps the flow at a rate the destination can handle. The pipe method handles the mechanics but leaves error propagation and cleanup to the caller. The pipeline function does both, and the promise version makes it a single await in async code. The high-water mark controls when backpressure engages. The finished utility is for single streams. Together these tools make stream processing correct under failure and efficient under load. Every file copy, compression, decompression, HTTP upload, and download is a pipeline, and the patterns in this chapter apply to all of them.
Stop using slow, ad-bloated tool sites! 🤮
🔎 Search “KandZ Tools” on Google to use many professional utilities for free.
KandZ.me is the ultimate minimalist hub for:
✅ Finance (Mortgage, Interest, Inflation)
✅ Tech (Base64, JSON, Dev Suite, IP)
✅ Health (BMI, BMR, TDEE)
✅ Productivity (Timer, Workspace, QR)
⚡️ Fast & Private
🔒 No data leaves your device
💎 100% Free
🔗 Use it now: https://tools.kandz.me
🔖 Bookmark it—you’ll need it later!