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.
| Type | Readable | Writable | Relationship |
|---|---|---|---|
| Readable | Yes | No | — |
| Writable | No | Yes | — |
| Duplex | Yes | Yes | Independent |
| Transform | Yes | Yes | Output 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.
| Type | High-water mark | Signal |
|---|---|---|
| Readable | readableHighWaterMark | readable event |
| Writable | writableHighWaterMark | drain 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.
| Event | Readable | Writable | When |
|---|---|---|---|
data | Yes | No | A chunk is available |
end | Yes | No | No more data |
readable | Yes | No | Data is available to read |
finish | No | Yes | All data has been flushed |
drain | No | Yes | The buffer has room again |
close | Yes | Yes | The underlying resource is closed |
error | Yes | Yes | An error occurred |
pipe | Yes | No | pipe connected the stream |
unpipe | Yes | No | pipe 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
| Type | Readable | Writable | Example |
|---|---|---|---|
| Readable | Yes | No | fs.createReadStream |
| Writable | No | Yes | fs.createWriteStream |
| Duplex | Yes | Yes | net.Socket |
| Transform | Yes | Yes | zlib.createGzip |
Common Events
| Event | Type | Meaning |
|---|---|---|
data | Readable | Chunk available |
end | Readable | No more data |
readable | Readable | Data ready to read |
finish | Writable | All data flushed |
drain | Writable | Buffer has room |
close | Both | Resource closed |
error | Both | Error occurred |
Methods
| Method | Purpose |
|---|---|
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
| Aspect | pipe | pipeline |
|---|---|---|
| Error propagation | No | Yes |
| Cleanup on error | No | Yes |
| Return value | Destination | Promise (promise version) |
| Recommended | Legacy | Yes |
High-Water Marks
| Stream | Default |
|---|---|
| Byte streams | 16 KB |
| Object mode | 16 objects |
Backpressure Signals
| Type | Signal |
|---|---|
| Writable full | write returns false |
| Writable drained | drain event |
| Readable has data | readable 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
| Pitfall | Why It Happens | Fix |
|---|---|---|
| Process crashes on error | No error listener | Add one |
| Memory grows unbounded | Backpressure ignored | Respect write return |
write after end | Wrote after end() | Check the stream state |
| Errors not propagated | Used pipe | Use pipeline |
| Slow consumer | Producer ignores drain | Wait for the drain event |
| Partial file written | Pipeline error not handled | Use pipeline with a catch |
| Stream never ends | Source does not call end | Ensure 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
| Item | Value |
|---|---|
| Readable | Produces data, emits data and end |
| Writable | Consumes data, emits finish and drain |
| Duplex | Both, independent directions |
| Transform | Both, output derived from input |
pipe | Connects readable to writable |
pipeline | Connects with error propagation and cleanup |
| Backpressure | write returns false, wait for drain |
| High-water mark | 16 KB for byte streams |
| Default encoding | utf8 for text streams |
| Object mode | Streams 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
dataandendevents are for Readable streams. Thefinishanddrainevents are for Writable streams. Theerrorevent is for both. - Use
pipelineinstead ofpipe. Thepipelinefunction propagates errors, cleans up on failure, and returns a promise. Thepipemethod handles backpressure but does not propagate errors. - Backpressure is automatic with
pipeandpipeline. When writing manually, respect the return value ofwriteand wait for thedrainevent. - The
writemethod returns a boolean.truemeans the buffer has room;falsemeans the producer should pause until thedrainevent. - 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!