| |

Node.js 19 🟢 Stream Fundamentals (Readable, Writable, Duplex, Transform)

A stream is an abstraction for data that arrives or leaves in chunks over time. Instead of loading an entire file into memory, a stream reads it piece by piece. Instead of building a complete response before sending it, a stream writes it as the data becomes available. Node.js provides four stream types: Readable, Writable, Duplex, and Transform. Each one models a different direction of data flow, and together they are the foundation of every I/O operation in Node.js. The fs.createReadStream, the HTTP request and response objects, the net.Socket, and the zlib compression functions are all streams.

The stream abstraction solves the memory problem. Reading a 2 GB file with fs.readFile would load 2 GB into memory, which is impossible on most machines. Reading it with fs.createReadStream processes it in 64 KB chunks, and the memory usage stays constant regardless of file size. The same applies to network data, which arrives in packets and cannot be buffered entirely. Streams also solve the latency problem: a stream can start processing data before the source has finished producing it, which is what makes streaming video, log tailing, and progressive downloads possible.

This chapter covers the four stream types, the events they emit, the pipe method, backpressure, the pipeline function, object mode, the stream lifecycle, and the patterns that make stream processing correct and efficient.

Key point: A Readable stream produces data, a Writable stream consumes it, a Duplex stream does both independently, and a Transform stream is a Duplex that modifies data as it passes through. Streams emit events (data, end, error, finish, close) and can be connected with pipe or, preferably, pipeline. Backpressure is the mechanism that prevents a fast producer from overwhelming a slow consumer.


Why streams exist

The memory problem. Reading an entire file into memory does not scale. A 10 GB log file cannot be loaded into a 4 GB machine. A stream reads it in chunks, so the memory usage is bounded by the chunk size, not the file size. The same applies to HTTP responses, database result sets, and any other data source of unknown or unbounded size.

The latency problem. A program that waits for the entire input before producing output has a latency equal to the time it takes to read everything. A stream can process each chunk as it arrives, so the output starts before the input is finished. This is what makes video streaming, log tailing, and progressive rendering possible.

The composition problem. Streams compose. A readable can be piped to a transform, which can be piped to another transform, which can be piped to a writable. Each stage processes the data independently, and the whole pipeline runs with constant memory. Without streams, each stage would need to buffer the entire dataset.

The backpressure problem. A fast producer — a file read — can produce data faster than a slow consumer — a network write — can accept it. Without backpressure, the internal buffers grow unbounded until the process runs out of memory. Streams implement backpressure: when the writable’s buffer is full, it signals the readable to pause until the buffer drains.

The uniformity problem. Files, sockets, HTTP bodies, compression, encryption, and process I/O all expose the same stream interface. A program that processes a stream does not need to know whether the data comes from a file, a network connection, or a generator. The stream is the common abstraction.


a. Readable streams

A Readable stream is a source of data. It produces chunks that consumers can read. The fs.createReadStream is the canonical example.

const fs = require('node:fs');

const stream = fs.createReadStream('large.txt', { encoding: 'utf8' });

stream.on('data', (chunk) => {
  process.stdout.write('.');
});

stream.on('end', () => {
  console.log('\ndone');
});

stream.on('error', (err) => {
  console.error('error:', err);
});

The data event fires for each chunk. The end event fires when the stream has no more data. The error event fires on failure, and it is mandatory: a stream that errors without a listener crashes the process.

A Readable has two modes. In flowing mode, the data event fires as data arrives. In paused mode, data is read explicitly with the read() method. Attaching a data listener switches the stream to flowing mode. Using pipe also switches to flowing mode. Using read() keeps it paused.

stream.pause();   // switch to paused mode
stream.resume();  // switch back to flowing mode

The readable event fires when data is available to read, which is the signal for the paused-mode consumer to call read().

stream.on('readable', () => {
  let chunk;
  while ((chunk = stream.read()) !== null) {
    process.stdout.write('.');
  }
});

