In the previous lesson you consumed streams that already existed: createReadStream produced chunks and createWriteStream received them. You understood backpressure, you saw that pipe() handles it for you and you discovered its flaw: it neither propagates errors nor cleans up what it leaves behind.
Now comes the other half. You are going to write your own streams — the middle link that turns one thing into another — and to assemble the complete processing pipe of Escena Viva: read data/sales.csv, turn every line into a sale object, filter by venue, aggregate by session and write the result, all with constant memory, whether the file is two kilobytes or eight hundred megabytes. And you will replace pipe() with stream.pipeline, the real tool: it propagates errors, destroys the whole chain when something fails and has a promise-based version. At the end you will see the most readable way of writing transformations, which is not a class but an async generator.
Contents
objectMode: streams that carry objects- The
Transformstream:_transformand_flush - A
Transformthat parses CSV lines - Custom
ReadableandWritablestreams stream.pipelineversuspipe()- The Escena Viva processing pipe
- Compression with
zlibin the pipe Readable.from()and async generators- Errors and cancellation with
AbortSignal
objectMode: streams that carry objects
objectMode: streams that carry objectsBy default, a stream carries bytes: Buffer objects or strings. That is the right thing for files and sockets, but useless when what flows through the pipe are sales, events or already interpreted records. objectMode: true changes that nature:
| Binary mode (default) | objectMode: true |
|
|---|---|---|
| What travels | Buffer or string |
Any value except null |
What the highWaterMark counts |
Bytes (64 KB) | Number of objects (16) |
| Chunks | Split and merged | Every push is one indivisible item |
The mode can be declared per side: new Transform({ readableObjectMode: true, writableObjectMode: false }) describes exactly what a parser does — it receives text and emits objects — while a plain objectMode: true switches on both sides. And null deserves a warning: in a stream, push(null) means "that's the end", so you cannot send null as data, not even in object mode; to represent the absence of something, use undefined or a marker object.
- The
Transform stream: _transform and _flush
Transform stream: _transform and _flushA Transform is a stream that reads on one side, does something and writes on the other. It is implemented with two methods:
| Method | When it is called | What for |
|---|---|---|
_transform(chunk, encoding, callback) |
For every chunk | Process and emit with this.push(...) |
_flush(callback) |
When the input closes | Emit whatever was left pending |
The mechanics are always the same: _transform receives a chunk, does its work, emits zero or more results with this.push(...) and announces it has finished by calling callback(); _flush runs exactly once at the end, when nothing more will come in, and is the last chance to emit.
Three rules that admit no exception:
callback()must be called exactly once per_transform. If you never call it, the pipe hangs forever with no message at all; if you call it twice, Node throwsERR_MULTIPLE_CALLBACK.- Errors are signaled with
callback(error), not withthrow: an exception thrown in an asynchronous context is caught by nobody. _flushis where the accumulated state lives. Everything you cannot emit until you see the end — the last line with no break, an aggregate total, the closing of a JSON document — comes out there. That is what tells aTransformapart from a plainmap: it can accumulate across chunks.
- A
Transform that parses CSV lines
Transform that parses CSV linesHere is the real problem: the file's chunks do not line up with the lines. A stateful Transform solves splitting and conversion in a single step.
// src/streams/parse-sales.js
// Transform: receives text from data/sales.csv and emits sale objects.
const { Transform } = require('node:stream');
class SalesParser extends Transform {
#leftover = ''; // fragment of a line left half-finished
#first = true; // so we can skip the header
#discarded = 0;
constructor(options = {}) {
// Text goes in, objects come out: the two sides differ.
super({ ...options, writableObjectMode: false, readableObjectMode: true });
}
#convert(line) {
const [ticketCode, sessionId, saleDate, priceCents, channel] = line.split(',');
if (!ticketCode || !sessionId || Number.isNaN(Number(priceCents))) {
this.#discarded += 1; // corrupt line: skipped, does not abort
return;
}
this.push({ ticketCode, sessionId, saleDate, priceCents: Number(priceCents), channel });
}
_transform(chunk, encoding, callback) {
const lines = (this.#leftover + chunk).split('\n');
this.#leftover = lines.pop(); // the last one may be cut in half
for (const line of lines) {
if (this.#first) { this.#first = false; continue; } // header
if (line.trim() !== '') this.#convert(line);
}
callback();
}
_flush(callback) {
// The last line, if the file does not end in a break.
if (this.#leftover.trim() !== '') this.#convert(this.#leftover);
if (this.#discarded > 0) console.error(`[sales] ${this.#discarded} discarded`);
callback();
}
}
module.exports = { SalesParser };Notice the role of #leftover: it is the problem readline solved for us in the previous lesson, now solved by hand because we need to emit objects, not lines. And _flush is where that pending leftover gets processed: without that line, the last sale would be silently lost whenever the file does not end in a break. It is a classic bug and hard to spot, because it fails on one record out of a million.
- Custom
Readable and Writable streams
Readable and Writable streamsA Readable produces data. It is implemented with _read(), which Node calls when it wants more material, and inside you emit with this.push(value); when there is nothing left to produce, this.push(null) signals the end of the stream. In practice you almost never need to write one: Readable.from() (section 8) covers nearly every case in a single line.
A Writable consumes data. It implements _write(chunk, encoding, callback), and that callback is the key piece: until you call it, the stream considers you still busy. That is where the backpressure you felt in the previous lesson is born.
const { Writable } = require('node:stream');
class SessionAccumulator extends Writable {
#bySession = new Map();
constructor() { super({ objectMode: true }); }
_write(sale, encoding, callback) {
const entry = this.#bySession.get(sale.sessionId) ??
{ sessionId: sale.sessionId, tickets: 0, revenueCents: 0 };
entry.tickets += 1;
entry.revenueCents += sale.priceCents;
this.#bySession.set(sale.sessionId, entry);
callback(); // ready for the next one
}
get result() {
return [...this.#bySession.values()].sort((a, b) => a.sessionId.localeCompare(b.sessionId));
}
}If the work in _write were asynchronous — saving to a database, calling an API — it would be enough to call the callback when it finished: the whole pipe slows down by itself in the meantime. There is also _writev for processing several items at once, useful when inserting in batches.
stream.pipeline versus pipe()
stream.pipeline versus pipe()You already saw the problem with pipe(). stream.pipeline fixes it, and its promise version — const { pipeline } = require('node:stream/promises') — is the one we will always use.
pipe() |
pipeline() |
|
|---|---|---|
| Backpressure | Yes | Yes |
| Propagates errors along the chain | No | Yes: one stage failing aborts everything |
| Destroys the streams on failure | No: descriptors stay open | Yes: destroy() on all of them |
| Tells you it has finished | By listening for finish by hand |
The promise resolves |
| Cancellation / errors | No / one on('error') per stream |
AbortSignal / a single try/catch |
The practical difference is enormous:
// With pipe: one on('error') per stage and manual cleanup.
source.on('error', handle);
transformation.on('error', handle);
destination.on('error', handle);
source.pipe(transformation).pipe(destination);
// With pipeline: one line and a try/catch.
try {
await pipeline(source, transformation, destination);
} catch (error) {
console.error(`[pipeline] failed: ${error.message}`);
}Course rule: never pipe() in production code. pipe is there to explain the concept; pipeline is what you write.
- The Escena Viva processing pipe
With the pieces above we assemble the complete process: from raw CSV to a JSON report, filtering by venue and never loading the whole file.
// src/reports/sales-pipeline.js
// Sales CSV -> objects -> filter by venue -> aggregation -> JSON.
const fs = require('node:fs');
const path = require('node:path');
const { Transform } = require('node:stream');
const { pipeline } = require('node:stream/promises');
const { SalesParser } = require('../streams/parse-sales.js');
const { getCatalog } = require('../catalog-data.js');
const { SALES_FILE, REPORTS_DIR } = require('../config/paths.js');
const { writeAtomic } = require('../utils/atomic-write.js');
// Object-to-object Transform: lets only the allowed sessions through.
// In shorthand form: options and a transform method, with no subclass.
const filterBySessions = (allowed) => new Transform({
objectMode: true,
transform(sale, encoding, callback) {
if (allowed.has(sale.sessionId)) this.push(sale);
callback();
}
});
async function generateSalesBySession({ venue = null } = {}) {
// The catalog cross-reference happens ONCE, not per CSV line.
const events = await getCatalog();
const relevant = venue ? events.filter((e) => e.venue === venue) : events;
const allowedSessions = new Set(relevant.flatMap((e) => e.sessions.map((s) => s.id)));
const accumulator = new SessionAccumulator();
await pipeline(
fs.createReadStream(SALES_FILE, { encoding: 'utf8' }),
new SalesParser(),
filterBySessions(allowedSessions),
accumulator
);
const sessions = accumulator.result;
const totals = sessions.reduce((t, s) => ({
tickets: t.tickets + s.tickets,
revenueCents: t.revenueCents + s.revenueCents
}), { tickets: 0, revenueCents: 0 });
const report = { generatedAt: new Date().toISOString(), venue: venue ?? 'all', sessions, totals };
const file = path.join(REPORTS_DIR, 'sales-by-session.json');
await writeAtomic(file, JSON.stringify(report, null, 2));
return { file, ...totals };
}
module.exports = { generateSalesBySession };node -e "require('./src/reports/sales-pipeline.js').generateSalesBySession().then(console.log)"
# { file: '.../reports/sales-by-session.json', tickets: 1811, revenueCents: 5389800 }The 1811 tickets match the seed, which is the check we were after. And look at the shape of the pipe: four stages, each with a single responsibility, chained in a pipeline that reads top to bottom; changing the filter, adding a validation or writing to another destination means touching one line. A design detail: the SessionAccumulator goes last because an aggregation cannot emit anything until it has seen the final item; if you needed to continue the pipe after aggregating, it would be a Transform emitting its result in _flush.
- Compression with
zlib in the pipe
zlib in the pipeArchiving the history means slotting in one more stage: zlib.createGzip() is a Transform just like yours.
// src/reports/archive-sales.js
const fs = require('node:fs');
const zlib = require('node:zlib');
const { pipeline } = require('node:stream/promises');
async function archiveSales(sourcePath, targetPath) {
await pipeline(
fs.createReadStream(sourcePath),
zlib.createGzip({ level: 9 }), // 1 = fast, 9 = maximum compression
fs.createWriteStream(targetPath)
);
}No encoding was given to the reader here: the data travels as Buffer objects from start to finish, which is right for compression — turning it into text and back into bytes would add nothing and cost CPU. Decompression is the same pipe in reverse: just slot zlib.createGunzip() between the .gz reader and the parser to process the compressed history without decompressing it to disk. That is what makes the model powerful: stages combine without any of them knowing anything about the others, and the parser has no idea its bytes came from a compressed file.
Readable.from() and async generators
Readable.from() and async generatorsReadable.from() turns any iterable — array, Map, generator — into a stream:
const { Readable } = require('node:stream');
// Any old array, already a stream in object mode.
const sales = Readable.from([{ sessionId: 'ses-001-1' }, { sessionId: 'ses-002-1' }]);It is invaluable for the tests in Module 9: you feed a pipe with fake data without touching the disk. But the best part comes now: an async generator can be used directly as a stage of a pipeline, and it is the most readable way of writing a transformation.
// The same logic as SalesParser, with no class and no callbacks.
// convert(line) returns the sale object, like the private method before.
async function* parseSales(source) {
let leftover = '';
let first = true;
for await (const chunk of source) {
const lines = (leftover + chunk).split('\n');
leftover = lines.pop();
for (const line of lines) {
if (first) { first = false; continue; }
if (line.trim() !== '') yield convert(line);
}
}
if (leftover.trim() !== '') yield convert(leftover); // the equivalent of _flush
}
// Filter: three lines, no class, no explicit objectMode.
const filterBySessions = (allowed) => async function* (source) {
for await (const sale of source) {
if (allowed.has(sale.sessionId)) yield sale;
}
};
await pipeline(
fs.createReadStream(SALES_FILE, { encoding: 'utf8' }),
parseSales,
filterBySessions(allowedSessions),
accumulator
);The comparison is revealing:
Transform class |
Async generator | |
|---|---|---|
| Lines of code | Constructor, _transform, _flush |
Far fewer |
objectMode |
Must be declared | Implicit |
| End of input | _flush |
The code after the for await |
| State across chunks | Class field | Local variable |
await inside / reusable as an object |
Awkward / yes | Natural / no |
For new transformations, always start with an async generator. Fall back to the Transform class when you need an object with state readable from outside — like our SessionAccumulator, whose result is read when it finishes — or when you have to integrate with an API expecting a specific stream.
- Errors and cancellation with
AbortSignal
AbortSignalA failure in any stage aborts the pipe and rejects the promise; what reaches the catch is the original error, with its code intact:
try {
await pipeline(source, parseSales, accumulator);
} catch (error) {
// Expected cases: reported and degraded. The rest is propagated.
if (error.code === 'ENOENT') return warn('the sales file does not exist');
if (error.name === 'AbortError') return warn('process cancelled');
throw error;
}And pipeline accepts a cancellation signal, just like fetch:
const controller = new AbortController();
// A report may not take longer than thirty seconds.
const limit = setTimeout(() => controller.abort(), 30_000);
try {
await pipeline(source, parseSales, accumulator, { signal: controller.signal });
} finally {
clearTimeout(limit); // do not leave timers alive
}On aborting, pipeline destroys every stream in the chain: descriptors are released and the promise is rejected with an AbortError. It is exactly the work that in the previous lesson had to be written by hand with crossed on('error') handlers and destroy(). In Module 4 you will come back to this: when an HTTP client closes the connection halfway through a download, the right move is to abort the pipe generating the response instead of working on for nobody.
Common Mistakes and Tips
- Forgetting
callback()in_transformor_write. The pipe hangs silently, with no error and no trace. It is the number one failure. - Using
throwinside_transforminstead ofcallback(error), or forgetting_flush, so everything accumulated — the last line, the total, the closing — is lost without warning. - Not declaring
objectModewhen emitting objects: Node complains withERR_INVALID_ARG_TYPEbecause it expects aBufferor a string. Andpush(null)is not data: it means end of stream, always. - Still using
pipe(). If an error from the first stage never reaches the last one, you have leaked descriptors and invisible failures. - Tip: one stage, one responsibility. A
Transformthat parses, filters and aggregates all at once is impossible to test; three chained stages are tested separately withReadable.from(). And if you hesitate between a class and a generator, start with the generator: turning it into a class later is easy, the other way round less so.
Exercises
Exercise 1: a validator in the pipe
Write an async generator validateSales that slots in between the parser and the filter, and checks every sale: ticketCode with the format EV-<year>-<6 digits>, sessionId shaped as ses-NNN-M, priceCents a positive integer and channel one of web, box-office and app. Valid sales pass through; invalid ones are counted and logged on stderr with their line number. At the end it must report how many were discarded.
Exercise 2: compressed report by channel
Write src/reports/sales-by-channel.js with a pipe that reads data/sales.csv, aggregates by channel (tickets, revenue and average ticket) and writes reports/sales-by-channel.json.gz compressed, without producing the uncompressed JSON on disk. Hint: Readable.from([JSON.stringify(report, null, 2)]) turns the result into the source of a second pipe.
Exercise 3: a pipe with a time limit and retries
Write a function generateWithRetry(options) that runs the generateSalesBySession pipe with a 10-second AbortSignal and, if it fails with an AbortError or a transient I/O error, retries it up to three times using retry and sleep from Module 2. It must not retry on ENOENT or on corrupt data: those failures are not fixed by repeating.
Solutions
Solution 1. The validator as a generator turns out almost declarative:
const CODE_PATTERN = /^EV-\d{4}-\d{6}$/;
const SESSION_PATTERN = /^ses-\d{3}-\d+$/;
const CHANNELS = new Set(['web', 'box-office', 'app']);
async function* validateSales(source) {
let number = 0;
let discarded = 0;
for await (const sale of source) {
number += 1;
const problems = [];
if (!CODE_PATTERN.test(sale.ticketCode)) problems.push('code');
if (!SESSION_PATTERN.test(sale.sessionId)) problems.push('session');
if (!Number.isInteger(sale.priceCents) || sale.priceCents <= 0) problems.push('price');
if (!CHANNELS.has(sale.channel)) problems.push('channel');
if (problems.length === 0) { yield sale; continue; }
discarded += 1;
console.error(`[validate] sale ${number} discarded (${problems.join(', ')})`);
}
if (discarded > 0) console.error(`[validate] ${discarded} sales discarded`);
}The final block, after the for await, is the exact equivalent of a Transform class's _flush: it runs when the input is exhausted. If you needed to parameterize the validator, wrap it in a function returning the generator, like the filter in section 8.
Solution 2. Two chained pipes, the second one with Readable.from:
const accumulator = new ChannelAccumulator();
await pipeline(fs.createReadStream(SALES_FILE, { encoding: 'utf8' }), parseSales, accumulator);
const report = accumulator.result.map((c) => ({
...c,
averageTicketCents: Math.round(c.revenueCents / c.tickets)
}));
await pipeline(
Readable.from([JSON.stringify(report, null, 2)]),
zlib.createGzip(),
fs.createWriteStream(path.join(REPORTS_DIR, 'sales-by-channel.json.gz'))
);The JSON never touches the disk uncompressed: it leaves memory as a stream, goes through the compressor and lands in the file. With a small report it makes no difference; with one of hundreds of megabytes it is the difference between needing twice the free space or not.
Solution 3. The key is to classify the errors before deciding whether to retry:
const NON_RETRYABLE = new Set(['ENOENT', 'EACCES', 'CORRUPT_DATA']);
const generateWithRetry = (options = {}) => retry(async () => {
const controller = new AbortController();
const limit = setTimeout(() => controller.abort(), 10_000);
try {
return await generateSalesBySession({ ...options, signal: controller.signal });
} catch (error) {
// We mark the definitive failures so retry does not insist.
if (NON_RETRYABLE.has(error.code ?? error.appCode)) error.permanent = true;
throw error;
} finally {
clearTimeout(limit);
}
}, { attempts: 3, delayMs: 500, shouldRetry: (e) => !e.permanent });Retrying an ENOENT is pointless: the file is not going to show up on its own and all you achieve is delaying the error message three times over. An AbortError from a time limit or a transient EBUSY, on the other hand, may well be resolved on the second attempt. Retrying without classifying the error is one of the most expensive ways of hiding a problem.
Conclusion
You no longer only consume streams: you write them. You know that objectMode turns a pipe of bytes into a pipe of objects, that the highWaterMark starts counting items and that readableObjectMode and writableObjectMode are declared separately because a parser receives text and emits objects. You have mastered Transform with _transform — calling callback() exactly once and signaling errors with callback(error) — and _flush, where everything accumulated is emitted, including that last line with no break that gets silently lost when it is missing. You have built a stateful CSV parser, an aggregating Writable and a filter, and you have chained them into the complete Escena Viva pipe: CSV → objects → filter by venue → aggregation → JSON, with the 1811 sales matching the seed and constant memory all the way through. And you have replaced pipe() with stream.pipeline from node:stream/promises, which propagates errors along the whole chain, destroys every stream on failure, integrates with try/catch and accepts an AbortSignal. The rule is firm: pipe to explain, pipeline to work.
You have also seen how zlib.createGzip slots in as one more stage — and how a pipe can read straight from a .gz without decompressing to disk —, how Readable.from() turns any iterable into a stream, and why an async generator is almost always better than a Transform class: less code, implicit objectMode, natural await and the end of input handled without ceremony. And you know how to cancel a pipe with AbortSignal without leaving descriptors alive.
One piece remains at the bottom of all this. Every time we wrote { encoding: 'utf8' } we were asking for a translation from bytes to text, and every time we left it out — when compressing, when copying — the data traveled as a Buffer. In Buffers and Binary Data we will finally look inside that box: what a Buffer is and why it lives outside the V8 heap, why allocUnsafe can show you somebody else's memory, the encodings and their uses, byte ordering, the classic mistake of slice sharing memory, why an emoji breaks a naive split and how StringDecoder avoids it. And we will apply it to Escena Viva: detecting by its magic numbers whether the poster an organizer uploads really is a PNG, and encoding a ticket's QR code in base64url.
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
