Overview
I once wrote a Node script that read a 4GB log file into memory, filtered it, and wrote the result back out. On my laptop it worked. On the production box with 2GB of RAM, it crashed the process and took down the service running on the same machine. That was the day I actually learned streams.
The API isn't friendly. The docs read like they were written for people who already understand streams. This is the version I wish I'd read first.
The mental model
A stream is a sequence of data that arrives over time, rather than all at once. Instead of holding the whole file in memory, you process it in chunks as they arrive.
Four stream types:
| Type | Reads or writes | Example |
|---|---|---|
| Readable | Reads | fs.createReadStream, HTTP request |
| Writable | Writes | fs.createWriteStream, HTTP response |
| Duplex | Both | TCP socket |
| Transform | Reads, transforms, writes | zlib.createGzip |
Every stream works the same way: data flows through, and you attach handlers for what to do when chunks arrive, when the stream ends, and when it errors.
Reading a file
const fs = require("fs");
const stream = fs.createReadStream("big.log", {
encoding: "utf8",
highWaterMark: 64 * 1024, // 64KB chunks
});
stream.on("data", (chunk) => {
console.log(`Got ${chunk.length} bytes`);
});
stream.on("end", () => {
console.log("Done");
});
stream.on("error", (err) => {
console.error("Failed:", err);
});
highWaterMark is the chunk size. The default is 64KB for files, 16KB for most other streams. Increasing it means fewer callbacks and less overhead; decreasing it means lower memory usage. For most cases, the default is fine.
The stream reads chunks as fast as you consume them, applying backpressure when you don't. Which brings us to the topic people skip.
Backpressure
When you write to a stream faster than the destination can handle, data builds up in memory. Node handles this automatically when you use pipe(), but not when you write manually.
// Broken — ignores backpressure, buffers everything
readStream.on("data", (chunk) => {
writeStream.write(chunk); // returns false when the buffer is full
});
// Correct — respects backpressure
readStream.on("data", (chunk) => {
const ok = writeStream.write(chunk);
if (!ok) {
readStream.pause();
writeStream.once("drain", () => readStream.resume());
}
});
write() returns false when the internal buffer exceeds highWaterMark. When that happens, you should stop reading until the writable emits drain, meaning the buffer has flushed and there's room again.
This is the entire reason pipe() exists. It handles backpressure so you don't have to.
pipe() and pipeline()
const fs = require("fs");
const zlib = require("zlib");
// Compress a file
fs.createReadStream("input.log")
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream("input.log.gz"));
This is clean and it handles backpressure automatically. The problem is error handling: if any stage fails, pipe() doesn't propagate the error and doesn't clean up the other stages.
// The classic pipe() problem
fs.createReadStream("input.log")
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream("output.gz"))
.on("error", (err) => {
// Only catches errors on the last stage
});
If the read fails, the write stream keeps running and the error goes unhandled. Use pipeline() instead:
const { pipeline } = require("stream/promises");
await pipeline(
fs.createReadStream("input.log"),
zlib.createGzip(),
fs.createWriteStream("output.gz")
);
Or the callback version if you're not in an async context:
const { pipeline } = require("stream");
pipeline(
fs.createReadStream("input.log"),
zlib.createGzip(),
fs.createWriteStream("output.gz"),
(err) => {
if (err) console.error("Pipeline failed:", err);
else console.log("Pipeline succeeded");
}
);
pipeline() propagates errors from any stage, cleans up the other streams when one fails, and works with await. There is no reason to use pipe() in new code.
Processing line by line
Reading a log file as raw chunks is awkward because lines get split across chunk boundaries. Node has a built-in helper:
const fs = require("fs");
const readline = require("readline");
const stream = fs.createReadStream("access.log");
const rl = readline.createInterface({ input: stream, crlfDelay: Infinity });
for await (const line of rl) {
if (line.includes("500")) {
console.log(line);
}
}
crlfDelay: Infinity makes readline treat \r\n as a single line break, which matters for files created on Windows. Without it, every line has a trailing \r.
The for await loop processes one line at a time and applies backpressure automatically. This is the pattern I use for almost every "process a big file" task now.
Transform streams
A Transform stream reads data, modifies it, and writes it downstream. It's how you build a processing stage that fits between a read and a write.
const { Transform } = require("stream");
class UppercaseTransform extends Transform {
_transform(chunk, encoding, callback) {
this.push(chunk.toString().toUpperCase());
callback();
}
}
await pipeline(
fs.createReadStream("input.txt"),
new UppercaseTransform(),
fs.createWriteStream("output.txt")
);
Three things to remember about _transform:
- Call
this.push(data)to emit output. You can call it zero, one, or many times per input chunk. - Call
callback()when done. If you call it with an error, the stream fails. - Don't call
callback()before you're done pushing. The stream will proceed before your output is ready.
A more realistic example — parse each line as JSON and filter:
class JsonFilterTransform extends Transform {
_transform(chunk, encoding, callback) {
const lines = chunk.toString().split("\n");
for (const line of lines) {
if (!line.trim()) continue;
try {
const obj = JSON.parse(line);
if (obj.level === "error") {
this.push(JSON.stringify(obj) + "\n");
}
} catch (err) {
// skip malformed lines
}
}
callback();
}
}
This has a subtle bug: chunks don't align with line boundaries. A line split across two chunks will fail to parse. For correctness, use readline or track a leftover buffer between calls. In practice, for line-oriented text streams, most chunks happen to align, but the failure is intermittent and confusing when it happens.
Reading HTTP request bodies
An incoming HTTP request is a Readable stream. If you don't consume it, the client waits.
// Express with body-parser — fine for small JSON bodies
app.use(express.json());
// For a large upload, stream it to disk
app.post("/upload", (req, res) => {
const filename = req.headers["x-filename"];
const writeStream = fs.createWriteStream(`/uploads/${filename}`);
req.pipe(writeStream);
writeStream.on("finish", () => {
res.json({ status: "ok" });
});
writeStream.on("error", (err) => {
res.status(500).json({ error: err.message });
});
});
Or with pipeline for proper error handling:
app.post("/upload", async (req, res) => {
try {
const writeStream = fs.createWriteStream(`/uploads/${req.headers["x-filename"]}`);
await pipeline(req, writeStream);
res.json({ status: "ok" });
} catch (err) {
res.status(500).json({ error: err.message });
}
});
This handles arbitrarily large uploads without holding them in memory. The express.json() middleware would buffer the whole body first, which is fine for a 10KB JSON payload and catastrophic for a 4GB video.
When streams are the wrong choice
- Small files. Reading a 5KB config file with
fs.readFileSyncis faster and simpler. Streams have overhead per chunk. - Data that needs random access. Streams are sequential. If you need to jump to a specific offset, use file descriptors directly.
- Data you need to process multiple times. Streams are consumed. If you need to read the input twice, either buffer it or re-open the source.
- JSON documents. You can't stream-parse arbitrary JSON with built-in Node. You either buffer, or use a library like
stream-jsonthat parses incrementally.
The last one catches people. Streaming a 100MB JSON array doesn't work with JSON.parse — you need a proper streaming parser, and those have their own learning curve.
The three things to remember
- Use
pipeline(), notpipe(). Always. - Use
readlinefor text files, not manual chunk splitting. - Streams apply backpressure. If you're writing manually instead of piping, you need to handle it or you'll buffer everything in memory without realizing it.