Most code uses flowing mode or pipe and never touches read() directly. The paused mode is for cases where the consumer needs precise control.


b. Writable streams

A Writable stream is a destination for data. It consumes chunks that producers write. The fs.createWriteStream is the canonical example.

const fs = require('node:fs');

const stream = fs.createWriteStream('output.txt');

stream.write('first line\n');
stream.write('second line\n');
stream.end('last line\n');

stream.on('finish', () => {
  console.log('written');
});

stream.on('error', (err) => {
  console.error('error:', err);
});

The write method sends a chunk. It returns true if the internal buffer has room and false if it is full. The end method signals that no more data will be written. The finish event fires when all data has been flushed.

The return value of write is the backpressure signal. When it returns false, the producer should stop writing until the drain event fires.

const ok = stream.write(chunk);
if (!ok) {
  stream.once('drain', () => {
    // resume writing
  });
}

When pipe is used, the backpressure is handled automatically. The readable is paused when the writable’s buffer is full and resumed when the drain event fires.

The error event is mandatory. A write error without a listener crashes the process.


c. Duplex and Transform streams

A Duplex stream is both Readable and Writable, with the two sides independent. The net.Socket is the canonical example: it receives data from the network and sends data to the network, and the two flows do not interact.

const net = require('node:net');

const socket = net.connect(8080, 'example.com');

socket.write('GET / HTTP/1.1\r\nHost: example.com\r\n\r\n');

socket.on('data', (chunk) => {
  process.stdout.write(chunk);
});

socket.on('end', () => {
  console.log('response complete');
});

The write sends data to the server, and the data event receives data from the server. The two directions are independent.

A Transform stream is a Duplex where the two sides are connected: data written to the writable side is transformed and appears on the readable side. The zlib.createGzip is the canonical example.

const zlib = require('node:zlib');
const fs = require('node:fs');

const gzip = zlib.createGzip();
const input = fs.createReadStream('input.txt');
const output = fs.createWriteStream('input.txt.gz');

input.pipe(gzip).pipe(output);

The gzip stream receives uncompressed data on its writable side and produces compressed data on its readable side. The pipe chain connects the source to the transform to the destination.

The difference between Duplex and Transform is the relationship between the two sides. In a Duplex, they are independent. In a Transform, the output is a function of the input.

TypeReadableWritableRelationship
ReadableYesNo—
WritableNoYes—
DuplexYesYesIndependent
TransformYesYesOutput derived from input

d. pipe and pipeline

The pipe method connects a readable to a writable. It handles backpressure automatically: when the writable’s buffer is full, the readable is paused; when the buffer drains, the readable is resumed.

input.pipe(output);

The method can be chained:

input.pipe(gzip).pipe(output);

The problem with pipe is error handling. When an error occurs in the middle of the chain, pipe does not propagate it to the other streams, and the streams may not be cleaned up. The pipeline function from node:stream/promises solves this.

const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
const zlib = require('node:zlib');

await pipeline(
  fs.createReadStream('input.txt'),
  zlib.createGzip(),
  fs.createWriteStream('input.txt.gz')
);

console.log('compressed');

The pipeline function propagates errors, cleans up all the streams if any step fails, and returns a promise when given the promise-based version. It is the recommended way to connect streams.

The callback version is also available:

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 promise version is preferred in async functions.


e. Backpressure

Backpressure is the mechanism that prevents a fast producer from overwhelming a slow consumer. Without it, the internal buffer of the writable would grow unbounded.

The mechanism is the return value of write. When the writable’s internal buffer is below the high-water mark, write returns true. When the buffer exceeds the high-water mark, write returns false, and the producer should stop writing until the drain event fires.

TypeHigh-water markSignal
ReadablereadableHighWaterMarkreadable event
WritablewritableHighWaterMarkdrain event

The default high-water mark is 16 KB for byte streams and 16 objects for object-mode streams.

