mitmario.dev

Streams verbinden

Node.js Sandbox 3 Min Lesezeit 3 BeispieleLektion 5 von 7

Ein Stream allein ist selten das Ziel. Interessant wird es, wenn mehrere hintereinanderhängen: lesen, umformen, schreiben. Genau dafür gibt es pipeline.

Zwei Ströme mit pipeline verbinden
import { writeFileSync, createReadStream, createWriteStream, readFileSync } from "node:fs";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";

writeFileSync("quelle.txt", "erste zeile\nzweite zeile\n");

const grossschreiben = new Transform({
  transform(stueck, kodierung, weiter) {
    weiter(null, stueck.toString("utf8").toUpperCase());
  },
});

await pipeline(
  createReadStream("quelle.txt"),
  grossschreiben,
  createWriteStream("ziel.txt"),
);

console.log(readFileSync("ziel.txt", "utf8").trim());

pipeline aus node:stream/promises nimmt beliebig viele Ströme entgegen: zuerst einen lesbaren, dann null oder mehr Umformer, zuletzt einen schreibbaren. Es verbindet sie, wartet, bis alles durch ist, und gibt dir ein Promise zurück, das du mit await abholst.

Ein Transform in seiner kleinsten Form ist eine Funktion mit drei Parametern: das Stück, seine Kodierung, und ein weiter, das du aufrufst, wenn du fertig bist. Der erste Parameter von weiter ist ein Fehler, der zweite das Ergebnis. Fehler zuerst, wie überall in Node.

Warum nicht .pipe()

In älterem Code steht meistens .pipe(). Das funktioniert und ist kürzer, hat aber ein Problem, das man erst im Fehlerfall sieht.

pipeline räumt hinter sich auf
import { writeFileSync, createReadStream, createWriteStream } from "node:fs";
import { Transform } from "node:stream";
import { pipeline } from "node:stream/promises";

writeFileSync("quelle.txt", "erste zeile\n");

const kaputt = new Transform({
  transform(stueck, kodierung, weiter) {
    weiter(new Error("Umformung fehlgeschlagen"));
  },
});

const ziel = createWriteStream("ziel.txt");

try {
  await pipeline(createReadStream("quelle.txt"), kaputt, ziel);
} catch (fehler) {
  console.log(`abgebrochen: ${fehler.message}`);
}

console.log(`Ziel aufgeraeumt: ${ziel.destroyed}`);

Der Umformer scheitert. pipeline wirft, du fängst es, und wichtiger: Der Schreibstrom ist geschlossen. Kein offenes Dateihandle, keine halb geschriebene Datei, die weiter im Weg liegt.

pipe lässt den Schreibstrom offen
import { writeFileSync, createReadStream, createWriteStream } from "node:fs";
import { Transform } from "node:stream";

writeFileSync("quelle.txt", "erste zeile\n");

const kaputt = new Transform({
  transform(stueck, kodierung, weiter) {
    weiter(new Error("Umformung fehlgeschlagen"));
  },
});

const ziel = createWriteStream("ziel.txt");

createReadStream("quelle.txt").pipe(kaputt).pipe(ziel);

kaputt.on("error", (fehler) => {
  console.log(`abgebrochen: ${fehler.message}`);
  console.log(`Ziel aufgeraeumt: ${ziel.destroyed}`);
});

Derselbe Fehler, .pipe() statt pipeline, und der Schreibstrom lebt weiter. Das Dateihandle bleibt offen, bis der Prozess endet.

Bei einem Skript merkt das niemand. Auf einem Server, der das pro Anfrage macht, sammeln sich diese Handles, und irgendwann steht in den Logs EMFILE: too many open files. Es ist eine der unangenehmeren Fehlersuchen, weil die Meldung an einer ganz anderen Stelle auftaucht als die Ursache.

Dazu kommt ein zweiter Punkt: .pipe() reicht Fehler nicht weiter. Jeder Strom der Kette braucht sein eigenes on("error", ...), sonst gilt die Regel aus Lektion 7.3 und der Prozess ist weg.

Nimm pipeline. Es ist der Grund, warum es das gibt.

Gegendruck, in drei Sätzen

Was passiert eigentlich, wenn der Empfänger langsamer ist als der Sender? Beim Lesen einer Datei von SSD und Schreiben über eine langsame Netzwerkverbindung ist das der Normalfall.

Ein Stream regelt das selbst. Der Schreibstrom meldet, wenn sein Puffer voll ist, der Lesestrom pausiert daraufhin, und sobald wieder Platz ist, geht es weiter. Das nennt sich Gegendruck.

Der Punkt daran ist nicht, wie es funktioniert, sondern was passiert, wenn man es selbst nachbaut und vergisst: Dann läuft der Puffer voll, der Speicherverbrauch steigt, und am Ende hat man doch wieder die ganze Datei im Arbeitsspeicher. Genau davor wollte man ja weg.

Wo dir das überall begegnet

Das Muster steckt an vielen Stellen, an denen man es nicht sofort vermutet.

Eine Datei komprimieren ist ein createReadStream, ein createGzip aus node:zlib und ein createWriteStream. Ein Datei-Upload ist ein Lesestrom aus der Anfrage in einen Schreibstrom auf die Platte. Und die Antwort eines HTTP-Servers ist selbst ein Schreibstrom, weshalb du eine große Datei direkt hineinleiten kannst, ohne sie vorher zu laden.

Das letzte kommt in Abschnitt 8 wieder, wenn der erste Server steht. Dann ist pipeline plötzlich kein Stream-Thema mehr, sondern der kürzeste Weg, eine Datei auszuliefern.

Zum Mitnehmen

pipeline räumt bei einem Fehler alle beteiligten Ströme auf. .pipe() lässt sie offen liegen, und genau deshalb nimmt man es nicht mehr.

Jetzt du

Basis Konto, kostenlos

Zu dieser Lektion gehört eine Aufgabe. Du schreibst den Code selbst, und nach jedem Lauf sagt dir eine Prüfliste, was schon stimmt.

Dafür brauchst du das Basis Konto. Es kostet nichts, und ein Passwort gibt es auch nicht.

In diesem Kurs läuft dein Code auf einem Server. Dafür hat das Basis Konto 1 Stunde im Monat, mehr Zeit gibt es mit dem Premium Konto.

Was in dieser Lektion steckt

  • Artikel mit 3 Beispielen zum Ausprobieren

    Steht hier, ohne Konto lesbar.

  • Aufgabe, dein Code läuft auf einem Server

    Öffnet sich mit dem Basis Konto.