AdvancedTypeScript · Lesson 5 of 10

Streams & Processing Large Data

Process files and network data piece by piece with streams and pipeline.

readFile loads a whole file into memory — fine for small files, a crash waiting to happen for a 5 GB log. Streams process data in chunks, so memory use stays flat no matter how large the input.

There are four kinds: Readable (source), Writable (destination), Duplex (both), and Transform (modifies data as it passes through). pipeline from node:stream/promises connects them and correctly handles errors and cleanup.

Streams handle backpressure: if the destination is slower than the source, reading pauses automatically. Readable streams are also async iterables, so for await...of works on them.

streams.tsTypeScript
import { createReadStream, createWriteStream } from "node:fs";
import { writeFile } from "node:fs/promises";
import { createInterface } from "node:readline";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";
import { createGzip } from "node:zlib";

async function main(): Promise<void> {
  await writeFile("results.csv", "name,score\nAmina,88\nJuma,42\nNeema,71\n");

  // 1. Line-by-line processing with constant memory
  const lines = createInterface({ input: createReadStream("results.csv"), crlfDelay: Infinity });
  let total = 0;
  let count = 0;
  for await (const line of lines) {
    const [, score] = line.split(",");
    if (score && !Number.isNaN(Number(score))) {
      total += Number(score);
      count++;
    }
  }
  console.log("Average:", (total / count).toFixed(1));

  // 2. Transform + gzip pipeline
  const upper = new Transform({
    transform(chunk: Buffer, _enc, callback) {
      callback(null, chunk.toString("utf8").toUpperCase());
    },
  });

  await pipeline(
    createReadStream("results.csv"),
    upper,
    createGzip(),
    createWriteStream("results.upper.csv.gz"),
  );
  console.log("Wrote results.upper.csv.gz");
}

main().catch((err) => {
  console.error(err);
  process.exitCode = 1;
});

Key points

  • Use streams whenever input size is unbounded or large.
  • pipeline wires streams together with proper error handling and cleanup.
  • Readable streams are async iterables — for await is often the simplest consumer.

Exercise

Generate a CSV with 1,000,000 rows of random student scores using a write stream (respect backpressure by awaiting the drain event). Then stream it back to compute the average and the grade distribution, logging process.memoryUsage().heapUsed to confirm memory stays low.

Show solution

Try the exercise yourself first — then compare your approach with this one.

When write returns false, the stream's buffer is full: wait for the drain event before writing more. That's backpressure, and it keeps memory flat even for a million rows. Reading back uses readline line by line, so the file is never fully in memory either.

million.tsTypeScript
import { createReadStream, createWriteStream } from "node:fs";
import { once } from "node:events";
import { createInterface } from "node:readline";

const FILE = "scores.csv";
const ROWS = 1_000_000;

async function generate(): Promise<void> {
  const out = createWriteStream(FILE);
  out.write("student,score\n");
  for (let i = 1; i <= ROWS; i++) {
    const ok = out.write("S" + i + "," + Math.floor(Math.random() * 101) + "\n");
    if (!ok) await once(out, "drain");          // respect backpressure
  }
  out.end();
  await once(out, "finish");
}

function grade(score: number): string {
  return score >= 75 ? "A" : score >= 65 ? "B" : score >= 45 ? "C" : score >= 30 ? "D" : "F";
}

async function analyse(): Promise<void> {
  const lines = createInterface({ input: createReadStream(FILE), crlfDelay: Infinity });
  let total = 0;
  let count = 0;
  const distribution: Record<string, number> = { A: 0, B: 0, C: 0, D: 0, F: 0 };

  for await (const line of lines) {
    const score = Number(line.split(",")[1]);
    if (Number.isNaN(score)) continue;          // header row
    total += score;
    count++;
    distribution[grade(score)]++;
  }

  console.log("Rows:", count, "average:", (total / count).toFixed(2));
  console.log(distribution);
  console.log("Heap used:", Math.round(process.memoryUsage().heapUsed / 1024 / 1024), "MB");
}

await generate();
await analyse();

Check your understanding

  1. Why use a stream instead of readFile for a 5 GB log file?

  2. What does it mean when writable.write() returns false?

  3. What advantage does pipeline() from node:stream/promises have over chaining .pipe()?

  4. Which stream type modifies data as it passes through?

Ask AI