With pipe and pipeline, the backpressure is handled automatically. The readable is paused when the writable returns false and resumed when the drain event fires. The producer does not need to manage the buffer.

When writing manually, the producer must respect the return value:

async function writeAll(stream, chunks) {
  for (const chunk of chunks) {
    if (!stream.write(chunk)) {
      await once(stream, 'drain');
    }
  }
  stream.end();
}

The once function from node:events 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.


f. The stream lifecycle

A stream goes through a defined lifecycle. The events that fire depend on the type.

EventReadableWritableWhen
dataYesNoA chunk is available
endYesNoNo more data
readableYesNoData is available to read
finishNoYesAll data has been flushed
drainNoYesThe buffer has room again
closeYesYesThe underlying resource is closed
errorYesYesAn error occurred
pipeYesNopipe connected the stream
unpipeYesNopipe disconnected the stream

The end event signals the end of the readable side. The finish event signals the end of the writable side. A Transform stream emits both: finish when the input side is done, and end when the output side is done.

The close event fires when the underlying resource — a file descriptor, a socket — is closed. It may fire after end or finish.

The error event can fire at any point. A stream that errors is no longer usable for reading or writing. The pipeline function destroys all the streams in the chain when an error occurs.


Complete Example Session

// ============================================
// PART 1: READABLE STREAM
// ============================================
const fs = require('node:fs');

const stream = fs.createReadStream('large.txt', { encoding: 'utf8' });

stream.on('data', (chunk) => process.stdout.write('.'));
stream.on('end', () => console.log('\ndone'));
stream.on('error', (err) => console.error(err));
// ============================================
// PART 2: WRITABLE STREAM
// ============================================
const out = fs.createWriteStream('output.txt');

out.write('first line\n');
out.write('second line\n');
out.end('last line\n');

out.on('finish', () => console.log('written'));
out.on('error', (err) => console.error(err));
// ============================================
// PART 3: PIPE
// ============================================
const input = fs.createReadStream('input.txt');
const output = fs.createWriteStream('output.txt');

input.pipe(output);
// ============================================
// PART 4: PIPE WITH TRANSFORM
// ============================================
const zlib = require('node:zlib');

fs.createReadStream('input.txt')
  .pipe(zlib.createGzip())
  .pipe(fs.createWriteStream('input.txt.gz'));
// ============================================
// PART 5: 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');
}

compress();
// ============================================
// PART 6: BACKPRESSURE WITH WRITE
// ============================================
const ok = out.write(chunk);
if (!ok) {
  out.once('drain', () => {
    // resume writing
  });
}
// ============================================
// PART 7: MANUAL BACKPRESSURE WITH ONCE
// ============================================
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 8: DUPLEX STREAM (NET SOCKET)
// ============================================
const net = require('node:net');

const socket = net.connect(8080, 'example.com');
socket.write('GET / HTTP/1.1\r\nHost: example.com\r\n\r\n');
socket.on('data', (chunk) => process.stdout.write(chunk));
socket.on('end', () => console.log('done'));
// ============================================
// PART 9: TRANSFORM STREAM WITH CUSTOM LOGIC
// ============================================
const { Transform } = require('node:stream');

const upper = new Transform({
  transform(chunk, encoding, callback) {
    this.push(chunk.toString().toUpperCase());
    callback();
  },
});

fs.createReadStream('input.txt')
  .pipe(upper)
  .pipe(fs.createWriteStream('output.txt'));
// ============================================
// PART 10: PIPELINE WITH CUSTOM TRANSFORM
// ============================================
async function processFile() {
  await pipeline(
    fs.createReadStream('input.txt'),
    upper,
    fs.createWriteStream('output.txt')
  );
  console.log('processed');
}

These ten parts cover the Readable stream, the Writable stream, pipe, pipe with a transform, pipeline with promises, backpressure with the write return value, manual backpressure with once, a Duplex stream, a custom Transform stream, and pipeline with a custom transform.


Quick Reference

Stream Types

