Node.js Streams Explained: A Practical Guide
Learn how Node.js streams process large files and network data with Readable, Writable, Duplex, Transform, pipeline, async iteration, and object mode.
14 min read
Introduction
Imagine you ordered a huge pizza. You do not wait for the whole pizza to arrive before eating. You eat slice by slice as it comes out of the oven.
That is the basic idea behind Node.js streams. Instead of waiting for an entire file, request body, or query result to exist in memory, a stream lets you work with data chunk by chunk as it arrives.
This guide explains the four stream types, how data moves between them, and the APIs you will use in real backend code.
A Simple Mental Model
A stream is an interface for data that arrives or leaves over time.
Most stream code has three parts:
source -> optional processing -> destinationFor example:
uploaded file -> validate and compress -> object storageThe source is usually a Readable stream. The destination is usually a Writable stream. Any processing step in the middle is usually a Transform stream.
A chunk is not necessarily a complete line, JSON object, or database row. It is simply the piece of data Node has available at that moment. Your code must not assume that application-level boundaries match chunk boundaries.
Why Streams Exist
Many Node.js APIs expose streams because backend data is often too large, too slow, or too unpredictable to handle as one value. Common examples include:
- File reads and writes
- HTTP request and response bodies
- TCP sockets
- Compression and encryption
- Database exports
process.stdinandprocess.stdout
Without streams, the simple approach is to load everything first:
import { readFile } from "node:fs/promises";
const data = await readFile("huge.csv");
await processFile(data);That can be perfectly reasonable for a small, bounded file. For a large or user-controlled file, however, the entire payload must fit in memory before processing begins.
A streamed version starts work as data arrives:
import { createReadStream } from "node:fs";
const source = createReadStream("huge.csv");
for await (const chunk of source) {
await processChunk(chunk);
}The file can be much larger than the amount held in memory at one time. Streams still use buffers, so memory is not zero or automatically fixed. The important difference is that the whole payload does not need to exist in memory at once.
The Four Stream Types
Node.js has four fundamental stream types:
| Stream type | Role | Backend example |
|---|---|---|
| Readable | Produces data for your code to consume | File read, HTTP request body |
| Writable | Accepts data and sends or stores it | File write, HTTP response |
| Duplex | Has independent readable and writable sides | TCP socket |
| Transform | Reads input and produces changed output | Gzip compression, CSV-to-NDJSON encoder |
The distinction is about what your code can do with the stream. An incoming HTTP request is readable from the server's point of view. An outgoing HTTP response is writable.
Readable Streams
A Readable stream is a source of data. fs.createReadStream() is a common example:
import { createReadStream } from "node:fs";
const source = createReadStream("large.csv");
for await (const chunk of source) {
console.log("Received", chunk.length, "bytes");
}Consuming a Readable with Async Iteration
for await...of is usually the clearest way to consume a Readable when each chunk requires asynchronous work. Errors from the stream become rejected promises that you can catch with normal try...catch:
import { createReadStream } from "node:fs";
async function countNewlines(path: string): Promise<number> {
let lines = 0;
try {
for await (const chunk of createReadStream(path)) {
for (const byte of chunk as Buffer) {
if (byte === 0x0a) lines++;
}
}
} catch (error) {
console.error("Could not read log file", { path, error });
throw error;
}
return lines;
}This counts line-feed bytes without collecting a large log file into one string.
Building a Custom Readable
You will usually consume streams created by Node.js or a library. When you do need a custom source, extend Readable and push null when no more data remains:
import { Readable } from "node:stream";
class CounterStream extends Readable {
private current = 0;
constructor(private readonly max: number) {
super();
}
_read() {
if (this.current >= this.max) {
this.push(null);
return;
}
this.push(`${this.current++}\n`);
}
}
const counter = new CounterStream(5);
counter.pipe(process.stdout);Writable Streams
A Writable stream is a destination. You send data with .write() and signal completion with .end():
import { createWriteStream } from "node:fs";
const output = createWriteStream("output.txt");
output.write("Hello, ");
output.write("streams!\n");
output.end();The write is asynchronous even though the method does not return a Promise. The stream may temporarily buffer data before the operating system or destination accepts it.
For one-off manual writes, wait for completion when the caller needs to know that the destination finished:
import { once } from "node:events";
import { createWriteStream } from "node:fs";
const output = createWriteStream("report.txt");
output.end("Report complete\n");
await once(output, "finish");When you are connecting a source to a destination, prefer pipeline() instead of managing these events yourself.
Duplex Streams
A Duplex stream is readable and writable. Its two sides are independent: writing data does not automatically make the same data available to read.
A TCP socket is the most useful mental model:
remote client -> readable side of socket
remote client <- writable side of socketimport { createServer } from "node:net";
const server = createServer((socket) => {
socket.setEncoding("utf8");
socket.on("data", (message) => {
socket.write(`received: ${message}`);
});
});
server.listen(4000);The same socket receives data from the client and sends data back, but each direction has its own buffer and lifecycle.
Transform Streams
A Transform stream is a special Duplex stream. You write input to one side, and it produces related output on the other side.
This makes transforms useful for compression, encryption, parsing, and serialization:
import { Transform } from "node:stream";
class UpperCaseTransform extends Transform {
_transform(chunk: Buffer, _encoding: string, callback: () => void) {
this.push(chunk.toString().toUpperCase());
callback();
}
}
process.stdin.pipe(new UpperCaseTransform()).pipe(process.stdout);Chunk Boundaries Are Not Record Boundaries
Suppose a file contains:
first line
second lineOne chunk could end after sec, and the next could begin with ond line. A line-processing transform must keep incomplete data until the next chunk arrives:
import { Transform } from "node:stream";
class LineSplitter extends Transform {
private remainder = "";
constructor() {
super({ readableObjectMode: true });
}
_transform(chunk: Buffer, _encoding: string, callback: () => void) {
const lines = (this.remainder + chunk.toString("utf8")).split("\n");
this.remainder = lines.pop() ?? "";
for (const line of lines) {
this.push(line);
}
callback();
}
_flush(callback: () => void) {
if (this.remainder) this.push(this.remainder);
callback();
}
}_flush() runs when the input ends. It gives the transform one final chance to emit buffered data.
Piping Streams Together
.pipe() connects a Readable to a Writable and returns the destination:
source.pipe(transform).pipe(destination);It also coordinates normal flow between the connected streams. This is convenient for short scripts, but error handling and cleanup become easy to get wrong when a chain has several stages.
For application code, use the Promise-based pipeline() utility:
import { createReadStream, createWriteStream } from "node:fs";
import { pipeline } from "node:stream/promises";
import { createGzip } from "node:zlib";
await pipeline(
createReadStream("input.txt"),
createGzip(),
createWriteStream("input.txt.gz")
);pipeline() is safer than a manual pipe chain because it:
- Forwards errors to one Promise
- Destroys the connected streams when the pipeline fails
- Resolves only after the pipeline finishes
- Coordinates data flow between each stage
Use .pipe() when its event-based lifecycle is genuinely what you want. Use pipeline() when the whole chain represents one operation that should succeed or fail together.
Error Handling
Wrap an awaited pipeline in try...catch and add useful context before rethrowing or returning an error:
import { createReadStream, createWriteStream } from "node:fs";
import { pipeline } from "node:stream/promises";
import { createGzip } from "node:zlib";
async function compressFile(inputPath: string, outputPath: string) {
try {
await pipeline(
createReadStream(inputPath),
createGzip(),
createWriteStream(outputPath)
);
} catch (error) {
console.error("Compression failed", { inputPath, outputPath, error });
throw error;
}
}If you use individual streams without pipeline(), every stream that can emit an error event needs an error-handling strategy. An unhandled stream error can terminate the Node.js process.
Avoid attaching an error listener and then pretending the operation succeeded. Logging, cleanup, and reporting failure are separate responsibilities.
Cancellation with AbortSignal
Long-running stream operations should stop when they are no longer useful. The Promise version of pipeline() accepts an AbortSignal:
import { createReadStream, createWriteStream } from "node:fs";
import { pipeline } from "node:stream/promises";
import { createGzip } from "node:zlib";
const controller = new AbortController();
const compression = pipeline(
createReadStream("large-export.csv"),
createGzip(),
createWriteStream("large-export.csv.gz"),
{ signal: controller.signal }
);
setTimeout(() => controller.abort(), 30_000);
try {
await compression;
} catch (error) {
if (error instanceof Error && error.name === "AbortError") {
console.log("Compression was cancelled");
} else {
throw error;
}
}Aborting rejects the pipeline with an AbortError and destroys the underlying stream chain. In a server, the abort can come from a request deadline, application shutdown, or a client disconnect instead of a timer.
Object Mode
Normal Node.js streams move bytes as Buffer, Uint8Array, or string values. Object mode lets a stream move JavaScript values such as database rows.
import { Transform } from "node:stream";
interface UserRow {
id: number;
email: string;
}
const toNDJSON = new Transform({
writableObjectMode: true,
transform(row: UserRow, _encoding, callback) {
callback(null, JSON.stringify(row) + "\n");
},
});This transform accepts objects but produces bytes, which makes it useful at the boundary between a row-producing database stream and an HTTP response or file.
Use object mode for structured values inside the application. Switch back to byte mode when serializing data for a file, socket, or response.
What highWaterMark Means
Readable and Writable streams keep internal buffers. highWaterMark is the threshold that influences when a Readable should stop fetching more data or when a Writable's .write() method starts returning false.
import { createReadStream } from "node:fs";
const source = createReadStream("large.bin", {
highWaterMark: 256 * 1024,
});For byte streams, the value represents bytes. For object-mode streams, it represents a number of objects.
Three details matter:
- It is a buffering threshold, not a guaranteed memory limit
- A larger value may reduce I/O overhead but increases per-stream buffering
- A smaller value may reduce buffering but increase the number of reads or writes
Defaults vary by stream implementation and Node.js version. Do not copy a value just because it appears in an example. Start with the default, measure the real workload, and tune only when the result solves an observed problem.
The advanced question of how these signals regulate producers across a full production system belongs in the backpressure guide linked below.
Practical Example: Stream a Large Export
Suppose a database library returns rows as an object-mode stream. The API should send NDJSON so the client can process one record per line without waiting for a giant JSON array.
import type { ServerResponse } from "node:http";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";
interface EventRow {
id: number;
type: string;
createdAt: Date;
}
async function sendEventExport(
rows: NodeJS.ReadableStream,
response: ServerResponse,
signal: AbortSignal
) {
const toNDJSON = new Transform({
writableObjectMode: true,
transform(row: EventRow, _encoding, callback) {
const record = {
id: row.id,
type: row.type,
createdAt: row.createdAt.toISOString(),
};
callback(null, JSON.stringify(record) + "\n");
},
});
response.writeHead(200, {
"Content-Type": "application/x-ndjson; charset=utf-8",
});
await pipeline(rows, toNDJSON, response, { signal });
}The operation has a clear shape:
database rows -> serialize each row -> HTTP responseIt starts sending useful data before the query has produced every row, avoids constructing one enormous array, reports errors through the pipeline Promise, and supports cancellation.
One practical warning: after response headers or body bytes have been sent, you may no longer be able to replace the response with a clean JSON error. Log the failure and let the connection close rather than attempting to send a second response.
Streams vs Loading Everything into Memory
Choose buffering when the input is small, bounded, and genuinely easier to use as one value. Choose streaming when data is large, arrives over time, is controlled by a user, or should be forwarded without waiting for the complete payload.
| Loading everything first | Streaming |
|---|---|
| Simpler for small payloads | Better fit for large or unknown payload sizes |
| Full value is immediately available | Code must handle chunks and partial records |
| Memory grows with the complete payload | Memory is mainly shaped by active buffers |
| Work begins after the full read finishes | Work can begin as soon as data arrives |
Streaming is not automatically faster. It mainly changes when work can begin and how much data must be retained at once.
Common Mistakes
Assuming One Chunk Equals One Record
File lines, JSON objects, and protocol messages can be split across chunks. Buffer incomplete records and finish them in the next chunk.
Using .pipe() Without an Error Plan
// Easy to overlook errors and cleanup across the chain
readable.pipe(transform).pipe(writable);
// One Promise represents the complete operation
await pipeline(readable, transform, writable);Collecting Every Chunk
If you push every chunk into an array and call Buffer.concat() at the end, you are buffering the complete payload again. That may be correct for a small, bounded input, but it removes the main memory advantage for a large one.
Forgetting _flush()
If a Transform holds a partial line or parser state, _flush() must emit or validate what remains when the input ends.
Mixing Object Mode and Byte Mode Accidentally
A stream that expects bytes cannot accept arbitrary objects. Make the boundary explicit with writableObjectMode, readableObjectMode, or a serialization transform.
Blocking the Event Loop
Streams limit how much data is retained; they do not make CPU-heavy work asynchronous. Expensive synchronous parsing or compression inside _transform() can still delay every request handled by the process.
Best Practices
- Use
pipeline()when several stages form one operation - Use
for await...ofwhen consuming a Readable with asynchronous code - Treat chunks as arbitrary pieces, not complete application records
- Add cancellation to long-running or request-scoped pipelines
- Use object mode for structured values and byte mode at I/O boundaries
- Start with the default
highWaterMarkand tune from measurements - Keep transforms focused on one job
- Propagate errors instead of only logging them
- Buffer the complete payload only when its maximum size is known and safe
Conclusion
Streams are the natural Node.js interface for data that moves over time.
The foundation is straightforward:
Readable -> Transform -> WritableUse a Duplex stream when data moves independently in both directions. Use for await...of to consume readable data. Use object mode for structured records. Use pipeline() when a chain should complete, fail, or be cancelled as one operation.
Once these mechanics are clear, the next problem is what happens when producers and consumers run at different speeds across streams, APIs, databases, and queues. Continue with the production backpressure guide for flow control, bounded work, admission control, and overload behavior.
References
Related
Written by
Faisal
Software engineer writing about backend systems, Node.js, system design, scalable applications, and modern web and mobile development.