Kodokon kodokon.com

Streams und Backpressure: große Dateien verarbeiten

Verarbeite Dateien von mehreren Gigabyte mit konstantem Speicherbedarf - dank Streams, pipeline und Backpressure.

9 Min. · 3 Fragen

Diese Lektion in Kodokon öffnen

readFile lädt eine ganze Datei in den Speicher: Bei einem Log von 4 GB fliegt dir dein Prozess um die Ohren. Streams verarbeiten die Daten in Stücken mit konstantem Speicherbedarf, begrenzt durch highWaterMark (standardmäßig 64 KiB bei einem Datei-Stream). Es gibt vier Typen: Readable, Writable, Duplex und Transform. Das zentrale Problem, das sie lösen, heißt Backpressure: Was tust du, wenn der Produzent schneller ist als der Konsument?

JAVASCRIPT
import { createWriteStream } from "node:fs";

const out = createWriteStream("big.txt");
let i = 0;

function writeChunks() {
  let ok = true;
  while (i < 1e6 && ok) {
    ok = out.write(`line ${i}\n`);
    i += 1;
  }
  if (i < 1e6) out.once("drain", writeChunks);
  else out.end();
}

writeChunks();
Der write/drain-Vertrag, von Hand umgesetzt.

Wenn write() false zurückgibt, hat der interne Puffer highWaterMark überschritten: Die Daten sind nicht verloren, aber weiterzuschreiben häuft alles im Speicher an - genau das, was Streams vermeiden sollen. Der Vertrag: aufhören zu schreiben und auf das Ereignis drain warten. Genau das erledigt pipe() für dich, und pipeline() macht es noch besser: Zusätzlich zum Backpressure leitet es Fehler weiter und zerstört alle Streams sauber, in beide Richtungen.

JAVASCRIPT
import { createReadStream, createWriteStream }
  from "node:fs";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";

const upper = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, chunk.toString().toUpperCase());
  },
});

await pipeline(
  createReadStream("big.txt"),
  upper,
  createWriteStream("big-upper.txt"),
);
console.log("done");
Transformation im Stream, konstanter Speicherbedarf.

Readable-Streams sind asynchron iterierbar: for await konsumiert die Stücke und respektiert dabei den Backpressure. Achtung, ein Stück ist ein beliebiger Ausschnitt aus Bytes: Nichts garantiert, dass es an einem Zeilenende aufhört. Für eine zeilenweise Verarbeitung kümmert sich node:readline um das Zusammensetzen der Teile, inklusive eines \r\n, das auf zwei Stücke verteilt ist - dank crlfDelay: Infinity.

JAVASCRIPT
import { createReadStream } from "node:fs";
import { createInterface } from "node:readline";

const rl = createInterface({
  input: createReadStream("big.txt"),
  crlfDelay: Infinity,
});

let count = 0;
for await (const line of rl) {
  if (line.includes("42")) count += 1;
}
console.log(count, "matching lines");
Eine große Datei Zeile für Zeile lesen.

Wissenscheck

Stelle sicher, dass du die wichtigsten Punkte dieser Lektion behalten hast.

  1. Was bedeutet es, wenn writable.write(chunk) false zurückgibt?
    • Das Stück ist verloren und muss neu geschrieben werden
    • Der interne Puffer überschreitet highWaterMark: nicht mehr schreiben bis zum Ereignis drain
    • Der Stream wurde durch einen Fehler zerstört
  2. Welchen entscheidenden Vorteil hat pipeline() gegenüber pipe()?
    • pipeline ist schneller, weil es mehrere Threads nutzt
    • pipe kümmert sich um den Backpressure, pipeline nicht
    • pipeline leitet Fehler weiter und zerstört alle beteiligten Streams
    • pipeline wandelt Kodierungen automatisch um
  3. Was stellt highWaterMark im objectMode dar?
    • Eine Anzahl wartender Objekte
    • Eine Anzahl Bytes, wie im Binärmodus
    • Die maximale Größe eines serialisierten Objekts