TypeReadableWritableExample
ReadableYesNofs.createReadStream
WritableNoYesfs.createWriteStream
DuplexYesYesnet.Socket
TransformYesYeszlib.createGzip

Common Events

EventTypeMeaning
dataReadableChunk available
endReadableNo more data
readableReadableData ready to read
finishWritableAll data flushed
drainWritableBuffer has room
closeBothResource closed
errorBothError occurred

Methods

MethodPurpose
read()Read from a paused Readable
write(chunk)Write to a Writable
end([chunk])Signal the end of writes
pipe(dest)Connect Readable to Writable
pause()Pause a flowing Readable
resume()Resume a paused Readable
destroy()Destroy the stream

pipeline vs pipe

Aspectpipepipeline
Error propagationNoYes
Cleanup on errorNoYes
Return valueDestinationPromise (promise version)
RecommendedLegacyYes

High-Water Marks

StreamDefault
Byte streams16 KB
Object mode16 objects

Backpressure Signals

TypeSignal
Writable fullwrite returns false
Writable draineddrain event
Readable has datareadable event

Best Practices

✅ Do This:

// Use pipeline instead of pipe
await pipeline(source, transform, destination);

// Handle the error event
stream.on('error', (err) => console.error(err));

// Respect backpressure
if (!stream.write(chunk)) {
  await once(stream, 'drain');
}

// Use streams for large files
fs.createReadStream('large.txt').pipe(destination);

// Use Transform for data processing
const upper = new Transform({ transform(chunk, enc, cb) { ... } });

❌ Don’t Do This:

// Use pipe without error handling
source.pipe(destination);  // ❌ errors not propagated

// Ignore the write return value
stream.write(chunk);  // ❌ may overflow the buffer

// Read a large file into memory
const data = await readFile('huge.log');  // ❌ out of memory

// Forget the error listener
stream.on('data', handler);  // ❌ crash on error

// Call write after end
stream.end();
stream.write('more');  // ❌ error

Common Pitfalls

PitfallWhy It HappensFix
Process crashes on errorNo error listenerAdd one
Memory grows unboundedBackpressure ignoredRespect write return
write after endWrote after end()Check the stream state
Errors not propagatedUsed pipeUse pipeline
Slow consumerProducer ignores drainWait for the drain event
Partial file writtenPipeline error not handledUse pipeline with a catch
Stream never endsSource does not call endEnsure the source finishes

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. Stream an HTTP Response

res.writeHead(200, { 'Content-Type': 'text/plain' });
fs.createReadStream('file.txt').pipe(res);

4. Read a Large Log Line by Line

const readline = require('node:readline');
const rl = readline.createInterface({
  input: fs.createReadStream('large.log'),
});
for await (const line of rl) {
  if (line.includes('ERROR')) console.log(line);
}

5. Transform Uppercase

const upper = new Transform({
  transform(chunk, encoding, callback) {
    this.push(chunk.toString().toUpperCase());
    callback();
  },
});

6. Stream a Request Body

app.post('/upload', (req, res) => {
  req.pipe(fs.createWriteStream('upload.bin'));
  req.on('end', () => res.send('uploaded'));
});

7. Parse a CSV Stream

const { parse } = require('csv-parse');
await pipeline(
  fs.createReadStream('data.csv'),
  parse({ columns: true }),
  async function* (source) {
    for await (const record of source) {
      yield process(record);
    }
  }
);

8. Chunked HTTP Response

for (const chunk of chunks) {
  if (!res.write(chunk)) {
    await once(res, 'drain');
  }
}
res.end();

9. Progress Reporting

let bytes = 0;
stream.on('data', (chunk) => {
  bytes += chunk.length;
  console.log(`${bytes} bytes received`);
});

10. Pause and Resume

stream.pause();
setTimeout(() => stream.resume(), 1000);

Visual

The Four Stream Types

