SALES_FILE has been declared since the previous lesson pointing at data/sales.csv, a file that does not exist yet. Its moment has come, and with it the problem none of the tools we have seen so far can solve.
The Escena Viva catalog takes two kilobytes: readFile swallows it whole without breaking a sweat. A season's sales history is another matter. Three venues, two hundred sessions a year, thousands of tickets per session, each with its code, purchase channel and timestamp: hundreds of megabytes. And readFile has no way of reading that without putting it all in memory.
Streams are Node's answer to that problem, and they are far more than an optimization: they are the model Node uses to think about input and output. The HTTP requests and responses of Module 4 are streams; compression, cryptography and network connections are streams. And — this will sound familiar — streams are EventEmitter objects, so everything from lesson 02-05 applies here directly.
Contents
- The problem: measuring what
readFilecosts - What a stream is and its four types
- Streams are
EventEmitterobjects - Flowing mode and paused mode
createReadStream,createWriteStreamandhighWaterMark- Backpressure: what
write()returns pipe()and its weak spot- The Escena Viva
data/sales.csvfile - Reading line by line with
readline for await...ofover a stream
- The problem: measuring what
readFile costs
readFile costsYou should not take this on faith. Let's measure it.
// src/lab/compare-memory.js
// Compares the memory usage of readFile against a stream.
const fs = require('node:fs');
const fsPromises = require('node:fs/promises');
const { SALES_FILE } = require('../config/paths.js');
function memoryMb() {
const { heapUsed, rss } = process.memoryUsage();
return { heapMb: +(heapUsed / 1048576).toFixed(1), rssMb: +(rss / 1048576).toFixed(1) };
}
async function withReadFile() {
const content = await fsPromises.readFile(SALES_FILE, 'utf8');
return { mode: 'readFile', lines: content.split('\n').length, ...memoryMb() };
}
function withStream() {
return new Promise((resolve, reject) => {
let lines = 0;
let leftover = '';
const reader = fs.createReadStream(SALES_FILE, { encoding: 'utf8' });
reader.on('data', (chunk) => {
const parts = (leftover + chunk).split('\n');
leftover = parts.pop(); // the last piece may be cut in half
lines += parts.length;
});
reader.on('end', () => resolve({ mode: 'stream', lines, ...memoryMb() }));
reader.on('error', reject);
});
}With an 800 MB sales.csv, the output is devastating:
┌─────────┬────────────┬──────────┬────────┬────────┐ │ (index) │ mode │ lines │ heapMb │ rssMb │ ├─────────┼────────────┼──────────┼────────┼────────┤ │ 0 │ 'readFile' │ 9200001 │ 1621.4 │ 1712.8 │ │ 1 │ 'stream' │ 9200001 │ 6.2 │ 58.3 │ └─────────┴────────────┴──────────┴────────┴────────┘
And if the file were any bigger, readFile would not even get to print: it would fail with ERR_STRING_TOO_LONG — a V8 string has a limit of roughly 512 MB — or with JavaScript heap out of memory, which kills the whole process. The difference is two orders of magnitude, and the explanation is simple: readFile needs the complete file in memory at once, while the stream processes it in 64 KB chunks and discards each one as soon as it has used it. The stream's usage does not depend on the size of the file: that is the key property, constant memory.
- What a stream is and its four types
A stream is a sequence of data processed in pieces, as they become available, instead of waiting to have them all. The water analogy is a good one: readFile is filling the whole bathtub before touching the water; a stream is opening the tap and working with whatever comes out.
| Type | What it does | Examples in Node |
|---|---|---|
| Readable | Produces data that you consume | fs.createReadStream, process.stdin, an HTTP request body |
| Writable | Consumes data that you produce | fs.createWriteStream, process.stdout, the HTTP response |
| Duplex | Independent reading and writing | A TCP socket (net.Socket) |
| Transform | A Duplex whose output is a function of its input | zlib.createGzip, encryption, your own CSV parser |
Transform streams are the subject of the next lesson. Here we concentrate on the first two, which are the foundation of everything else.
- Streams are
EventEmitter objects
EventEmitter objectsThis is not an analogy: Readable and Writable inherit from EventEmitter. on, once, emit, off and everything from lesson 02-05 works as is, including the golden rule that now takes on its true meaning.
| Event | Stream | When it fires |
|---|---|---|
data |
Readable | A chunk is available. Subscribing switches on flowing mode |
end |
Readable | There is no more data to read |
error |
Both | The operation failed. If you do not listen, it kills the process |
close |
Both | The underlying descriptor has been closed and released |
finish |
Writable | end() was called and everything pending has been flushed |
drain |
Writable | The internal buffer has emptied: you can write again |
Two warnings worth half a lesson. The first: end and finish are not the same. end belongs to the reading side ("the input is over") and finish to the writing side ("I have finished flushing"); mixing them up leads to closing a file before its data is on disk. The second: a stream's error event follows the EventEmitter rule you already know — an error with no listeners throws an uncaught exception and kills the process. A stream without an on('error') is a production outage waiting to happen, and the usual cause is not a broken disk: it is an ENOENT because the file was not there.
- Flowing mode and paused mode
A Readable has two ways of delivering data, and this duality is the source of half the trouble people hit when starting with streams.
| Flowing mode | Paused mode | |
|---|---|---|
| How it is switched on | on('data'), pipe() or resume() |
It is the initial state; you return with pause() |
| Who sets the pace | The stream: it pushes the data | You: you ask for it with read() |
| How you consume | The data callback |
on('readable') + read() |
| Risk | If you are slow to process, data keeps arriving | Forgetting to call read() and stalling |
// Flowing mode: the stream pushes.
reader.on('data', (chunk) => process(chunk));
// Paused mode: you pull.
reader.on('readable', () => {
let chunk;
while ((chunk = reader.read()) !== null) process(chunk);
});Do not mix them. Subscribing to data and calling read() on the same stream produces behavior that is hard to reason about: lost chunks, an end that never arrives, scrambled ordering. Pick one mode per stream and stick to it — although in practice the recommendation is to use neither directly: pipe, pipeline or for await...of solve the problem better. And a third detail about flowing mode: if you subscribe to data after an await, you can lose chunks, because the stream started flowing before your listener arrived. Always subscribe synchronously, right after creating the stream.
createReadStream, createWriteStream and highWaterMark
createReadStream, createWriteStream and highWaterMarkconst fs = require('node:fs');
const reader = fs.createReadStream(SALES_FILE, {
encoding: 'utf8', // without this, chunks are Buffer objects
highWaterMark: 64 * 1024 // internal buffer size: 64 KB (the default)
});
const writer = fs.createWriteStream(outputPath, { flags: 'a' }); // 'a': appendThe highWaterMark is the size of the internal buffer: how many bytes the stream accumulates before deciding it has enough. Lowering it (16 KB) spends less memory per stream in exchange for more system calls and more events; raising it (1 MB) does the opposite. The default 64 KB is reasonable almost always: lowering it makes sense when you handle thousands of simultaneous streams — a server with many connections — and raising it when you process a few very large files. In object mode (which you will see in the next lesson) the highWaterMark counts objects, not bytes, and defaults to 16.
A warning about encoding: if you do not set it, chunks arrive as Buffer objects; and if you set it to 'utf8', the stream makes sure no multibyte character is split between two chunks — a real problem we will look at in depth in lesson 03-06 with StringDecoder.
- Backpressure: what
write() returns
write() returnsHere is the concept that separates people who use streams from people who understand them. Imagine you read from a fast SSD and write to a slow disk or a network connection: data comes in at 500 MB/s and goes out at 50 MB/s. Where does the difference go? Into the write stream's internal buffer, which grows without stopping until memory runs out. In other words: you have swapped a readFile that consumed 800 MB for a stream that consumes... 800 MB.
Backpressure is the mechanism that prevents that, and it rests on a detail almost everybody ignores:
const canContinue = writer.write(chunk);
// true -> the buffer has room, keep writing
// false -> the buffer is FULL. Stop writing and wait for the 'drain' eventwrite() returns a boolean. And that boolean is not informational: it is an order.
// BAD: ignoring the return value.
reader.on('data', (chunk) => {
writer.write(chunk); // returns false and nobody cares
});That code works with small files and blows up with large ones: the writer's buffer grows indefinitely because nobody stops feeding it. It is the most common mistake when writing streams by hand, and the hardest to diagnose, because in development — small file, fast disk — it never shows up.
// GOOD: respecting backpressure.
reader.on('data', (chunk) => {
if (!writer.write(chunk)) {
reader.pause(); // stop reading
writer.once('drain', () => reader.resume()); // resume when it empties
}
});
reader.on('end', () => writer.end());The full cycle is: data → write() returns false → pause() → the buffer empties and the writer emits drain → resume() → back to data. The result is that the slow consumer sets the pace of the fast producer, and memory stays bounded by the highWaterMark. This idea is not exclusive to files: it is the same one that governs TCP and the one that will keep a slow client from taking down your HTTP server in Module 4.
pipe() and its weak spot
pipe() and its weak spotWriting the pause/resume/drain dance by hand in every pipe would be unbearable. pipe() does it for you:
That line is equivalent to the whole previous block: it connects both streams, manages backpressure and calls writer.end() when the reader finishes. They can be chained — a.pipe(b).pipe(c) — because pipe returns the destination stream.
But pipe has a serious and little-known flaw: it neither propagates errors nor cleans up what it leaves behind. If reader fails — an ENOENT, a disk with errors — the error never reaches writer, which stays open with its descriptor unreleased; and if nobody listens for error on reader, the process dies.
With a pipe of three or four stages the problem multiplies: you have to subscribe on('error') on every stream and manually close whichever are left hanging. It is tedious, it is easy to forget one, and that oversight is paid in leaked descriptors until the EMFILE.
That is why the modern recommendation is blunt: use pipe() to understand the concept and stream.pipeline to do the work. pipeline propagates errors, destroys every stream in the chain when something fails and has a promise-based version. It is the central subject of the next lesson.
- The Escena Viva
data/sales.csv file
data/sales.csv fileWe need data to work with. This is the format of the sales history:
ticketCode,sessionId,saleDate,priceCents,channel EV-2026-000001,ses-001-1,2026-06-14T10:22:41,2500,web EV-2026-000002,ses-001-1,2026-06-14T10:23:07,2500,web EV-2026-000003,ses-002-1,2026-06-14T11:02:55,1800,box-office EV-2026-000004,ses-003-2,2026-06-15T09:41:12,4200,app
The columns follow the project conventions: ticketCode with the format EV-<year>-<6 digits>, sessionId shaped as ses-NNN-M, saleDate as an ISO string with no zone, priceCents an integer and channel one of three values (web, box-office, app). We generate a file consistent with the seed: exactly 1811 sales, distributed according to each session's sold count.
// src/lab/generate-sales.js
// Generates data/sales.csv consistent with data/events.json.
const fs = require('node:fs');
const { getCatalog } = require('../catalog-data.js');
const { SALES_FILE } = require('../config/paths.js');
const CHANNELS = ['web', 'box-office', 'app'];
async function main() {
const events = await getCatalog();
const writer = fs.createWriteStream(SALES_FILE, { encoding: 'utf8' });
writer.write('ticketCode,sessionId,saleDate,priceCents,channel\n');
let number = 0;
for (const event of events) {
for (const session of event.sessions) {
for (let i = 0; i < session.sold; i += 1) {
number += 1;
const code = `EV-2026-${String(number).padStart(6, '0')}`;
const channel = CHANNELS[number % CHANNELS.length];
// Sales spread by minutes since the box office opened.
const date = new Date(2026, 5, 14, 10, number % 600).toISOString().slice(0, 19);
// Backpressure is ignored here on purpose: it is 1811 lines.
writer.write(`${code},${session.id},${date},${session.priceCents},${channel}\n`);
}
}
}
writer.end();
// 'finish' arrives when EVERYTHING has been flushed to disk, not when end() returns.
writer.on('finish', () => console.error(`[sales] ${number} sales written`));
}
if (require.main === module) main();node src/lab/generate-sales.js
# [sales] 1811 sales written
wc -l data/sales.csv
# 1812 data/sales.csv (1811 sales + header)The comment about backpressure is deliberate: with 1811 lines there is no problem, but the correct code for real volumes would have to respect the return of write() or, better, be built with pipeline and a generator, as we will do in the next lesson.
- Reading line by line with
readline
readlineA CSV file is processed by lines, but a stream's chunks do not line up with lines: a 64 KB chunk cuts whatever line falls on its edge in half. In the lab of section 1 we solved it by hand with a leftover variable; the node:readline module does it properly and with no code of your own.
// src/reports/sales-by-session.js
// Walks data/sales.csv line by line and accumulates the total per session.
const fs = require('node:fs');
const readline = require('node:readline');
const { SALES_FILE } = require('../config/paths.js');
async function totalBySession(file = SALES_FILE) {
const reader = readline.createInterface({
input: fs.createReadStream(file, { encoding: 'utf8' }),
crlfDelay: Infinity // treats \r\n as a single break (Windows files)
});
const bySession = new Map();
let lineNumber = 0;
for await (const line of reader) {
lineNumber += 1;
if (lineNumber === 1 || line.trim() === '') continue; // header and blanks
const [, sessionId, , priceCents, channel] = line.split(',');
if (!sessionId || Number.isNaN(Number(priceCents))) {
// Diagnostics on stderr: one corrupt line does not abort the report.
console.error(`[sales] line ${lineNumber} skipped: ${line}`);
continue;
}
const entry = bySession.get(sessionId) ??
{ sessionId, tickets: 0, revenueCents: 0, channels: new Set() };
entry.tickets += 1;
entry.revenueCents += Number(priceCents);
entry.channels.add(channel);
bySession.set(sessionId, entry);
}
return [...bySession.values()]
.map((e) => ({ ...e, channels: [...e.channels].join('/') }))
.sort((a, b) => a.sessionId.localeCompare(b.sessionId));
}
module.exports = { totalBySession };node -e "require('./src/reports/sales-by-session.js').totalBySession().then(console.table)"
# ses-001-1 180 tickets 450000 cents web/box-office/app
# ses-001-2 96 tickets 211200 cents web/box-office/app
# ses-002-1 118 tickets 212400 cents web/box-office/app
# ...Three decisions worth mentioning. crlfDelay: Infinity keeps a file generated on Windows from producing lines ending in an invisible \r that ruins the last field. A corrupt line is logged and skipped, it does not abort the report: in a history of millions of records there is always garbage, and a report that dies on line 4,000,000 is worth nothing. And memory usage is constant: bySession holds one entry per session (seven), not per sale.
for await...of over a stream
for await...of over a streamThat for await (const line of reader) deserves an explanation, because it is the most readable way of consuming a stream and you already know it from lesson 02-05, when you saw events.on().
Every Readable is an async iterable. You can walk it directly:
// Consumption with for await: no callbacks, no manual pause/resume.
const reader = fs.createReadStream(SALES_FILE, { encoding: 'utf8' });
let bytes = 0;
for await (const chunk of reader) {
bytes += chunk.length;
}
console.log(`${bytes} characters read`);Its advantages over on('data') are concrete: backpressure is automatic (the loop does not ask for the next chunk until it finishes the body) instead of manual; errors are caught with a normal try/catch around the loop instead of a separate on('error'); you can use await inside; and a break destroys the stream on its own, with no explicit destroy().
That await inside the loop is the decisive difference. If every CSV line requires a database query, with on('data') the queries would all fire at once and sink the server; with for await, reading stops by itself until the query finishes, so backpressure comes for free. The price is raw throughput: for await is somewhat slower because each iteration creates a promise. For a batch process that is irrelevant; on a hot path with heavy traffic, measure it.
Common Mistakes and Tips
- Not subscribing
on('error'). Anerrorwith no listeners on anEventEmitterkills the process, and a stream is anEventEmitter. The definitive solution ispipeline. - Ignoring the return of
write(). It works in development and exhausts memory in production. If you write by hand, respect thedrain. - Confusing
endwithfinish.endis the reader exhausted;finishis the writer flushed. Closing on the wrong event truncates data. - Assuming a chunk is a line. It never is. Use
readlineor aTransformthat splits. - Mixing flowing and paused mode, or subscribing to
dataafter anawait. In the first case behavior is erratic; in the second you lose chunks. - Accumulating every chunk in an array to join them at the end. That is
readFilewith extra steps: you have lost the stream's only advantage. - Tip: if the file fits comfortably in memory and you only read it once,
readFileis simpler and faster; streams are for the huge, the endless, or whatever arrives bit by bit. And when you doubt whether backpressure is happening, printwriter.writableLengthnow and then: if it grows without stopping, you are ignoring it.
Exercises
Exercise 1: sales report by channel and by hour
Extend src/reports/sales-by-session.js with a summarizeByChannel() function that walks data/sales.csv a single time and returns, per channel, the number of tickets, the total revenue and the peak hour (the hour of the day with the most sales). It must use readline, keep constant memory and check that the sum of tickets across all channels is 1811.
Exercise 2: copying with manual backpressure
Write src/lab/copy-with-backpressure.js that copies a large file without using pipe, respecting the drain, and prints every second: megabytes copied, the writer's writableLength and the number of backpressure pauses. Then run a version that ignores the return of write() and compare memory usage with process.memoryUsage().rss.
Exercise 3: filtering the sales of one venue
Write src/lab/filter-sales.js that reads data/sales.csv, keeps only the sales whose sessions belong to a given venue (--venue="Sala Boveda", cross-referencing the catalog) and writes the result to data/sales-<venue>.csv preserving the header. Use readline to read and a createWriteStream to write, respecting backpressure.
Solutions
Solution 1. The key is to accumulate on two levels during the same pass:
for await (const line of reader) {
// ...skip the header and blank lines as before...
const [, , saleDate, priceCents, channel] = line.split(',');
const entry = byChannel.get(channel) ??
{ channel, tickets: 0, revenueCents: 0, byHour: new Array(24).fill(0) };
entry.tickets += 1;
entry.revenueCents += Number(priceCents);
// 'YYYY-MM-DDTHH:mm:ss' -> the hour sits at positions 11 and 12.
entry.byHour[Number(saleDate.slice(11, 13))] += 1;
byChannel.set(channel, entry);
}
return [...byChannel.values()].map(({ byHour, ...rest }) => ({
...rest,
peakHour: byHour.indexOf(Math.max(...byHour))
}));The 24-slot array per channel is the trick that keeps memory constant: sales are not stored, only the counter for each hour. A Map of full dates would grow with the file.
Solution 2. The core of the manual copy:
reader.on('data', (chunk) => {
copied += chunk.length;
if (!writer.write(chunk)) {
pauses += 1;
reader.pause();
writer.once('drain', () => reader.resume());
}
});
reader.on('end', () => writer.end());
reader.on('error', (e) => { console.error(e); writer.destroy(); });
writer.on('error', (e) => { console.error(e); reader.destroy(); });With backpressure respected, writableLength hovers around the highWaterMark and rss stays flat. Without it, writableLength grows without limit and rss climbs toward the size of the file. Notice as well the two crossed on('error') handlers needed to avoid leaking descriptors: that manual work is exactly what pipeline automates.
Solution 3. What matters is cross-referencing the catalog before starting to read the CSV, so it does not happen per line:
const events = await getCatalog();
const venueSessions = new Set(
events.filter((e) => e.venue === venue).flatMap((e) => e.sessions.map((s) => s.id))
);
const writer = fs.createWriteStream(outputPath, { encoding: 'utf8' });
let first = true;
for await (const line of reader) {
if (first) { writer.write(`${line}\n`); first = false; continue; }
if (!venueSessions.has(line.split(',')[1])) continue;
// for await already applies backpressure to READING; here we apply it
// to WRITING by waiting for the drain when the buffer fills up.
if (!writer.write(`${line}\n`)) {
await new Promise((resolve) => writer.once('drain', resolve));
}
}
writer.end();That await new Promise(... 'drain' ...) is backpressure translated into the world of async/await, and it is a pattern worth memorizing: inside a for await, waiting for the drain also stops the reading, because the loop does not ask for the next item until it finishes the body. A Set with the session ids avoids walking the catalog 1811 times.
Conclusion
You now know why streams exist and exactly what problem they solve: you have measured the difference between readFile's 1621 MB and a stream's 6 MB on the same file, and you understand that a stream's usage does not depend on the size of the input. You know the four types — Readable, Writable, Duplex and Transform — and the fundamental fact that they are EventEmitter objects, with data, end, error, close, finish and drain, without confusing end with finish and knowing that an error with no listeners takes the process down.
You tell flowing mode from paused mode and why they must not be mixed; you know what the highWaterMark is and what you gain and lose by moving it. And, above all, you understand backpressure: that write() returns a boolean which is an order, that ignoring it turns your stream into a readFile in disguise, and that the pause → drain → resume cycle is what lets the slow consumer set the pace. You know pipe() automates that dance but neither propagates errors nor cleans up descriptors, and that is why it is not the final tool. Escena Viva, meanwhile, already has its data/sales.csv with the seed's 1811 sales and a src/reports/sales-by-session.js that walks it line by line with readline, accumulating tickets, revenue and channels per session with constant memory. And you have seen the most readable way of consuming any stream: for await...of, which throws in backpressure for free and lets you await inside the loop.
In Transform Streams and pipeline we close the circle. You will learn to write your own streams with Transform (_transform and _flush), to assemble complete processing pipes — raw CSV to objects, filtering by venue, aggregating by session, writing the report — and to replace pipe() with stream.pipeline, which propagates errors, destroys the whole chain when something fails and has a promise-based version. You will also see Readable.from(), async generators as transformations, zlib compression slotted into the pipe and cancellation with AbortSignal.
Node.js Course: From Beginner to Advanced
Module 1: Introduction to Node.js
- What Is Node.js?
- Installing and Setting Up the Environment
- Your First Node.js Program
- The Node.js REPL
- Modern JavaScript for Node.js
- The Course Project: the Escena Viva Platform
Module 2: Core Concepts
- Node.js Architecture
- The Event Loop
- Callbacks and Asynchronous Programming
- Promises and async/await
- Events and EventEmitter
- CommonJS Modules and require()
- ES Modules and Interoperability
Module 3: File System and I/O
- Reading and Writing Files
- The fs Module in Depth
- Cross-Platform Paths with the path Module
- Working with Streams
- Transform Streams and pipeline
- Buffers and Binary Data
Module 4: HTTP and Web Servers
- Creating a Simple HTTP Server
- Handling Requests and Responses
- Manual Routing
- Serving Static Files
- Receiving Data: Request Bodies and JSON
- Consuming External APIs from Node.js
Module 5: NPM and Package Management
- Introduction to NPM and package.json
- Installing and Using Packages
- Semantic Versioning and package-lock
- npm Scripts and Project Automation
- Creating and Publishing Packages
- Dependency Security and Maintenance
Module 6: The Express.js Framework
- Introduction to Express.js
- Setting Up an Express Application
- Routing in Express
- Middleware
- Essential Third-Party Middleware
- Input Data Validation
- Error Handling
Module 7: Databases and ORMs
- Introduction to Databases
- Using MongoDB with Mongoose
- CRUD Operations
- Relationships, Population and Advanced Queries
- Using SQL Databases with Sequelize
- Migrations, Transactions and Seed Data
Module 8: Authentication and Authorization
- Introduction to Authentication
- User Registration and Password Hashing
- Sessions and Cookies with Passport.js
- Authentication with JWT
- Role-Based Access Control
- API Security Best Practices
Module 9: Testing and Debugging
- Introduction to Testing
- Unit Testing with Mocha and Chai
- Test Doubles with Sinon
- Integration Testing
- Coverage and Test Automation
- Debugging Node.js Applications
Module 10: Advanced Topics
- The Cluster Module
- Worker Threads
- Caching and Job Queues with Redis
- Performance Optimization
- Building RESTful APIs
- GraphQL with Node.js
Module 11: Deployment and DevOps
- Configuration and Environment Variables
- Logging and Monitoring in Production
- Using PM2 for Process Management
- Packaging with Docker
- Deploying to Heroku and Other PaaS
- Continuous Integration and Deployment
