| |

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:

ErrorCause
ENOENTInput file not found
EACCESPermission denied
EISDIRInput is a directory
ERR_STREAM_PREMATURE_CLOSEA stream closed before finishing
ERR_STREAM_DESTROYEDA stream was destroyed
ERR_INVALID_ARG_TYPEWrong 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

Aspectpipepipeline
Error propagationNoYes
Cleanup on errorNoYes
Return valueDestinationPromise (promise version)
BackpressureYesYes
Multiple streamsChainedSingle call
RecommendedLegacyYes

pipeline Imports

FormImportUsage
Promisenode:stream/promisesawait pipeline(...)
Callbacknode:streampipeline(..., cb)

Backpressure Signals

SignalMeaning
write returns falseBuffer above high-water mark
drain eventBuffer below high-water mark
pipe pauses sourceAutomatic backpressure
pipe resumes sourceAutomatic after drain

High-Water Mark Defaults

StreamDefault
Byte streams16 KB
Object mode16 objects

Common Errors

ErrorCause
ENOENTFile not found
EACCESPermission denied
EISDIRInput is a directory
ERR_STREAM_PREMATURE_CLOSEClosed before finishing
ERR_STREAM_DESTROYEDStream destroyed

finished Utility

FormImport
Promisenode:stream/promises
Callbacknode: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

PitfallWhy It HappensFix
Process crashes on errorUsed pipe without listenersUse pipeline
Memory grows unboundedBackpressure ignoredRespect write return
finished never resolvesStream not flowingCall resume() or attach a data listener
ERR_STREAM_PREMATURE_CLOSESource closed before endHandle in the catch
Slow pipelineHigh-water mark too smallIncrease it
write after endWrote after end()Check the stream state
Errors lost in the middleUsed pipeUse 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

ItemValue
pipeConnects streams, handles backpressure
pipe limitationNo error propagation, no cleanup
pipelineConnects streams, propagates errors, cleans up
Promise importnode:stream/promises
Callback importnode:stream
Backpressure signalwrite returns false
Resume signaldrain event
High-water mark16 KB (byte streams), 16 (object mode)
finishedWaits for a single stream
Last argumentMust be Writable or a function

Key takeaways:

  • pipe connects 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.
  • pipeline is 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 from node:stream/promises.
  • Backpressure is the flow control between a fast producer and a slow consumer. The writable signals with write returning false, and the readable pauses until the drain event fires. pipe and pipeline handle 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 finished utility waits for a single stream. It resolves when the stream emits end or finish and rejects on error. The stream must be flowing for the promise to resolve.
  • Use try/catch with await pipeline. The promise rejects with the error from any stage, and all the streams are destroyed automatically.
  • The last argument to pipeline must 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!