┌──────────────────────────────────────────────────────────────┐
│  READABLE:                                                   │
│  source ──▶ [readable] ──▶ consumer                          │
│                                                              │
│  WRITABLE:                                                   │
│  producer ──▶ [writable] ──▶ destination                     │
│                                                              │
│  DUPLEX:                                                     │
│  source ──▶ [readable] ──▶ consumer                          │
│  producer ──▶ [writable] ──▶ destination                     │
│  (independent)                                               │
│                                                              │
│  TRANSFORM:                                                  │
│  producer ──▶ [transform] ──▶ consumer                       │
│  (output derived from input)                                 │
└──────────────────────────────────────────────────────────────┘

Pipeline Flow

┌──────────────────────────────────────────────────────────────┐
│  source ──▶ transform 1 ──▶ transform 2 ──▶ destination      │
│                                                              │
│  Each stage processes chunks as they arrive.                 │
│  Backpressure propagates from destination to source.         │
│  Memory usage is bounded by the chunk size, not the data.    │
└──────────────────────────────────────────────────────────────┘

Backpressure

┌──────────────────────────────────────────────────────────────┐
│  FAST PRODUCER:                                              │
│  ──chunk──chunk──chunk──chunk──chunk──▶                      │
│                                                              │
│  SLOW CONSUMER:                                              │
│  ──chunk──────────chunk──────────chunk──▶                    │
│                                                              │
│  Without backpressure:                                       │
│  ┌────────────────────────────────────────────────────────┐  │
│  │  Buffer grows: [chunk][chunk][chunk][chunk]...         │  │
│  └────────────────────────────────────────────────────────┘  │
│                                                              │
│  With backpressure:                                          │
│  write returns false → producer waits → drain → resume      │
└──────────────────────────────────────────────────────────────┘

Stream Lifecycle

┌──────────────────────────────────────────────────────────────┐
│  READABLE:                                                   │
│  data → data → data → end → close                            │
│                                                              │
│  WRITABLE:                                                   │
│  write → write → write → finish → close                      │
│                                                              │
│  TRANSFORM:                                                  │
│  write → transform → data → write → transform → data →      │
│  finish → end → close                                        │
│                                                              │
│  ERROR:                                                      │
│  any stage → error → close                                   │
└──────────────────────────────────────────────────────────────┘

Summary

ItemValue
ReadableProduces data, emits data and end
WritableConsumes data, emits finish and drain
DuplexBoth, independent directions
TransformBoth, output derived from input
pipeConnects readable to writable
pipelineConnects with error propagation and cleanup
Backpressurewrite returns false, wait for drain
High-water mark16 KB for byte streams
Default encodingutf8 for text streams
Object modeStreams of JavaScript objects

Key takeaways:

  • A stream processes data in chunks. It does not load the entire dataset into memory, so it handles files and network data of any size with bounded memory usage.
  • There are four stream types. Readable produces data, Writable consumes it, Duplex does both independently, and Transform derives its output from its input.
  • The data and end events are for Readable streams. The finish and drain events are for Writable streams. The error event is for both.
  • Use pipeline instead of pipe. The pipeline function propagates errors, cleans up on failure, and returns a promise. The pipe method handles backpressure but does not propagate errors.
  • Backpressure is automatic with pipe and pipeline. When writing manually, respect the return value of write and wait for the drain event.
  • The write method returns a boolean. true means the buffer has room; false means the producer should pause until the drain event.
  • Streams are the foundation of Node.js I/O. Files, sockets, HTTP requests and responses, compression, and encryption all expose the stream interface.

Remember: Streams are the abstraction that makes Node.js suitable for I/O-heavy workloads. A stream processes data in chunks, so the memory usage is bounded by the chunk size, not the total data size. The four stream types cover the four directions of data flow: Readable, Writable, Duplex, and Transform. The pipe method connects them, and pipeline does so with error handling. Backpressure is the mechanism that prevents a fast producer from overwhelming a slow consumer. When you process a file, a network response, or a compressed stream, you are working with streams, 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!