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

  1. objectMode: streams that carry objects
  2. The Transform stream: _transform and _flush
  3. A Transform that parses CSV lines
  4. Custom Readable and Writable streams
  5. stream.pipeline versus pipe()
  6. The Escena Viva processing pipe
  7. Compression with zlib in the pipe
  8. Readable.from() and async generators
  9. Errors and cancellation with AbortSignal

  1. objectMode: streams that carry objects

By 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.

  1. The Transform stream: _transform and _flush

A 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 throws ERR_MULTIPLE_CALLBACK.
  • Errors are signaled with callback(error), not with throw: an exception thrown in an asynchronous context is caught by nobody.
  • _flush is 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 a Transform apart from a plain map: it can accumulate across chunks.

  1. A Transform that parses CSV lines

Here 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.

  1. Custom Readable and Writable streams

A 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.

  1. 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.

  1. 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.

  1. Compression with zlib in the pipe

Archiving 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)
  );
}
ls -l data/sales*
# sales.csv 112340   |   sales-2026.csv.gz 14208   <- 87 % less

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.

  1. Readable.from() and async generators

Readable.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.

  1. Errors and cancellation with AbortSignal

A 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 _transform or _write. The pipe hangs silently, with no error and no trace. It is the number one failure.
  • Using throw inside _transform instead of callback(error), or forgetting _flush, so everything accumulated — the last line, the total, the closing — is lost without warning.
  • Not declaring objectMode when emitting objects: Node complains with ERR_INVALID_ARG_TYPE because it expects a Buffer or a string. And push(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 Transform that parses, filters and aggregates all at once is impossible to test; three chained stages are tested separately with Readable.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

Module 2: Core Concepts

Module 3: File System and I/O

Module 4: HTTP and Web Servers

Module 5: NPM and Package Management

Module 6: The Express.js Framework

Module 7: Databases and ORMs

Module 8: Authentication and Authorization

Module 9: Testing and Debugging

Module 10: Advanced Topics

Module 11: Deployment and DevOps

Module 12: Real-World Projects

© Copyright 2026. All rights reserved