Verständnis von Node.js Streams: Das Problem, das sie tatsächlich lösen
Erfahren Sie, warum Node.js-Streams existieren, wie das Piping intern funktioniert und was Backpressure tatsächlich bedeutet, um große Daten effizient zu verarbeiten.
Die meisten Tutorials zu Streams beginnen mit den API-Elementen .pipe(), Readable, Writable und Transform, ohne jemals das eigentliche Problem anzusprechen, für dessen Lösung diese Werkzeuge entwickelt wurden. Auf diese Weise wirken Streams wie willkürliche Regeln, die man auswendig lernen muss. Wenn man den Ansatz umkehrt und mit dem Problem beginnt, erklärt sich die API meist von selbst.
Das Problem: Einige Daten sind zu groß, um sie auf einmal zu speichern
Stellen Sie sich eine Aufgabe vor, bei der ein 4GB großer CSV-Export in JSON umgewandelt werden muss. Der naive Ansatz sieht so aus:
const fs = require("fs");
const data = fs.readFileSync("export.csv", "utf8");
const rows = data.split("\n").map(parseRow);
fs.writeFileSync("export.json", JSON.stringify(rows));
readFileSync gibt den Kontrollfluss an das Programm erst zurück, nachdem die gesamte 4GB große Datei als einziger JavaScript-String in den Speicher geladen wurde, und diese String-Repräsentation verbraucht in der Regel deutlich mehr Speicher, als die reine Dateigröße vermuten ließe. Auf einem Rechner mit nur 2GB verfügbarem RAM läuft dieser Code nicht nur ineffizient, er stürzt sogar ab oder lässt alle anderen Prozesse, die um denselben Speicherpool konkurrieren, hungern. Mit der Logik an sich ist alles in Ordnung; alle Umwandlungsschritte sind korrekt. Der eigentliche Fehler liegt in einer versteckten Annahme innerhalb von readFileSync: nämlich dass es in Ordnung ist, eine ganze Datei unabhängig von ihrer Größe gleichzeitig im Speicher zu halten.
Die eigentliche Idee hinter einem Stream
Anstatt die gesamte Datensammlung von vornherein anzufordern, arbeitet ein Stream nach einem anderen Prinzip: Zuerst wird ein Teil angefordert, dieser verarbeitet, anschließend der nächste Teil. Zu keinem Zeitpunkt befindet sich die gesamte Datei im Speicher. Immer nur ein kleiner Teil existiert, und dieser wird verarbeitet sowie freigegeben, bevor der nächste Teil eintrifft.
const fs = require("fs");
const readStream = fs.createReadStream("export.csv", { encoding: "utf8" });
readStream.on("data", (chunk) => {
console.log(`Received ${chunk.length} characters`);
});
readStream.on("end", () => {
console.log("Done reading the whole file, piece by piece");
});
Beachten Sie, dass die Gesamtkapazität der Datei in diesem Code keinerlei Rolle spielt. Egal, ob die Quelldatei 4 GB oder 4 KB groß ist – dieser Codeabschnitt läuft identisch ab, obwohl der maximale Speicherverbrauch in den beiden Fällen erheblich unterschiedlich ist. Genau das ist der Sinn des Streamings: Anstelle der Erfordernis einer vollständigen Datensammlung vor dem Start wird die Möglichkeit genutzt, sofort mit der Verarbeitung zu beginnen und die Aufgabe abzuschließen, ohne jemals mehr als einen sehr kleinen Datensatz im Speicher zu halten.
Piping: Eine Quelle direkt mit einem Ziel verbinden
Das Manuelle Herunterladen von Datenblöcken nacheinander ist ein nützliches Grundprinzip, doch in der Praxis ist es weitaus üblicher, einen lesbaren Stream direkt mit einem schreibbaren zu verbinden, sodass die Daten automatisch vom Ausgangspunkt zum Zielort fließen können, anstatt jeden Block manuell weiterzuleiten:
const fs = require("fs");
fs.createReadStream("export.csv")
.pipe(fs.createWriteStream("export-copy.csv"));
Zwei Zeilen reichen aus, um eine Datei beliebiger Größe zu kopieren, ohne sie vollständig in den Speicher laden zu müssen. .pipe() bewirkt hier keine Magie – es leitet lediglich das data-Event, das von der Quelle ausgesendet wird, in einen write-Aufruf am Zielort um und kümmert sich außerdem um ein weiteres Detail, das sich als wichtiger erweist als der Kopiervorgang selbst.
Der Teil, den fast niemand richtig erklärt: Rückdruck
Dies ist das eigentliche Problem, dem .pipe() begegnet – jenseits bloßer Bequemlichkeit. Stellen Sie sich eine Situation vor, in der das Lesen schnell erfolgt, beispielsweise von einer lokalen Festplatte, während das Schreiben langsam abläuft, möglicherweise über eine Netzwerkverbindung mit begrenztem Bandbreitenangebot.
readStream.on("data", (chunk) => {
writeStream.write(chunk); // what happens if this can't keep up?
});
Wenn .write() schneller aufgerufen wird, als das Ziel tatsächlich Daten aufnehmen kann, führt dies weder dazu, dass der schreibbare Stream einen Fehler auslöst noch die Ausführung blockiert. Stattdessen sammelt er den Überschuss stillschweigend in einem internen Speicherpuffer an und hält ihn dort, bis er die Gelegenheit hat, ihn auszuspeichern. Wenn dieses Ungleichgewicht lange anhält und der Unterschied zwischen Lesegeschwindigkeit und Schreibgeschwindigkeit groß genug ist, entsteht das genau gleiche Problem, das Streaming eigentlich beseitigen sollte: Der Speicherverbrauch steigt unbegrenzt an – lediglich verzögert statt sofort – und ist in der Regel viel schwieriger zu erkennen, bevor ein Absturz ausgelöst wird.
Backpressure ist der Schutz vor genau diesem Fehlermodus: Sie gibt einem schreibbaren Stream eine Möglichkeit, anzuzeigen, dass er seine Kapazität erreicht hat und der Produzent langsamer vorgehen muss. Ein ordnungsgemäß implementierter Produzent achtet auf dieses Signal und stoppt die Übertragung anstelle weiterer Datenversendungen.
readStream.on("data", (chunk) => {
const canContinue = writeStream.write(chunk);
if (!canContinue) {
readStream.pause(); // stop reading until the writable side catches up
}
});
writeStream.on("drain", () => {
readStream.resume(); // writable side is ready for more
});
Sobald der interne Puffer eines schreibbaren Streams seine konfigurierte Grenze überschreitet, gibt .write() false zurück – das dient als Signal, die Produktion zu stoppen, bis der Stream ein drain-Event auslöst, das anzeigt, dass der Rückstand beseitigt wurde und er wieder mehr Daten entgegennehmen kann. Dieses Muster des Anhaltens und Wiederaufnahmens wird genau von .pipe() im Hintergrund automatisch gesteuert:
readStream.pipe(writeStream); // handles backpressure for you, silently, correctly
Das ist, und nicht die Kürze, der eigentliche Grund dafür, .pipe() gegenüber dem manuellen Verbinden von data- und write-Listenern vorzuziehen. Das Manuell-Weiterleiten von Datenblöcken unter Ignorieren des Rückgabewerts von .write() führt erneut zu dem Problem der unbegrenzten Speichernutzung, das Streams ursprünglich vermeiden sollten – es handelt sich dabei lediglich um einen Schritt entfernt vom offensichtlichen Fehler, der in readFileSync bereits vorhanden ist.
Transform-Streams: Daten während der Verarbeitung bearbeiten
In einigen Fällen reicht es nicht aus, Daten unverändert von einem Ort zum anderen zu übertragen – man muss sie unterwegs auch umformen. Ein Transform-Stream wurde genau dafür konzipiert: Er befindet sich in der Mitte einer Pipe-Kette, nimmt eingehende Datenblöcke entgegen, wendet an jedem eine bestimmte Operation an und leitet das Ergebnis an den nächsten Schritt weiter:
const { Transform } = require("stream");
const upperCaseTransform = new Transform({
transform(chunk, encoding, callback) {
callback(null, chunk.toString().toUpperCase());
},
});
fs.createReadStream("input.txt")
.pipe(upperCaseTransform)
.pipe(fs.createWriteStream("output.txt"));
Jeder Datenblock wird während seines Durchlaufs durch die Pipeline in Großbuchstaben umgewandelt, und zu keinem Zeitpunkt muss die gesamte Datei – weder in ihrer ursprünglichen noch in der umgewandelten Form – vollständig im Speicher vorhanden sein. Genau dieses Prinzip steckt hinter Node’s eingebauter Funktion zlib.createGzip(). Es handelt sich dabei lediglich um einen Transform-Stream, der jeden ankommenden Datenblock komprimiert, und wie jede Transformation kann er genauso wie im Beispiel mit den Großbuchstaben direkt in eine Pipeline-Kette eingefügt werden:
const zlib = require("zlib");
fs.createReadStream("export.csv")
.pipe(zlib.createGzip())
.pipe(fs.createWriteStream("export.csv.gz"));
Lesen, Komprimieren und Schreiben finden gleichzeitig statt, wobei jeweils nur kleine Datenschnitte verarbeitet werden; dabei muss zu keinem Zeitpunkt die gesamte Datei gleichzeitig geladen werden.
Warum es sich lohnt, das wirklich zu verstehen
Streams gelten als einer der unhandlicheren Bereiche der Node-API-Oberfläche, und ehrlich gesagt wirkt die rohe, ereignisgesteuerte Schnittstelle tatsächlich unpraktisch – daher ist dieses Image nicht völlig unbegründet. Doch das zugrundeliegende Konzept ist einfach: Man vermeidet es, den gesamten Datensatz in den Speicher zu laden, arbeitet stattdessen Schritt für Schritt damit und sorgt dafür, dass ein schneller Produzent dabei niemals einen langsamen Verbraucher unbemerkt überfluten kann. Sobald man dieses Konzept als grundlegendes Modell verinnerlicht hat, wirken .pipe(), Transform sowie Backpressure nicht mehr wie drei unabhängige APIs, die man getrennt auswendig lernen muss. Tatsächlich handelt es sich dabei um eine einzige Idee, die über drei miteinander verbundene Bestandteile der API zum Vorschein kommt und ein Problem löst, mit dem sich Funktionen wie readFileSync niemals auseinandersetzen mussten, da sie ursprünglich nie dazu gedacht waren, etwas anderes als kleine Dateien zu verarbeiten.
Verwandte Literatur
- Node.js Streams jenseits der Grundlagen: Speicher, Rückdruck und echte Fehler — Erfahren Sie, wie Node.js Streams mit Web Streams interagieren, welche Speichereinsparungen durch echte Benchmarks erzielt werden können und welche Produktionsfehler nur unter Last auftreten.
- Erklärung zu Node.js Streams: Behebung von Speicherfehlern bei Dateien — Erfahren Sie, warum das Laden ganzer Dateien in den Speicher zu Abstürzen von Node.js-Servern führt und wie lesbare, schreibbare, duplex- sowie transformierende Streams dies mithilfe von Rückdruck beheben.