Node.js v26.0.0 Dokumentation
- Node.js v26.0.0
- Inhaltsverzeichnis
- Stream
- Gliederung dieses Dokuments
- Arten von Streams
- API für Stream-Consumer
- Schreibbare Streams (Writable)
- Klasse:
stream.Writable- Ereignis:
'close' - Ereignis:
'drain' - Ereignis:
'error' - Ereignis:
'finish' - Ereignis:
'pipe' - Ereignis:
'unpipe' writable.cork()writable.destroy([error])writable.closedwritable.destroyedwritable.end([chunk[, encoding]][, callback])writable.setDefaultEncoding(encoding)writable.uncork()writable.writablewritable.writableAbortedwritable.writableEndedwritable.writableCorkedwritable.erroredwritable.writableFinishedwritable.writableHighWaterMarkwritable.writableLengthwritable.writableNeedDrainwritable.writableObjectModewritable[Symbol.asyncDispose]()writable.write(chunk[, encoding][, callback])
- Ereignis:
- Klasse:
- Lesbare Streams (Readable)
- Zwei Lesemodi
- Drei Zustände
- Wählen Sie einen API-Stil
- Klasse:
stream.Readable- Ereignis:
'close' - Ereignis:
'data' - Ereignis:
'end' - Ereignis:
'error' - Ereignis:
'pause' - Ereignis:
'readable' - Ereignis:
'resume' readable.destroy([error])readable.closedreadable.destroyedreadable.isPaused()readable.pause()readable.pipe(destination[, options])readable.read([size])readable.readablereadable.readableAbortedreadable.readableDidReadreadable.readableEncodingreadable.readableEndedreadable.erroredreadable.readableFlowingreadable.readableHighWaterMarkreadable.readableLengthreadable.readableObjectModereadable.resume()readable.setEncoding(encoding)readable.unpipe([destination])readable.unshift(chunk[, encoding])readable.wrap(stream)readable[Symbol.asyncIterator]()readable[Symbol.asyncDispose]()readable.compose(stream[, options])readable.iterator([options])readable.map(fn[, options])readable.filter(fn[, options])readable.forEach(fn[, options])readable.toArray([options])readable.some(fn[, options])readable.find(fn[, options])readable.every(fn[, options])readable.flatMap(fn[, options])readable.drop(limit[, options])readable.take(limit[, options])readable.reduce(fn[, initial[, options]])
- Ereignis:
- Duplex- und Transform-Streams
stream.finished(stream[, options], callback)stream.pipeline(source[, ...transforms], destination, callback)stream.pipeline(streams, callback)stream.compose(...streams)stream.isErrored(stream)stream.isReadable(stream)stream.isWritable(stream)stream.Readable.from(iterable[, options])stream.Readable.fromWeb(readableStream[, options])stream.Readable.isDisturbed(stream)stream.Readable.toWeb(streamReadable[, options])stream.Writable.fromWeb(writableStream[, options])stream.Writable.toWeb(streamWritable)stream.Duplex.from(src)stream.Duplex.fromWeb(pair[, options])stream.Duplex.toWeb(streamDuplex[, options])stream.addAbortSignal(signal, stream)stream.getDefaultHighWaterMark(objectMode)stream.setDefaultHighWaterMark(objectMode, value)
- Schreibbare Streams (Writable)
- API für Stream-Implementierer
- Vereinfachte Konstruktion
- Implementierung eines schreibbaren Streams
new stream.Writable([options])writable._construct(callback)writable._write(chunk, encoding, callback)writable._writev(chunks, callback)writable._destroy(err, callback)writable._final(callback)- Fehler beim Schreiben
- Beispiel für einen schreibbaren Stream
- Dekodieren von Puffern in einem schreibbaren Stream
- Implementierung eines lesbaren Streams
- Implementierung eines Duplex-Streams
- Implementierung eines Transform-Streams
- Zusätzliche Hinweise
- Stream
- Index
- Über diese Dokumentation
- Verwendung und Beispiel
- Assertion-Tests
- Asynchrone Kontextverfolgung
- Async-Hooks
- Buffer
- C++-Add-ons
- C/C++-Add-ons mit Node-API
- C++-Embedder-API
- Kindprozesse
- Cluster
- Kommandozeilenoptionen
- Konsole
- Crypto
- Debugger
- Veraltete APIs
- Diagnostics Channel
- DNS
- Domain
- Umgebungsvariablen
- Fehler
- Events
- Dateisystem
- Globale Objekte
- HTTP
- HTTP/2
- HTTPS
- Inspector
- Internationalisierung
- Module: CommonJS-Module
- Module: ECMAScript-Module
- Module:
node:module-API - Module: Pakete
- Module: TypeScript
- Net
- Iterable Streams API
- OS
- Path
- Performance-Hooks
- Berechtigungen
- Process
- Punycode
- Query-Strings
- Readline
- REPL
- Report
- Einzelne ausführbare Anwendungen
- SQLite
- Stream
- String-Decoder
- Test-Runner
- Timer
- TLS/SSL
- Trace-Events
- TTY
- UDP/Datagramm
- URL
- Utilities
- V8
- VM
- WASI
- Web Crypto API
- Web Streams API
- Worker-Threads
- Zlib
- Zlib Iterable Compression
- Andere Versionen
- Optionen
Stream#
Stabilität: 2 - Stabil
Ein Stream ist eine abstrakte Schnittstelle für die Arbeit mit Streaming-Daten in Node.js. Das node:stream-Modul stellt eine API zur Implementierung der Stream-Schnittstelle bereit.
Es gibt viele von Node.js bereitgestellte Stream-Objekte. Zum Beispiel sind eine Anfrage an einen HTTP-Server und process.stdout beides Stream-Instanzen.
Streams können lesbar, schreibbar oder beides sein. Alle Streams sind Instanzen des EventEmitter.
Zugriff auf das node:stream-Modul
const stream = require('node:stream');
Das node:stream-Modul ist nützlich für die Erstellung neuer Arten von Stream-Instanzen. Es ist normalerweise nicht notwendig, das node:stream-Modul zu verwenden, um Streams zu konsumieren.
Gliederung dieses Dokuments#
Dieses Dokument enthält zwei Hauptabschnitte und einen dritten Abschnitt für Hinweise. Der erste Abschnitt erläutert, wie bestehende Streams innerhalb einer Anwendung verwendet werden. Der zweite Abschnitt erklärt, wie man neue Arten von Streams erstellt.
Arten von Streams#
Es gibt vier grundlegende Stream-Typen innerhalb von Node.js
Writable: Streams, in die Daten geschrieben werden können (z. B.fs.createWriteStream()).Readable: Streams, aus denen Daten gelesen werden können (z. B.fs.createReadStream()).Duplex: Streams, die sowohlReadableals auchWritablesind (z. B.net.Socket).Transform:Duplex-Streams, die die Daten beim Schreiben und Lesen modifizieren oder transformieren können (z. B.zlib.createDeflate()).
Zusätzlich enthält dieses Modul die Hilfsfunktionen stream.duplexPair(), stream.pipeline(), stream.finished(), stream.Readable.from() und stream.addAbortSignal().
Streams Promises-API#
Die stream/promises-API bietet einen alternativen Satz asynchroner Hilfsfunktionen für Streams, die Promise-Objekte zurückgeben, anstatt Callbacks zu verwenden. Die API ist über require('node:stream/promises') oder require('node:stream').promises zugänglich.
stream.pipeline(streams[, options])#
stream.pipeline(source[, ...transforms], destination[, options])#
streams<Stream[]>|<Iterable[]>|<AsyncIterable[]>|<Function[]>source<Stream>|<Iterable>|<AsyncIterable>|<Function>- Rückgabe:
<Promise>|<AsyncIterable>
- Rückgabe:
...transforms<Stream>|<Function>source<AsyncIterable>- Rückgabe:
<Promise>|<AsyncIterable>
destination<Stream>|<Function>source<AsyncIterable>- Rückgabe:
<Promise>|<AsyncIterable>
options<Object>Pipeline-Optionensignal<AbortSignal>end<boolean>Beende den Ziel-Stream, wenn der Quell-Stream endet. Transform-Streams werden immer beendet, selbst wenn dieser Wertfalseist. Standard:true.
- Rückgabe:
<Promise>Wird erfüllt, wenn die Pipeline abgeschlossen ist.
const { pipeline } = require('node:stream/promises'); const fs = require('node:fs'); const zlib = require('node:zlib'); async function run() { await pipeline( fs.createReadStream('archive.tar'), zlib.createGzip(), fs.createWriteStream('archive.tar.gz'), ); console.log('Pipeline succeeded.'); } run().catch(console.error);import { pipeline } from 'node:stream/promises'; import { createReadStream, createWriteStream } from 'node:fs'; import { createGzip } from 'node:zlib'; await pipeline( createReadStream('archive.tar'), createGzip(), createWriteStream('archive.tar.gz'), ); console.log('Pipeline succeeded.');
Um ein AbortSignal zu verwenden, übergeben Sie es innerhalb eines Optionen-Objekts als letztes Argument. Wenn das Signal abgebrochen wird, wird destroy für die zugrunde liegende Pipeline mit einem AbortError aufgerufen.
const { pipeline } = require('node:stream/promises'); const fs = require('node:fs'); const zlib = require('node:zlib'); async function run() { const ac = new AbortController(); const signal = ac.signal; setImmediate(() => ac.abort()); await pipeline( fs.createReadStream('archive.tar'), zlib.createGzip(), fs.createWriteStream('archive.tar.gz'), { signal }, ); } run().catch(console.error); // AbortErrorimport { pipeline } from 'node:stream/promises'; import { createReadStream, createWriteStream } from 'node:fs'; import { createGzip } from 'node:zlib'; const ac = new AbortController(); const { signal } = ac; setImmediate(() => ac.abort()); try { await pipeline( createReadStream('archive.tar'), createGzip(), createWriteStream('archive.tar.gz'), { signal }, ); } catch (err) { console.error(err); // AbortError }
Die pipeline-API unterstützt auch asynchrone Generatoren
const { pipeline } = require('node:stream/promises'); const fs = require('node:fs'); async function run() { await pipeline( fs.createReadStream('lowercase.txt'), async function* (source, { signal }) { source.setEncoding('utf8'); // Work with strings rather than `Buffer`s. for await (const chunk of source) { yield await processChunk(chunk, { signal }); } }, fs.createWriteStream('uppercase.txt'), ); console.log('Pipeline succeeded.'); } run().catch(console.error);import { pipeline } from 'node:stream/promises'; import { createReadStream, createWriteStream } from 'node:fs'; await pipeline( createReadStream('lowercase.txt'), async function* (source, { signal }) { source.setEncoding('utf8'); // Work with strings rather than `Buffer`s. for await (const chunk of source) { yield await processChunk(chunk, { signal }); } }, createWriteStream('uppercase.txt'), ); console.log('Pipeline succeeded.');
Denken Sie daran, das an den asynchronen Generator übergebene signal-Argument zu verarbeiten. Besonders in dem Fall, in dem der asynchrone Generator die Quelle für die Pipeline (d. h. das erste Argument) ist oder die Pipeline niemals abgeschlossen wird.
const { pipeline } = require('node:stream/promises'); const fs = require('node:fs'); async function run() { await pipeline( async function* ({ signal }) { await someLongRunningfn({ signal }); yield 'asd'; }, fs.createWriteStream('uppercase.txt'), ); console.log('Pipeline succeeded.'); } run().catch(console.error);import { pipeline } from 'node:stream/promises'; import fs from 'node:fs'; await pipeline( async function* ({ signal }) { await someLongRunningfn({ signal }); yield 'asd'; }, fs.createWriteStream('uppercase.txt'), ); console.log('Pipeline succeeded.');
Die pipeline-API bietet eine Callback-Version
stream.finished(stream[, options])#
stream<Stream>|<ReadableStream>|<WritableStream>Ein lesbarer und/oder schreibbarer Stream/Webstream.options<Object>error<boolean>|<undefined>readable<boolean>|<undefined>writable<boolean>|<undefined>signal<AbortSignal>|<undefined>cleanup<boolean>|<undefined>Wenntrue, werden die von dieser Funktion registrierten Listener entfernt, bevor das Promise erfüllt wird. Standard:false.
- Rückgabe:
<Promise>Wird erfüllt, wenn der Stream nicht mehr lesbar oder schreibbar ist.
const { finished } = require('node:stream/promises'); const fs = require('node:fs'); const rs = fs.createReadStream('archive.tar'); async function run() { await finished(rs); console.log('Stream is done reading.'); } run().catch(console.error); rs.resume(); // Drain the stream.import { finished } from 'node:stream/promises'; import { createReadStream } from 'node:fs'; const rs = createReadStream('archive.tar'); async function run() { await finished(rs); console.log('Stream is done reading.'); } run().catch(console.error); rs.resume(); // Drain the stream.
Die finished-API bietet auch eine Callback-Version.
stream.finished() hinterlässt hängende Ereignis-Listener (insbesondere 'error', 'end', 'finish' und 'close'), nachdem das zurückgegebene Promise aufgelöst oder abgelehnt wurde. Der Grund dafür ist, dass unerwartete 'error'-Ereignisse (aufgrund fehlerhafter Stream-Implementierungen) keine unerwarteten Abstürze verursachen sollen. Wenn dies unerwünscht ist, sollte options.cleanup auf true gesetzt werden.
await finished(rs, { cleanup: true });
Objektmodus#
Alle von Node.js-APIs erstellten Streams operieren ausschließlich auf Strings, <Buffer>, <TypedArray> und <DataView> Objekten.
StringsundBufferssind die am häufigsten mit Streams verwendeten Typen.TypedArrayundDataViewermöglichen die Verarbeitung binärer Daten mit Typen wieInt32ArrayoderUint8Array. Wenn Sie ein TypedArray oder DataView in einen Stream schreiben, verarbeitet Node.js die Rohbytes.
Es ist jedoch möglich, dass Stream-Implementierungen mit anderen Arten von JavaScript-Werten arbeiten (mit Ausnahme von null, das innerhalb von Streams eine besondere Bedeutung hat). Solche Streams gelten als im "Objektmodus" arbeitend.
Stream-Instanzen werden bei ihrer Erstellung mithilfe der objectMode-Option in den Objektmodus geschaltet. Der Versuch, einen bestehenden Stream in den Objektmodus zu schalten, ist nicht sicher.
Pufferung#
Sowohl Writable- als auch Readable-Streams speichern Daten in einem internen Puffer.
Die Menge der potenziell gepufferten Daten hängt von der highWaterMark-Option ab, die an den Konstruktor des Streams übergeben wird. Bei normalen Streams gibt die highWaterMark-Option eine Gesamtanzahl von Bytes an. Bei Streams, die im Objektmodus arbeiten, gibt die highWaterMark eine Gesamtanzahl von Objekten an. Bei Streams, die auf Strings operieren (aber diese nicht dekodieren), gibt die highWaterMark eine Gesamtanzahl von UTF-16-Code-Einheiten an.
Daten werden in Readable-Streams gepuffert, wenn die Implementierung stream.push(chunk) aufruft. Wenn der Consumer des Streams stream.read() nicht aufruft, verbleiben die Daten in der internen Warteschlange, bis sie konsumiert werden.
Sobald die Gesamtgröße des internen Lesepuffers den durch highWaterMark angegebenen Schwellenwert erreicht, stoppt der Stream vorübergehend das Lesen von Daten aus der zugrunde liegenden Ressource, bis die aktuell gepufferten Daten konsumiert werden können (das heißt, der Stream hört auf, die interne readable._read()-Methode aufzurufen, die zum Füllen des Lesepuffers verwendet wird).
Daten werden in Writable-Streams gepuffert, wenn die Methode writable.write(chunk) wiederholt aufgerufen wird. Solange die Gesamtgröße des internen Schreibpuffers unter dem durch highWaterMark gesetzten Schwellenwert liegt, geben Aufrufe von writable.write() den Wert true zurück. Sobald die Größe des internen Puffers den highWaterMark erreicht oder überschreitet, wird false zurückgegeben.
Ein Hauptziel der stream-API, insbesondere der Methode stream.pipe(), ist es, die Pufferung von Daten auf ein akzeptables Maß zu begrenzen, sodass Quellen und Ziele unterschiedlicher Geschwindigkeit den verfügbaren Speicher nicht überfordern.
Die highWaterMark-Option ist ein Schwellenwert, kein Limit: Sie diktiert die Menge an Daten, die ein Stream puffert, bevor er aufhört, nach weiteren Daten zu fragen. Sie erzwingt im Allgemeinen keine strikte Speicherbeschränkung. Spezifische Stream-Implementierungen können sich dafür entscheiden, strengere Grenzwerte durchzusetzen, aber dies ist optional.
Da Duplex- und Transform-Streams sowohl Readable als auch Writable sind, unterhält jeder zwei separate interne Puffer zum Lesen und Schreiben. Dies ermöglicht es jeder Seite, unabhängig von der anderen zu arbeiten und gleichzeitig einen angemessenen und effizienten Datenfluss aufrechtzuerhalten. Zum Beispiel sind net.Socket-Instanzen Duplex-Streams, deren Readable-Seite den Verbrauch von vom Socket empfangenen Daten ermöglicht und deren Writable-Seite das Schreiben von Daten auf den Socket erlaubt. Da Daten möglicherweise schneller oder langsamer auf den Socket geschrieben als empfangen werden können, sollte jede Seite unabhängig von der anderen arbeiten (und puffern).
Die Mechanik der internen Pufferung ist ein internes Implementierungsdetail und kann jederzeit geändert werden. Für bestimmte fortgeschrittene Implementierungen können die internen Puffer jedoch mithilfe von writable.writableBuffer oder readable.readableBuffer abgerufen werden. Von der Verwendung dieser nicht dokumentierten Eigenschaften wird abgeraten.
API für Stream-Consumer#
Fast alle Node.js-Anwendungen, egal wie einfach sie sind, verwenden auf irgendeine Weise Streams. Das Folgende ist ein Beispiel für die Verwendung von Streams in einer Node.js-Anwendung, die einen HTTP-Server implementiert.
const http = require('node:http');
const server = http.createServer((req, res) => {
// `req` is an http.IncomingMessage, which is a readable stream.
// `res` is an http.ServerResponse, which is a writable stream.
let body = '';
// Get the data as utf8 strings.
// If an encoding is not set, Buffer objects will be received.
req.setEncoding('utf8');
// Readable streams emit 'data' events once a listener is added.
req.on('data', (chunk) => {
body += chunk;
});
// The 'end' event indicates that the entire body has been received.
req.on('end', () => {
try {
const data = JSON.parse(body);
// Write back something interesting to the user:
res.write(typeof data);
res.end();
} catch (er) {
// uh oh! bad json!
res.statusCode = 400;
return res.end(`error: ${er.message}`);
}
});
});
server.listen(1337);
// $ curl localhost:1337 -d "{}"
// object
// $ curl localhost:1337 -d "\"foo\""
// string
// $ curl localhost:1337 -d "not json"
// error: Unexpected token 'o', "not json" is not valid JSON
Writable-Streams (wie res im Beispiel) stellen Methoden wie write() und end() bereit, die verwendet werden, um Daten in den Stream zu schreiben.
Readable-Streams verwenden die EventEmitter-API, um Anwendungscode zu benachrichtigen, wenn Daten zum Lesen aus dem Stream verfügbar sind. Diese verfügbaren Daten können auf verschiedene Weise aus dem Stream gelesen werden.
Sowohl Writable- als auch Readable-Streams verwenden die EventEmitter-API auf verschiedene Weise, um den aktuellen Status des Streams zu kommunizieren.
Duplex- und Transform-Streams sind sowohl Writable als auch Readable.
Anwendungen, die Daten in einen Stream schreiben oder aus diesem konsumieren, sind nicht verpflichtet, die Stream-Schnittstellen direkt zu implementieren, und haben im Allgemeinen keinen Grund, require('node:stream') aufzurufen.
Entwickler, die neue Arten von Streams implementieren möchten, sollten sich auf den Abschnitt API für Stream-Implementierer beziehen.
Schreibbare Streams#
Schreibbare Streams sind eine Abstraktion für ein Ziel, in das Daten geschrieben werden.
Beispiele für Writable-Streams sind:
- HTTP-Anfragen auf der Client-Seite
- HTTP-Antworten auf der Server-Seite
- fs-Schreib-Streams
- zlib-Streams
- crypto-Streams
- TCP-Sockets
- Standardeingabe (stdin) eines Kindprozesses
process.stdout,process.stderr
Einige dieser Beispiele sind eigentlich Duplex-Streams, die die Writable-Schnittstelle implementieren.
Alle Writable-Streams implementieren die durch die Klasse stream.Writable definierte Schnittstelle.
Während sich spezifische Instanzen von Writable-Streams in verschiedener Weise unterscheiden können, folgen alle Writable-Streams demselben grundlegenden Nutzungsmuster, wie im folgenden Beispiel illustriert:
const myStream = getWritableStreamSomehow();
myStream.write('some data');
myStream.write('some more data');
myStream.end('done writing data');
Klasse: stream.Writable#
Ereignis: 'close'#
Das 'close'-Ereignis wird emittiert, wenn der Stream und alle seine zugrunde liegenden Ressourcen (z. B. ein Dateideskriptor) geschlossen wurden. Das Ereignis zeigt an, dass keine weiteren Ereignisse mehr emittiert werden und keine weiteren Berechnungen stattfinden werden.
Ein Writable-Stream emittiert immer das 'close'-Ereignis, wenn er mit der emitClose-Option erstellt wurde.
Ereignis: 'drain'#
Wenn ein Aufruf von stream.write(chunk) den Wert false zurückgibt, wird das 'drain'-Ereignis emittiert, wenn es angebracht ist, das Schreiben von Daten in den Stream fortzusetzen.
// Write the data to the supplied writable stream one million times.
// Be attentive to back-pressure.
function writeOneMillionTimes(writer, data, encoding, callback) {
let i = 1000000;
write();
function write() {
let ok = true;
do {
i--;
if (i === 0) {
// Last time!
writer.write(data, encoding, callback);
} else {
// See if we should continue, or wait.
// Don't pass the callback, because we're not done yet.
ok = writer.write(data, encoding);
}
} while (i > 0 && ok);
if (i > 0) {
// Had to stop early!
// Write some more once it drains.
writer.once('drain', write);
}
}
}
Ereignis: 'error'#
- Typ:
<Error>
Das 'error'-Ereignis wird emittiert, wenn beim Schreiben oder Pipelining von Daten ein Fehler aufgetreten ist. Der Listener-Callback erhält bei Aufruf ein einzelnes Error-Argument.
Der Stream wird geschlossen, wenn das 'error'-Ereignis emittiert wird, es sei denn, die Option autoDestroy wurde beim Erstellen des Streams auf false gesetzt.
Nach 'error' sollten keine weiteren Ereignisse außer 'close' emittiert werden (einschließlich 'error'-Ereignisse).
Ereignis: 'finish'#
Das 'finish'-Ereignis wird emittiert, nachdem die Methode stream.end() aufgerufen wurde und alle Daten an das zugrunde liegende System geleert (flushed) wurden.
const writer = getWritableStreamSomehow();
for (let i = 0; i < 100; i++) {
writer.write(`hello, #${i}!\n`);
}
writer.on('finish', () => {
console.log('All writes are now complete.');
});
writer.end('This is the end\n');
Ereignis: 'pipe'#
src<stream.Readable>Quell-Stream, der in diesen Writable-Stream gepiped wird.
Das 'pipe'-Ereignis wird emittiert, wenn die Methode stream.pipe() auf einem lesbaren Stream aufgerufen wird und diesen Writable-Stream zu seinen Zielen hinzufügt.
const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('pipe', (src) => {
console.log('Something is piping into the writer.');
assert.equal(src, reader);
});
reader.pipe(writer);
Ereignis: 'unpipe'#
src<stream.Readable>Der Quell-Stream, der diesen Writable-Stream unpiped (die Verbindung trennte).
Das 'unpipe'-Ereignis wird emittiert, wenn die Methode stream.unpipe() auf einem Readable-Stream aufgerufen wird und diesen Writable-Stream aus seinen Zielen entfernt.
Dies wird auch emittiert, falls dieser Writable-Stream einen Fehler emittiert, während ein Readable-Stream in ihn hineinpiped.
const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
console.log('Something has stopped piping into the writer.');
assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer);
writable.cork()#
Die Methode writable.cork() erzwingt, dass alle geschriebenen Daten im Speicher gepuffert werden. Die gepufferten Daten werden geleert, wenn entweder die Methode stream.uncork() oder stream.end() aufgerufen wird.
Die primäre Absicht von writable.cork() ist es, Situationen zu bewältigen, in denen mehrere kleine Datenstücke in schneller Folge in den Stream geschrieben werden. Anstatt sie sofort an das zugrunde liegende Ziel weiterzuleiten, puffert writable.cork() alle Datenstücke, bis writable.uncork() aufgerufen wird, was diese dann alle an writable._writev() übergibt, falls vorhanden. Dies verhindert eine Head-of-Line-Blocking-Situation, in der Daten gepuffert werden, während auf die Verarbeitung des ersten kleinen Datenstücks gewartet wird. Die Verwendung von writable.cork() ohne die Implementierung von writable._writev() kann jedoch nachteilige Auswirkungen auf den Durchsatz haben.
Siehe auch: writable.uncork(), writable._writev().
writable.destroy([error])#
error<Error>Optional, ein Fehler, der mit dem'error'-Ereignis emittiert werden soll.- Rückgabe:
<this>
Zerstöre den Stream. Emittiere optional ein 'error'-Ereignis und ein 'close'-Ereignis (außer emitClose ist auf false gesetzt). Nach diesem Aufruf ist der schreibbare Stream beendet, und nachfolgende Aufrufe von write() oder end() führen zu einem ERR_STREAM_DESTROYED-Fehler. Dies ist eine destruktive und sofortige Methode, um einen Stream zu zerstören. Frühere Aufrufe von write() wurden möglicherweise noch nicht geleert und können einen ERR_STREAM_DESTROYED-Fehler auslösen. Verwenden Sie end() anstelle von destroy(), wenn Daten vor dem Schließen geleert werden sollen, oder warten Sie auf das 'drain'-Ereignis, bevor Sie den Stream zerstören.
const { Writable } = require('node:stream');
const myStream = new Writable();
const fooErr = new Error('foo error');
myStream.destroy(fooErr);
myStream.on('error', (fooErr) => console.error(fooErr.message)); // foo error
const { Writable } = require('node:stream');
const myStream = new Writable();
myStream.destroy();
myStream.on('error', function wontHappen() {});
const { Writable } = require('node:stream');
const myStream = new Writable();
myStream.destroy();
myStream.write('foo', (error) => console.error(error.code));
// ERR_STREAM_DESTROYED
Sobald destroy() aufgerufen wurde, sind alle weiteren Aufrufe No-Ops, und außer von _destroy() sollten keine weiteren Fehler mehr als 'error' emittiert werden.
Implementierer sollten diese Methode nicht überschreiben, sondern stattdessen writable._destroy() implementieren.
writable.closed#
- Typ:
<boolean>
Ist true, nachdem 'close' emittiert wurde.
writable.destroyed#
- Typ:
<boolean>
Ist true, nachdem writable.destroy() aufgerufen wurde.
const { Writable } = require('node:stream');
const myStream = new Writable();
console.log(myStream.destroyed); // false
myStream.destroy();
console.log(myStream.destroyed); // true
writable.end([chunk[, encoding]][, callback])#
chunk<string>|<Buffer>|<TypedArray>|<DataView>|<any>Optionale Daten zum Schreiben. Für Streams, die nicht im Objektmodus operieren, musschunkein<string>,<Buffer>,<TypedArray>oder<DataView>sein. Für Streams im Objektmodus kannchunkjeder beliebige JavaScript-Wert außernullsein.encoding<string>Die Kodierung, wennchunkein String ist.callback<Function>Callback für den Abschluss des Streams.- Rückgabe:
<this>
Der Aufruf der Methode writable.end() signalisiert, dass keine weiteren Daten mehr in den Writable geschrieben werden. Die optionalen chunk- und encoding-Argumente erlauben es, ein letztes zusätzliches Datenstück unmittelbar vor dem Schließen des Streams zu schreiben.
Der Aufruf der Methode stream.write() nach dem Aufruf von stream.end() löst einen Fehler aus.
// Write 'hello, ' and then end with 'world!'.
const fs = require('node:fs');
const file = fs.createWriteStream('example.txt');
file.write('hello, ');
file.end('world!');
// Writing more now is not allowed!
writable.setDefaultEncoding(encoding)#
Die Methode writable.setDefaultEncoding() setzt die Standard-encoding für einen Writable-Stream.
writable.uncork()#
Die Methode writable.uncork() leert alle seit dem Aufruf von stream.cork() gepufferten Daten.
Bei der Verwendung von writable.cork() und writable.uncork() zur Verwaltung der Pufferung von Schreibvorgängen in einen Stream sollten Aufrufe von writable.uncork() mit process.nextTick() verzögert werden. Dies ermöglicht die Stapelverarbeitung (Batching) aller writable.write()-Aufrufe, die innerhalb einer bestimmten Node.js-Ereignisschleifenphase auftreten.
stream.cork();
stream.write('some ');
stream.write('data ');
process.nextTick(() => stream.uncork());
Wenn die Methode writable.cork() mehrfach auf einem Stream aufgerufen wird, muss die gleiche Anzahl an Aufrufen von writable.uncork() erfolgen, um die gepufferten Daten zu leeren.
stream.cork();
stream.write('some ');
stream.cork();
stream.write('data ');
process.nextTick(() => {
stream.uncork();
// The data will not be flushed until uncork() is called a second time.
stream.uncork();
});
Siehe auch: writable.cork().
writable.writable#
- Typ:
<boolean>
Ist true, wenn es sicher ist, writable.write() aufzurufen, was bedeutet, dass der Stream weder zerstört noch fehlerhaft oder beendet wurde.
writable.writableAborted#
- Typ:
<boolean>
Gibt zurück, ob der Stream zerstört wurde oder einen Fehler hatte, bevor 'finish' emittiert wurde.
writable.writableEnded#
- Typ:
<boolean>
Ist true, nachdem writable.end() aufgerufen wurde. Diese Eigenschaft gibt nicht an, ob die Daten bereits geleert wurden; verwenden Sie dafür stattdessen writable.writableFinished.
writable.writableCorked#
- Typ:
<integer>
Anzahl der Male, die writable.uncork() aufgerufen werden muss, um den Stream vollständig zu entkorken.
writable.errored#
- Typ:
<Error>
Gibt einen Fehler zurück, wenn der Stream mit einem Fehler zerstört wurde.
writable.writableFinished#
- Typ:
<boolean>
Wird unmittelbar vor dem Emittieren des 'finish'-Ereignisses auf true gesetzt.
writable.writableHighWaterMark#
- Typ:
<number>
Gibt den Wert von highWaterMark zurück, der bei der Erstellung dieses Writable-Streams übergeben wurde.
writable.writableLength#
- Typ:
<number>
Diese Eigenschaft enthält die Anzahl der Bytes (oder Objekte) in der Warteschlange, die bereit sind, geschrieben zu werden. Der Wert liefert Introspektionsdaten zum Status des highWaterMark.
writable.writableNeedDrain#
- Typ:
<boolean>
Ist true, wenn der Puffer des Streams voll war und der Stream ein 'drain'-Ereignis emittieren wird.
writable.writableObjectMode#
- Typ:
<boolean>
Getter für die Eigenschaft objectMode eines gegebenen Writable-Streams.
writable[Symbol.asyncDispose]()#
Ruft writable.destroy() mit einem AbortError auf und gibt ein Promise zurück, das erfüllt wird, wenn der Stream beendet ist.
writable.write(chunk[, encoding][, callback])#
chunk<string>|<Buffer>|<TypedArray>|<DataView>|<any>Optionale Daten zum Schreiben. Für Streams, die nicht im Objektmodus operieren, musschunkein<string>,<Buffer>,<TypedArray>oder<DataView>sein. Für Streams im Objektmodus kannchunkjeder beliebige JavaScript-Wert außernullsein.encoding<string>|<null>Die Kodierung, fallschunkein String ist. Standard:'utf8'callback<Function>Callback für den Zeitpunkt, an dem dieses Datenstück geleert wurde.- Rückgabe:
<boolean>false, wenn der Stream möchte, dass der aufrufende Code auf das Eintreten des Ereignisses'drain'wartet, bevor er weitere Daten schreibt; andernfallstrue.
Die Methode writable.write() schreibt Daten in den Stream und ruft den angegebenen callback auf, sobald die Daten vollständig verarbeitet wurden. Wenn ein Fehler auftritt, wird der callback mit dem Fehler als erstem Argument aufgerufen. Der callback wird asynchron und vor dem Emittieren von 'error' aufgerufen.
Der Rückgabewert ist true, wenn der interne Puffer nach der Aufnahme von chunk kleiner als der bei der Erstellung des Streams konfigurierte highWaterMark ist. Wenn false zurückgegeben wird, sollten weitere Versuche, Daten in den Stream zu schreiben, gestoppt werden, bis das 'drain'-Ereignis emittiert wird.
Während ein Stream nicht "draint", puffern Aufrufe von write() den chunk und geben false zurück. Sobald alle aktuell gepufferten Stücke geleert (vom Betriebssystem zur Auslieferung akzeptiert) wurden, wird das 'drain'-Ereignis emittiert. Sobald write() false zurückgibt, schreiben Sie keine weiteren Stücke, bis das 'drain'-Ereignis emittiert wird. Während das Aufrufen von write() auf einem Stream, der nicht draint, erlaubt ist, puffert Node.js alle geschriebenen Stücke, bis der maximale Speicherverbrauch erreicht ist, woraufhin es bedingungslos abbricht. Selbst vor dem Abbruch führt ein hoher Speicherverbrauch zu einer schlechten Garbage-Collector-Leistung und einem hohen RSS (das typischerweise nicht an das System zurückgegeben wird, selbst wenn der Speicher nicht mehr benötigt wird). Da TCP-Sockets möglicherweise niemals drainen, wenn der entfernte Partner die Daten nicht liest, kann das Schreiben in einen Socket, der nicht draint, zu einer remote ausnutzbaren Sicherheitslücke führen.
Das Schreiben von Daten, während der Stream nicht draint, ist besonders problematisch für einen Transform-Stream, da diese standardmäßig pausiert sind, bis sie gepiped werden oder ein 'data'- oder 'readable'-Ereignis-Listener hinzugefügt wird.
Wenn die zu schreibenden Daten bei Bedarf generiert oder abgerufen werden können, wird empfohlen, die Logik in einen Readable-Stream zu kapseln und stream.pipe() zu verwenden. Wenn das Aufrufen von write() jedoch bevorzugt wird, ist es möglich, Backpressure zu respektieren und Speicherprobleme durch Verwendung des 'drain'-Ereignisses zu vermeiden.
function write(data, cb) {
if (!stream.write(data)) {
stream.once('drain', cb);
} else {
process.nextTick(cb);
}
}
// Wait for cb to be called before doing any other write.
write('hello', () => {
console.log('Write completed, do more writes now.');
});
Ein Writable-Stream im Objektmodus ignoriert das encoding-Argument immer.
Lesbare Streams#
Lesbare Streams sind eine Abstraktion für eine Quelle, aus der Daten konsumiert werden.
Beispiele für Readable-Streams sind:
- HTTP-Antworten auf der Client-Seite
- HTTP-Anfragen auf der Server-Seite
- fs-Lese-Streams
- zlib-Streams
- crypto-Streams
- TCP-Sockets
- stdout und stderr von Kindprozessen
process.stdin
Alle Readable-Streams implementieren die Schnittstelle der Klasse stream.Readable.
Zwei Lesemodi#
Readable-Streams arbeiten effektiv in einem von zwei Modi: fließend (flowing) und pausiert (paused). Diese Modi sind vom Objektmodus getrennt. Ein Readable-Stream kann sich im Objektmodus befinden oder nicht, unabhängig davon, ob er sich im fließenden oder pausierten Modus befindet.
-
Im fließenden Modus werden Daten automatisch aus dem zugrunde liegenden System gelesen und einer Anwendung so schnell wie möglich über Ereignisse mittels der
EventEmitter-Schnittstelle bereitgestellt. -
Im pausierten Modus muss die Methode
stream.read()explizit aufgerufen werden, um Datenstücke aus dem Stream zu lesen.
Alle Readable-Streams beginnen im pausierten Modus, können aber auf eine der folgenden Weisen in den fließenden Modus geschaltet werden:
- Hinzufügen eines
'data'-Ereignis-Listeners. - Aufrufen der Methode
stream.resume(). - Aufrufen der Methode
stream.pipe(), um die Daten an einenWritable-Stream zu senden.
Der Readable-Stream kann auf eine der folgenden Weisen zurück in den pausierten Modus schalten:
- Wenn keine Pipe-Ziele existieren, durch Aufrufen der Methode
stream.pause(). - Wenn Pipe-Ziele existieren, durch Entfernen aller Pipe-Ziele. Mehrere Pipe-Ziele können durch Aufrufen der Methode
stream.unpipe()entfernt werden.
Das wichtige Konzept, das man sich merken muss, ist, dass ein Readable-Stream keine Daten generiert, bis ein Mechanismus zum Konsumieren oder Ignorieren dieser Daten bereitgestellt wird. Wenn der konsumierende Mechanismus deaktiviert oder entfernt wird, versucht der Readable-Stream, die Generierung der Daten zu stoppen.
Aus Gründen der Abwärtskompatibilität führt das Entfernen von 'data'-Ereignis-Listenern nicht automatisch zum Pausieren des Streams. Wenn zudem Pipe-Ziele vorhanden sind, garantiert das Aufrufen von stream.pause() nicht, dass der Stream pausiert bleibt, sobald diese Ziele drainen und nach weiteren Daten fragen.
Wenn ein Readable in den fließenden Modus geschaltet wird und keine Consumer zur Verarbeitung der Daten verfügbar sind, gehen diese Daten verloren. Dies kann beispielsweise passieren, wenn die Methode readable.resume() aufgerufen wird, ohne dass ein Listener an das 'data'-Ereignis angehängt ist, oder wenn ein 'data'-Ereignis-Listener vom Stream entfernt wird.
Das Hinzufügen eines 'readable'-Ereignis-Listeners bewirkt automatisch, dass der Stream aufhört zu fließen, und die Daten müssen dann über readable.read() konsumiert werden. Wenn der 'readable'-Ereignis-Listener entfernt wird, beginnt der Stream wieder zu fließen, sofern ein 'data'-Ereignis-Listener vorhanden ist.
Drei Zustände#
Die "zwei Modi" des Betriebs für einen Readable-Stream sind eine vereinfachte Abstraktion für die kompliziertere interne Zustandsverwaltung, die innerhalb der Readable-Stream-Implementierung stattfindet.
Genauer gesagt befindet sich jeder Readable-Stream zu jedem Zeitpunkt in einem von drei möglichen Zuständen:
readable.readableFlowing === nullreadable.readableFlowing === falsereadable.readableFlowing === true
Wenn readable.readableFlowing den Wert null hat, ist kein Mechanismus zum Konsumieren der Daten des Streams bereitgestellt. Daher generiert der Stream keine Daten. Während er sich in diesem Zustand befindet, schaltet das Anhängen eines Listeners für das 'data'-Ereignis, das Aufrufen der readable.pipe()-Methode oder das Aufrufen der readable.resume()-Methode readable.readableFlowing auf true, wodurch der Readable-Stream beginnt, aktiv Ereignisse zu emittieren, während Daten generiert werden.
Das Aufrufen von readable.pause(), readable.unpipe() oder das Empfangen von Backpressure führt dazu, dass readable.readableFlowing auf false gesetzt wird, wodurch der Fluss von Ereignissen vorübergehend angehalten, aber nicht die Generierung von Daten gestoppt wird. Während er sich in diesem Zustand befindet, schaltet das Anhängen eines Listeners für das 'data'-Ereignis readable.readableFlowing nicht auf true.
const { PassThrough, Writable } = require('node:stream');
const pass = new PassThrough();
const writable = new Writable();
pass.pipe(writable);
pass.unpipe(writable);
// readableFlowing is now false.
pass.on('data', (chunk) => { console.log(chunk.toString()); });
// readableFlowing is still false.
pass.write('ok'); // Will not emit 'data'.
pass.resume(); // Must be called to make stream emit 'data'.
// readableFlowing is now true.
Während readable.readableFlowing den Wert false hat, können sich Daten im internen Puffer des Streams ansammeln.
Wählen Sie einen API-Stil#
Die Readable-Stream-API hat sich über mehrere Node.js-Versionen hinweg weiterentwickelt und bietet mehrere Methoden zum Konsumieren von Stream-Daten. Im Allgemeinen sollten Entwickler eine Methode zum Konsumieren von Daten wählen und niemals mehrere Methoden verwenden, um Daten aus einem einzigen Stream zu konsumieren. Insbesondere die Verwendung einer Kombination aus on('data'), on('readable'), pipe() oder asynchronen Iteratoren könnte zu unintuitivem Verhalten führen.
Klasse: stream.Readable#
Ereignis: 'close'#
Das 'close'-Ereignis wird emittiert, wenn der Stream und alle seine zugrunde liegenden Ressourcen (z. B. ein Dateideskriptor) geschlossen wurden. Das Ereignis zeigt an, dass keine weiteren Ereignisse mehr emittiert werden und keine weiteren Berechnungen stattfinden werden.
Ein Readable-Stream emittiert immer das 'close'-Ereignis, wenn er mit der emitClose-Option erstellt wurde.
Ereignis: 'data'#
chunk<Buffer>|<string>|<any>Das Datenstück. Für Streams, die nicht im Objektmodus operieren, ist das Stück entweder ein String oder einBuffer. Für Streams im Objektmodus kann das Stück jeder JavaScript-Wert außernullsein.
Das 'data'-Ereignis wird jedes Mal emittiert, wenn der Stream die Kontrolle über ein Datenstück an einen Consumer abgibt. Dies kann geschehen, wenn der Stream in den fließenden Modus geschaltet wird, indem readable.pipe() oder readable.resume() aufgerufen wird oder ein Listener-Callback an das 'data'-Ereignis angehängt wird. Das 'data'-Ereignis wird auch jedes Mal emittiert, wenn die Methode readable.read() aufgerufen wird und ein Datenstück zur Rückgabe verfügbar ist.
Das Anhängen eines 'data'-Ereignis-Listeners an einen Stream, der nicht explizit pausiert wurde, schaltet den Stream in den fließenden Modus. Daten werden dann weitergegeben, sobald sie verfügbar sind.
Der Listener-Callback erhält das Datenstück als String, wenn für den Stream eine Standardkodierung mithilfe der readable.setEncoding()-Methode festgelegt wurde; andernfalls werden die Daten als Buffer übergeben.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
});
Ereignis: 'end'#
Das 'end'-Ereignis wird emittiert, wenn keine Daten mehr aus dem Stream konsumiert werden können.
Das 'end'-Ereignis wird nicht emittiert, sofern die Daten nicht vollständig konsumiert wurden. Dies kann erreicht werden, indem der Stream in den fließenden Modus geschaltet wird oder stream.read() wiederholt aufgerufen wird, bis alle Daten konsumiert wurden.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
});
readable.on('end', () => {
console.log('There will be no more data.');
});
Ereignis: 'error'#
- Typ:
<Error>
Das 'error'-Ereignis kann von einer Readable-Implementierung jederzeit emittiert werden. Typischerweise kann dies auftreten, wenn der zugrunde liegende Stream aufgrund eines internen Fehlers keine Daten generieren kann oder wenn eine Stream-Implementierung versucht, ein ungültiges Datenstück zu pushen.
Dem Listener-Callback wird ein einzelnes Error-Objekt übergeben.
Ereignis: 'pause'#
Das 'pause'-Ereignis wird emittiert, wenn stream.pause() aufgerufen wird und readableFlowing nicht false ist.
Ereignis: 'readable'#
Das 'readable'-Ereignis wird emittiert, wenn Daten aus dem Stream gelesen werden können, bis hin zum konfigurierten High Water Mark (state.highWaterMark). Effektiv zeigt es an, dass der Stream neue Informationen im Puffer hat. Wenn in diesem Puffer Daten verfügbar sind, kann stream.read() aufgerufen werden, um diese Daten abzurufen. Zusätzlich kann das 'readable'-Ereignis auch emittiert werden, wenn das Ende des Streams erreicht wurde.
const readable = getReadableStreamSomehow();
readable.on('readable', function() {
// There is some data to read now.
let data;
while ((data = this.read()) !== null) {
console.log(data);
}
});
Wenn das Ende des Streams erreicht wurde, gibt der Aufruf von stream.read() den Wert null zurück und löst das 'end'-Ereignis aus. Dies gilt auch, wenn nie Daten zum Lesen vorhanden waren. Im folgenden Beispiel ist foo.txt zum Beispiel eine leere Datei:
const fs = require('node:fs');
const rr = fs.createReadStream('foo.txt');
rr.on('readable', () => {
console.log(`readable: ${rr.read()}`);
});
rr.on('end', () => {
console.log('end');
});
Die Ausgabe beim Ausführen dieses Skripts ist:
$ node test.js
readable: null
end
In einigen Fällen führt das Anhängen eines Listeners für das 'readable'-Ereignis dazu, dass eine gewisse Menge an Daten in einen internen Puffer gelesen wird.
Im Allgemeinen sind die readable.pipe()- und 'data'-Ereignis-Mechanismen einfacher zu verstehen als das 'readable'-Ereignis. Die Handhabung von 'readable' könnte jedoch zu einem erhöhten Durchsatz führen.
Wenn sowohl 'readable' als auch 'data' gleichzeitig verwendet werden, hat 'readable' Vorrang bei der Steuerung des Flusses, d. h. 'data' wird nur emittiert, wenn stream.read() aufgerufen wird. Die Eigenschaft readableFlowing wird dann zu false. Wenn 'data'-Listener vorhanden sind, während 'readable' entfernt wird, beginnt der Stream zu fließen, d. h. 'data'-Ereignisse werden ohne Aufruf von .resume() emittiert.
Ereignis: 'resume'#
Das 'resume'-Ereignis wird emittiert, wenn stream.resume() aufgerufen wird und readableFlowing nicht true ist.
readable.destroy([error])#
Zerstöre den Stream. Emittiere optional ein 'error'-Ereignis und ein 'close'-Ereignis (außer emitClose ist auf false gesetzt). Nach diesem Aufruf gibt der lesbare Stream alle internen Ressourcen frei und nachfolgende Aufrufe von push() werden ignoriert.
Sobald destroy() aufgerufen wurde, sind alle weiteren Aufrufe No-Ops, und außer von _destroy() sollten keine weiteren Fehler mehr als 'error' emittiert werden.
Implementierer sollten diese Methode nicht überschreiben, sondern stattdessen readable._destroy() implementieren.
readable.closed#
- Typ:
<boolean>
Ist true, nachdem 'close' emittiert wurde.
readable.destroyed#
- Typ:
<boolean>
Ist true, nachdem readable.destroy() aufgerufen wurde.
readable.isPaused()#
- Rückgabewert:
<boolean>
Die Methode readable.isPaused() gibt den aktuellen Betriebszustand des Readable-Streams zurück. Dies wird hauptsächlich von dem Mechanismus verwendet, der der readable.pipe()-Methode zugrunde liegt. In den meisten typischen Fällen gibt es keinen Grund, diese Methode direkt zu verwenden.
const readable = new stream.Readable();
readable.isPaused(); // === false
readable.pause();
readable.isPaused(); // === true
readable.resume();
readable.isPaused(); // === false
readable.pause()#
- Rückgabe:
<this>
Die Methode readable.pause() bewirkt, dass ein Stream im fließenden Modus aufhört, 'data'-Ereignisse zu emittieren, und aus dem fließenden Modus schaltet. Alle Daten, die verfügbar werden, verbleiben im internen Puffer.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
readable.pause();
console.log('There will be no additional data for 1 second.');
setTimeout(() => {
console.log('Now data will start flowing again.');
readable.resume();
}, 1000);
});
Die Methode readable.pause() hat keine Wirkung, wenn ein 'readable'-Ereignis-Listener vorhanden ist.
readable.pipe(destination[, options])#
destination<stream.Writable>Das Ziel für das Schreiben von Daten.options<Object>Pipe-Optionen.end<boolean>Beende den Writer, wenn der Reader endet. Standard:true.
- Rückgabe:
<stream.Writable>Das Ziel, was eine Kette von Pipes ermöglicht, wenn es einDuplex- oderTransform-Stream ist.
Die Methode readable.pipe() hängt einen Writable-Stream an den Readable-Stream an, wodurch dieser automatisch in den fließenden Modus schaltet und alle seine Daten an den angehängten Writable-Stream weiterleitet. Der Datenfluss wird automatisch verwaltet, sodass der Ziel-Writable-Stream nicht von einem schnelleren Readable-Stream überfordert wird.
Das folgende Beispiel piped alle Daten vom readable in eine Datei namens file.txt:
const fs = require('node:fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt'.
readable.pipe(writable);
Es ist möglich, mehrere Writable-Streams an einen einzigen Readable-Stream anzuhängen.
Die Methode readable.pipe() gibt eine Referenz auf den Ziel-Stream zurück, wodurch es möglich ist, Ketten von gepipten Streams einzurichten.
const fs = require('node:fs');
const zlib = require('node:zlib');
const r = fs.createReadStream('file.txt');
const z = zlib.createGzip();
const w = fs.createWriteStream('file.txt.gz');
r.pipe(z).pipe(w);
Standardmäßig wird stream.end() auf dem Ziel-Writable-Stream aufgerufen, wenn der Quell-Readable-Stream 'end' emittiert, sodass das Ziel nicht mehr beschreibbar ist. Um dieses Standardverhalten zu deaktivieren, kann die Option end auf false gesetzt werden, wodurch der Ziel-Stream offen bleibt.
reader.pipe(writer, { end: false });
reader.on('end', () => {
writer.end('Goodbye\n');
});
Ein wichtiger Hinweis ist, dass der Writable-Ziel-Stream nicht automatisch geschlossen wird, wenn der Readable-Stream während der Verarbeitung einen Fehler emittiert. Wenn ein Fehler auftritt, ist es notwendig, jeden Stream manuell zu schließen, um Speicherlecks zu vermeiden.
Die process.stderr- und process.stdout Writable-Streams werden niemals geschlossen, bevor der Node.js-Prozess beendet wird, unabhängig von den angegebenen Optionen.
readable.read([size])#
size<number>Optionales Argument zur Angabe der zu lesenden Datenmenge.- Rückgabe:
<string>|<Buffer>|<null>|<any>
Die Methode readable.read() liest Daten aus dem internen Puffer und gibt sie zurück. Wenn keine Daten zum Lesen verfügbar sind, wird null zurückgegeben. Standardmäßig werden die Daten als Buffer-Objekt zurückgegeben, sofern keine Kodierung mithilfe der readable.setEncoding()-Methode angegeben wurde oder der Stream im Objektmodus arbeitet.
Das optionale size-Argument gibt eine spezifische Anzahl von Bytes zum Lesen an. Wenn size Bytes nicht zum Lesen verfügbar sind, wird null zurückgegeben, außer der Stream wurde bereits beendet, in welchem Fall alle im internen Puffer verbliebenen Daten zurückgegeben werden.
Wenn das size-Argument nicht angegeben ist, werden alle im internen Puffer enthaltenen Daten zurückgegeben.
Das size-Argument muss kleiner oder gleich 1 GiB sein.
Die Methode readable.read() sollte nur auf Readable-Streams aufgerufen werden, die im pausierten Modus arbeiten. Im fließenden Modus wird readable.read() automatisch aufgerufen, bis der interne Puffer vollständig geleert ist.
const readable = getReadableStreamSomehow();
// 'readable' may be triggered multiple times as data is buffered in
readable.on('readable', () => {
let chunk;
console.log('Stream is readable (new data received in buffer)');
// Use a loop to make sure we read all currently available data
while (null !== (chunk = readable.read())) {
console.log(`Read ${chunk.length} bytes of data...`);
}
});
// 'end' will be triggered once when there is no more data available
readable.on('end', () => {
console.log('Reached end of stream.');
});
Jeder Aufruf von readable.read() gibt ein Datenstück oder null zurück, was signalisiert, dass in diesem Moment keine weiteren Daten zum Lesen vorhanden sind. Diese Stücke werden nicht automatisch verkettet. Da ein einzelner read()-Aufruf nicht alle Daten zurückgibt, kann eine while-Schleife erforderlich sein, um kontinuierlich Datenstücke zu lesen, bis alle Daten abgerufen wurden. Beim Lesen einer großen Datei könnte .read() vorübergehend null zurückgeben, was darauf hinweist, dass der gesamte gepufferte Inhalt konsumiert wurde, aber möglicherweise noch weitere Daten zu puffern sind. In solchen Fällen wird ein neues 'readable'-Ereignis emittiert, sobald weitere Daten im Puffer sind, und das 'end'-Ereignis signalisiert das Ende der Datenübertragung.
Um also den gesamten Inhalt einer Datei aus einem readable-Stream zu lesen, ist es notwendig, Datenstücke über mehrere 'readable'-Ereignisse hinweg zu sammeln:
const chunks = [];
readable.on('readable', () => {
let chunk;
while (null !== (chunk = readable.read())) {
chunks.push(chunk);
}
});
readable.on('end', () => {
const content = chunks.join('');
});
Ein Readable-Stream im Objektmodus gibt bei einem Aufruf von readable.read(size) immer ein einzelnes Element zurück, unabhängig vom Wert des size-Arguments.
Wenn die Methode readable.read() ein Datenstück zurückgibt, wird auch ein 'data'-Ereignis emittiert.
Das Aufrufen von stream.read([size]) nach dem Emittieren des 'end'-Ereignisses gibt null zurück. Es wird kein Laufzeitfehler ausgelöst.
readable.readable#
- Typ:
<boolean>
Ist true, wenn es sicher ist, readable.read() aufzurufen, was bedeutet, dass der Stream nicht zerstört wurde und weder 'error' noch 'end' emittiert hat.
readable.readableAborted#
- Typ:
<boolean>
Gibt zurück, ob der Stream zerstört wurde oder einen Fehler hatte, bevor 'end' emittiert wurde.
readable.readableDidRead#
- Typ:
<boolean>
Gibt zurück, ob 'data' emittiert wurde.
readable.readableEncoding#
Getter für die Eigenschaft encoding eines gegebenen Readable-Streams. Die encoding-Eigenschaft kann mithilfe der readable.setEncoding()-Methode gesetzt werden.
readable.readableEnded#
- Typ:
<boolean>
Wird true, wenn das 'end'-Ereignis emittiert wird.
readable.errored#
- Typ:
<Error>
Gibt einen Fehler zurück, wenn der Stream mit einem Fehler zerstört wurde.
readable.readableFlowing#
- Typ:
<boolean>
Diese Eigenschaft spiegelt den aktuellen Zustand eines Readable-Streams wider, wie im Abschnitt Drei Zustände beschrieben.
readable.readableHighWaterMark#
- Typ:
<number>
Gibt den Wert von highWaterMark zurück, der bei der Erstellung dieses Readable-Streams übergeben wurde.
readable.readableLength#
- Typ:
<number>
Diese Eigenschaft enthält die Anzahl der Bytes (oder Objekte) in der Warteschlange, die bereit zum Lesen sind. Der Wert liefert Introspektionsdaten zum Status des highWaterMark.
readable.readableObjectMode#
- Typ:
<boolean>
Getter für die Eigenschaft objectMode eines gegebenen Readable-Streams.
readable.resume()#
- Rückgabe:
<this>
Die Methode readable.resume() bewirkt, dass ein explizit pausierter Readable-Stream wieder 'data'-Ereignisse emittiert und in den fließenden Modus schaltet.
Die Methode readable.resume() kann verwendet werden, um die Daten eines Streams vollständig zu konsumieren, ohne diese Daten tatsächlich zu verarbeiten.
getReadableStreamSomehow()
.resume()
.on('end', () => {
console.log('Reached the end, but did not read anything.');
});
Die Methode readable.resume() hat keine Wirkung, wenn ein 'readable'-Ereignis-Listener vorhanden ist.
readable.setEncoding(encoding)#
Die Methode readable.setEncoding() legt die Zeichenkodierung für Daten fest, die aus dem Readable-Stream gelesen werden.
Standardmäßig ist keine Kodierung zugewiesen, und Stream-Daten werden als Buffer-Objekte zurückgegeben. Das Festlegen einer Kodierung bewirkt, dass Stream-Daten als Strings der angegebenen Kodierung anstelle von Buffer-Objekten zurückgegeben werden. Zum Beispiel führt der Aufruf von readable.setEncoding('utf8') dazu, dass die Ausgabedaten als UTF-8 interpretiert und als Strings übergeben werden. Der Aufruf von readable.setEncoding('hex') bewirkt, dass die Daten im hexadezimalen String-Format kodiert werden.
Der Readable-Stream handhabt korrekt Multibyte-Zeichen, die durch den Stream übertragen werden, welche andernfalls falsch dekodiert würden, wenn sie einfach als Buffer-Objekte aus dem Stream abgerufen würden.
const readable = getReadableStreamSomehow();
readable.setEncoding('utf8');
readable.on('data', (chunk) => {
assert.equal(typeof chunk, 'string');
console.log('Got %d characters of string data:', chunk.length);
});
readable.unpipe([destination])#
destination<stream.Writable>Optionaler spezifischer Stream, dessen Pipe getrennt werden soll.- Rückgabe:
<this>
Die Methode readable.unpipe() trennt einen Writable-Stream, der zuvor mithilfe der stream.pipe()-Methode angehängt wurde.
Wenn das destination-Argument nicht angegeben ist, werden alle Pipes getrennt.
Wenn das destination-Argument angegeben ist, aber keine Pipe für diesen Stream eingerichtet ist, bewirkt die Methode nichts.
const fs = require('node:fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt',
// but only for the first second.
readable.pipe(writable);
setTimeout(() => {
console.log('Stop writing to file.txt.');
readable.unpipe(writable);
console.log('Manually close the file stream.');
writable.end();
}, 1000);
readable.unshift(chunk[, encoding])#
chunk<Buffer>|<TypedArray>|<DataView>|<string>|<null>|<any>Datenstück, das auf die Lesewarteschlange zurückgelegt (unshift) werden soll. Für Streams, die nicht im Objektmodus arbeiten, musschunkein<string>,<Buffer>,<TypedArray>,<DataView>odernullsein. Für Streams im Objektmodus kannchunkjeder JavaScript-Wert sein.encoding<string>Kodierung von String-Datenstücken. Muss eine gültigeBuffer-Kodierung sein, wie z. B.'utf8'oder'ascii'.
Das Übergeben von null als chunk signalisiert das Ende des Streams (EOF) und verhält sich identisch mit readable.push(null), wonach keine weiteren Daten mehr geschrieben werden können. Das EOF-Signal wird an das Ende des Puffers gestellt, und alle gepufferten Daten werden dennoch geleert.
Die Methode readable.unshift() schiebt ein Datenstück zurück in den internen Puffer. Dies ist nützlich in Situationen, in denen ein Stream von Code konsumiert wird, der einige der bereits optimistisch aus der Quelle gelesenen Daten "ent-konsumieren" muss, damit diese an eine andere Partei weitergegeben werden können.
Die Methode stream.unshift(chunk) kann nicht aufgerufen werden, nachdem das 'end'-Ereignis emittiert wurde; andernfalls wird ein Laufzeitfehler geworfen.
Entwickler, die stream.unshift() verwenden, sollten oft in Erwägung ziehen, stattdessen einen Transform-Stream zu verwenden. Siehe den Abschnitt API für Stream-Implementierer für weitere Informationen.
// Pull off a header delimited by \n\n.
// Use unshift() if we get too much.
// Call the callback with (error, header, stream).
const { StringDecoder } = require('node:string_decoder');
function parseHeader(stream, callback) {
stream.on('error', callback);
stream.on('readable', onReadable);
const decoder = new StringDecoder('utf8');
let header = '';
function onReadable() {
let chunk;
while (null !== (chunk = stream.read())) {
const str = decoder.write(chunk);
if (str.includes('\n\n')) {
// Found the header boundary.
const split = str.split(/\n\n/);
header += split.shift();
const remaining = split.join('\n\n');
const buf = Buffer.from(remaining, 'utf8');
stream.removeListener('error', callback);
// Remove the 'readable' listener before unshifting.
stream.removeListener('readable', onReadable);
if (buf.length)
stream.unshift(buf);
// Now the body of the message can be read from the stream.
callback(null, header, stream);
return;
}
// Still reading the header.
header += str;
}
}
}
Anders als stream.push(chunk) beendet stream.unshift(chunk) den Lesevorgang nicht durch das Zurücksetzen des internen Lesestatus des Streams. Dies kann unerwartete Ergebnisse verursachen, wenn readable.unshift() während eines Lesevorgangs aufgerufen wird (d. h. innerhalb einer stream._read()-Implementierung eines benutzerdefinierten Streams). Das Folgen des Aufrufs von readable.unshift() mit einem unmittelbaren stream.push('') setzt den Lesestatus korrekt zurück, es ist jedoch am besten, den Aufruf von readable.unshift() während eines laufenden Lesevorgangs einfach zu vermeiden.
readable.wrap(stream)#
Vor Node.js 0.10 implementierten Streams nicht die gesamte node:stream-Modul-API, wie sie derzeit definiert ist. (Siehe Kompatibilität für weitere Informationen.)
Bei der Verwendung einer älteren Node.js-Bibliothek, die 'data'-Ereignisse emittiert und eine stream.pause()-Methode hat, die nur beratend ist, kann die Methode readable.wrap() verwendet werden, um einen Readable-Stream zu erstellen, der den alten Stream als Datenquelle verwendet.
Es wird selten notwendig sein, readable.wrap() zu verwenden, aber die Methode wurde als Komfortfunktion für die Interaktion mit älteren Node.js-Anwendungen und -Bibliotheken bereitgestellt.
const { OldReader } = require('./old-api-module.js');
const { Readable } = require('node:stream');
const oreader = new OldReader();
const myReader = new Readable().wrap(oreader);
myReader.on('readable', () => {
myReader.read(); // etc.
});
readable[Symbol.asyncIterator]()#
- Rückgabe:
<AsyncIterator>, um den Stream vollständig zu konsumieren.
const fs = require('node:fs');
async function print(readable) {
readable.setEncoding('utf8');
let data = '';
for await (const chunk of readable) {
data += chunk;
}
console.log(data);
}
print(fs.createReadStream('file')).catch(console.error);
Wenn die Schleife mit break, return oder einem throw endet, wird der Stream zerstört. Mit anderen Worten: Das Iterieren über einen Stream konsumiert den Stream vollständig. Der Stream wird in Stücken gelesen, die der Größe der highWaterMark-Option entsprechen. Im Codebeispiel oben sind die Daten in einem einzigen Stück, wenn die Datei weniger als 64 KiB Daten enthält, da an fs.createReadStream() keine highWaterMark-Option übergeben wurde.
readable[Symbol.asyncDispose]()#
Ruft readable.destroy() mit einem AbortError auf und gibt ein Promise zurück, das erfüllt wird, wenn der Stream beendet ist.
readable.compose(stream[, options])#
stream<Writable>|<Duplex>|<WritableStream>|<TransformStream>|<Function>options<Object>signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Duplex>ein mit dem Streamstreamzusammengesetzter Stream.
import { Readable } from 'node:stream';
async function* splitToWords(source) {
for await (const chunk of source) {
const words = String(chunk).split(' ');
for (const word of words) {
yield word;
}
}
}
const wordsStream = Readable.from(['text passed through', 'composed stream']).compose(splitToWords);
const words = await wordsStream.toArray();
console.log(words); // prints ['text', 'passed', 'through', 'composed', 'stream']
readable.compose(s) ist äquivalent zu stream.compose(readable, s).
Diese Methode ermöglicht es auch, ein <AbortSignal> bereitzustellen, das den zusammengesetzten Stream bei Abbruch zerstört.
Weitere Informationen finden Sie unter stream.compose(...streams).
readable.iterator([options])#
options<Object>destroyOnReturn<boolean>Wenn auffalsegesetzt, wird der Stream nicht zerstört, wennreturnfür den Async-Iterator aufgerufen wird oder einefor await...of-Iteration mittelsbreak,returnoderthrowverlassen wird. Standard:true.
- Rückgabe:
<AsyncIterator>zum Konsumieren des Streams.
Der von dieser Methode erstellte Iterator gibt Benutzern die Möglichkeit, die Zerstörung des Streams zu stornieren, wenn die for await...of-Schleife durch return, break oder throw verlassen wird, oder wenn der Iterator den Stream zerstören soll, falls der Stream während der Iteration einen Fehler ausgegeben hat.
const { Readable } = require('node:stream');
async function printIterator(readable) {
for await (const chunk of readable.iterator({ destroyOnReturn: false })) {
console.log(chunk); // 1
break;
}
console.log(readable.destroyed); // false
for await (const chunk of readable.iterator({ destroyOnReturn: false })) {
console.log(chunk); // Will print 2 and then 3
}
console.log(readable.destroyed); // True, stream was totally consumed
}
async function printSymbolAsyncIterator(readable) {
for await (const chunk of readable) {
console.log(chunk); // 1
break;
}
console.log(readable.destroyed); // true
}
async function showBoth() {
await printIterator(Readable.from([1, 2, 3]));
await printSymbolAsyncIterator(Readable.from([1, 2, 3]));
}
showBoth();
readable.map(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Funktion, die auf jedes Datenpaket (Chunk) im Stream angewendet wird.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.highWaterMark<number>wie viele Elemente zwischengespeichert werden sollen, während auf die Konsumierung der zugeordneten Elemente gewartet wird. Standard:concurrency * 2 - 1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Readable>ein mit der Funktionfnzugeordneter Stream.
Diese Methode ermöglicht das Abbilden (Mapping) über den Stream. Die Funktion fn wird für jedes Datenpaket im Stream aufgerufen. Wenn die Funktion fn ein Promise zurückgibt, wird dieses Promise awaited, bevor es an den Ergebnis-Stream weitergeleitet wird.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).map((x) => x * 2)) {
console.log(chunk); // 2, 4, 6, 8
}
// With an asynchronous mapper, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).map((domain) => resolver.resolve4(domain), { concurrency: 2 });
for await (const result of dnsResults) {
console.log(result); // Logs the DNS result of resolver.resolve4.
}
readable.filter(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Funktion zum Filtern von Datenpaketen aus dem Stream.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.highWaterMark<number>wie viele Elemente zwischengespeichert werden sollen, während auf die Konsumierung der gefilterten Elemente gewartet wird. Standard:concurrency * 2 - 1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Readable>ein mit dem Prädikatfngefilterter Stream.
Diese Methode ermöglicht das Filtern des Streams. Für jedes Datenpaket im Stream wird die Funktion fn aufgerufen; wenn sie einen wahrheitsgemäßen (truthy) Wert zurückgibt, wird das Paket an den Ergebnis-Stream weitergegeben. Wenn die Funktion fn ein Promise zurückgibt, wird dieses Promise awaited.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).filter(async (domain) => {
const { address } = await resolver.resolve4(domain, { ttl: true });
return address.ttl > 60;
}, { concurrency: 2 });
for await (const result of dnsResults) {
// Logs domains with more than 60 seconds on the resolved dns record.
console.log(result);
}
readable.forEach(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Funktion, die für jedes Datenpaket des Streams aufgerufen wird.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Promise>ein Promise für den Zeitpunkt, an dem der Stream beendet ist.
Diese Methode ermöglicht das Iterieren über einen Stream. Für jedes Datenpaket im Stream wird die Funktion fn aufgerufen. Wenn die Funktion fn ein Promise zurückgibt, wird dieses Promise awaited.
Diese Methode unterscheidet sich von for await...of-Schleifen dadurch, dass sie Datenpakete optional gleichzeitig verarbeiten kann. Darüber hinaus kann eine forEach-Iteration nur gestoppt werden, indem eine signal-Option übergeben und der zugehörige AbortController abgebrochen wird, während for await...of mit break oder return gestoppt werden kann. In beiden Fällen wird der Stream zerstört.
Diese Methode unterscheidet sich vom Lauschen auf das 'data'-Ereignis dadurch, dass sie das readable-Ereignis in der zugrunde liegenden Mechanik verwendet und die Anzahl gleichzeitiger fn-Aufrufe begrenzen kann.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).map(async (domain) => {
const { address } = await resolver.resolve4(domain, { ttl: true });
return address;
}, { concurrency: 2 });
await dnsResults.forEach((result) => {
// Logs result, similar to `for await (const result of dnsResults)`
console.log(result);
});
console.log('done'); // Stream has finished
readable.toArray([options])#
Stabilität: 1 - Experimentell
options<Object>signal<AbortSignal>ermöglicht den Abbruch dertoArray-Operation, wenn das Signal abgebrochen wird.
- Rückgabe:
<Promise>ein Promise, das ein Array mit dem Inhalt des Streams enthält.
Diese Methode ermöglicht es, den Inhalt eines Streams einfach abzurufen.
Da diese Methode den gesamten Stream in den Speicher liest, macht sie die Vorteile von Streams zunichte. Sie ist für Interoperabilität und Komfort gedacht, nicht als primäre Methode zum Konsumieren von Streams.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
await Readable.from([1, 2, 3, 4]).toArray(); // [1, 2, 3, 4]
const resolver = new Resolver();
// Make dns queries concurrently using .map and collect
// the results into an array using toArray
const dnsResults = await Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).map(async (domain) => {
const { address } = await resolver.resolve4(domain, { ttl: true });
return address;
}, { concurrency: 2 }).toArray();
readable.some(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Funktion, die für jedes Datenpaket des Streams aufgerufen wird.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Promise>ein Promise, das zutrueausgewertet wird, wennfnfür mindestens eines der Datenpakete einen wahrheitsgemäßen Wert zurückgegeben hat.
Diese Methode ähnelt Array.prototype.some und ruft fn für jedes Datenpaket im Stream auf, bis der erwartete Rückgabewert true (oder ein wahrheitsgemäßer Wert) ist. Sobald der erwartete Rückgabewert eines fn-Aufrufs für ein Datenpaket wahrheitsgemäß ist, wird der Stream zerstört und das Promise mit true erfüllt. Wenn keiner der fn-Aufrufe für die Datenpakete einen wahrheitsgemäßen Wert zurückgibt, wird das Promise mit false erfüllt.
import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';
// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).some((x) => x > 2); // true
await Readable.from([1, 2, 3, 4]).some((x) => x < 0); // false
// With an asynchronous predicate, making at most 2 file checks at a time.
const anyBigFile = await Readable.from([
'file1',
'file2',
'file3',
]).some(async (fileName) => {
const stats = await stat(fileName);
return stats.size > 1024 * 1024;
}, { concurrency: 2 });
console.log(anyBigFile); // `true` if any file in the list is bigger than 1MB
console.log('done'); // Stream has finished
readable.find(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Funktion, die für jedes Datenpaket des Streams aufgerufen wird.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Promise>ein Promise, das zu dem ersten Datenpaket ausgewertet wird, für dasfneinen wahrheitsgemäßen Wert zurückgegeben hat, oderundefined, wenn kein Element gefunden wurde.
Diese Methode ähnelt Array.prototype.find und ruft fn für jedes Datenpaket im Stream auf, um ein Paket mit einem wahrheitsgemäßen fn-Wert zu finden. Sobald der erwartete Rückgabewert eines fn-Aufrufs wahrheitsgemäß ist, wird der Stream zerstört und das Promise mit dem Wert erfüllt, für den fn einen wahrheitsgemäßen Wert zurückgegeben hat. Wenn alle fn-Aufrufe für die Datenpakete einen falschen Wert zurückgeben, wird das Promise mit undefined erfüllt.
import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';
// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).find((x) => x > 2); // 3
await Readable.from([1, 2, 3, 4]).find((x) => x > 0); // 1
await Readable.from([1, 2, 3, 4]).find((x) => x > 10); // undefined
// With an asynchronous predicate, making at most 2 file checks at a time.
const foundBigFile = await Readable.from([
'file1',
'file2',
'file3',
]).find(async (fileName) => {
const stats = await stat(fileName);
return stats.size > 1024 * 1024;
}, { concurrency: 2 });
console.log(foundBigFile); // File name of large file, if any file in the list is bigger than 1MB
console.log('done'); // Stream has finished
readable.every(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Funktion, die für jedes Datenpaket des Streams aufgerufen wird.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Promise>ein Promise, das zutrueausgewertet wird, wennfnfür alle Datenpakete einen wahrheitsgemäßen Wert zurückgegeben hat.
Diese Methode ähnelt Array.prototype.every und ruft fn für jedes Datenpaket im Stream auf, um zu prüfen, ob alle erwarteten Rückgabewerte für fn wahrheitsgemäß sind. Sobald der erwartete Rückgabewert eines fn-Aufrufs falsch ist, wird der Stream zerstört und das Promise mit false erfüllt. Wenn alle fn-Aufrufe für die Datenpakete einen wahrheitsgemäßen Wert zurückgeben, wird das Promise mit true erfüllt.
import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';
// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).every((x) => x > 2); // false
await Readable.from([1, 2, 3, 4]).every((x) => x > 0); // true
// With an asynchronous predicate, making at most 2 file checks at a time.
const allBigFiles = await Readable.from([
'file1',
'file2',
'file3',
]).every(async (fileName) => {
const stats = await stat(fileName);
return stats.size > 1024 * 1024;
}, { concurrency: 2 });
// `true` if all files in the list are bigger than 1MiB
console.log(allBigFiles);
console.log('done'); // Stream has finished
readable.flatMap(fn[, options])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncGeneratorFunction>|<AsyncFunction>eine Funktion, die über jedes Datenpaket im Stream gemappt wird.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
options<Object>concurrency<number>die maximale Anzahl gleichzeitiger Aufrufe vonfnfür den Stream. Standard:1.signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Readable>ein mit der Funktionfnflach-abgebildeter Stream.
Diese Methode gibt einen neuen Stream zurück, indem der angegebene Callback auf jedes Datenpaket des Streams angewendet und das Ergebnis anschließend abgeflacht wird.
Es ist möglich, von fn einen Stream oder ein anderes iterierbares oder asynchron iterierbares Objekt zurückzugeben; die Ergebnis-Streams werden dann in den zurückgegebenen Stream zusammengeführt (abgeflacht).
import { Readable } from 'node:stream';
import { createReadStream } from 'node:fs';
// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).flatMap((x) => [x, x])) {
console.log(chunk); // 1, 1, 2, 2, 3, 3, 4, 4
}
// With an asynchronous mapper, combine the contents of 4 files
const concatResult = Readable.from([
'./1.mjs',
'./2.mjs',
'./3.mjs',
'./4.mjs',
]).flatMap((fileName) => createReadStream(fileName));
for await (const result of concatResult) {
// This will contain the contents (all chunks) of all 4 files
console.log(result);
}
readable.drop(limit[, options])#
Stabilität: 1 - Experimentell
limit<number>die Anzahl der aus dem Readable zu löschenden Datenpakete.options<Object>signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Readable>ein Stream, bei demlimitDatenpakete gelöscht wurden.
Diese Methode gibt einen neuen Stream zurück, bei dem die ersten limit Datenpakete gelöscht wurden.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).drop(2).toArray(); // [3, 4]
readable.take(limit[, options])#
Stabilität: 1 - Experimentell
limit<number>die Anzahl der aus dem Readable zu entnehmenden Datenpakete.options<Object>signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Readable>ein Stream, bei demlimitDatenpakete entnommen wurden.
Diese Methode gibt einen neuen Stream mit den ersten limit Datenpaketen zurück.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).take(2).toArray(); // [1, 2]
readable.reduce(fn[, initial[, options]])#
Stabilität: 1 - Experimentell
fn<Function>|<AsyncFunction>eine Reduzierfunktion, die für jedes Datenpaket im Stream aufgerufen wird.previous<any>der Wert, der aus dem letzten Aufruf vonfnoder deminitial-Wert (falls angegeben) oder sonst dem ersten Datenpaket des Streams erhalten wurde.data<any>ein Datenpaket aus dem Stream.options<Object>signal<AbortSignal>wird abgebrochen, wenn der Stream zerstört wird, wodurch der Aufruf vonfnvorzeitig abgebrochen werden kann.
initial<any>der anfängliche Wert für die Reduktion.options<Object>signal<AbortSignal>ermöglicht die Zerstörung des Streams, wenn das Signal abgebrochen wird.
- Rückgabe:
<Promise>ein Promise für den Endwert der Reduktion.
Diese Methode ruft fn nacheinander für jedes Datenpaket des Streams auf und übergibt das Ergebnis der Berechnung des vorherigen Elements. Sie gibt ein Promise für den Endwert der Reduktion zurück.
Wenn kein initial-Wert angegeben wird, wird das erste Datenpaket des Streams als Anfangswert verwendet. Wenn der Stream leer ist, wird das Promise mit einem TypeError mit dem Fehlercode ERR_INVALID_ARGS abgelehnt.
import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
const directoryPath = './src';
const filesInDir = await readdir(directoryPath);
const folderSize = await Readable.from(filesInDir)
.reduce(async (totalSize, file) => {
const { size } = await stat(join(directoryPath, file));
return totalSize + size;
}, 0);
console.log(folderSize);
Die Reduzierfunktion iteriert den Stream Element für Element, was bedeutet, dass es keinen concurrency-Parameter oder Parallelismus gibt. Um eine reduce-Operation gleichzeitig durchzuführen, können Sie die asynchrone Funktion in die Methode readable.map auslagern.
import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
const directoryPath = './src';
const filesInDir = await readdir(directoryPath);
const folderSize = await Readable.from(filesInDir)
.map((file) => stat(join(directoryPath, file)), { concurrency: 2 })
.reduce((totalSize, { size }) => totalSize + size, 0);
console.log(folderSize);
Duplex- und Transform-Streams#
Klasse: stream.Duplex#
Duplex-Streams sind Streams, die sowohl die Schnittstellen Readable als auch Writable implementieren.
Beispiele für Duplex-Streams sind
duplex.allowHalfOpen#
- Typ:
<boolean>
Wenn false, beendet der Stream die Writable-Seite automatisch, wenn die Readable-Seite endet. Anfangs durch die allowHalfOpen-Konstruktoroption gesetzt, die standardmäßig true ist.
Dies kann manuell geändert werden, um das Half-Open-Verhalten einer bestehenden Duplex-Stream-Instanz zu ändern, muss jedoch vor dem Emittieren des 'end'-Ereignisses geschehen.
Klasse: stream.Transform#
Transform-Streams sind Duplex-Streams, bei denen die Ausgabe in irgendeiner Weise mit der Eingabe zusammenhängt. Wie alle Duplex-Streams implementieren Transform-Streams sowohl die Schnittstellen Readable als auch Writable.
Beispiele für Transform-Streams sind
transform.destroy([error])#
Zerstört den Stream und emittiert optional ein 'error'-Ereignis. Nach diesem Aufruf gibt der Transform-Stream alle internen Ressourcen frei. Implementierer sollten diese Methode nicht überschreiben, sondern stattdessen readable._destroy() implementieren. Die Standardimplementierung von _destroy() für Transform emittiert ebenfalls 'close', es sei denn, emitClose ist auf false gesetzt.
Sobald destroy() aufgerufen wurde, sind weitere Aufrufe ein No-Op, und außer von _destroy() können keine weiteren Fehler als 'error' emittiert werden.
stream.duplexPair([options])#
options<Object>Ein Wert, der an beideDuplex-Konstruktoren übergeben wird, um Optionen wie Pufferung festzulegen.- Rückgabe:
<Array>von zweiDuplex-Instanzen.
Die Hilfsfunktion duplexPair gibt ein Array mit zwei Elementen zurück, die jeweils ein mit der anderen Seite verbundener Duplex-Stream sind.
const [ sideA, sideB ] = duplexPair();
Was immer in einen Stream geschrieben wird, ist auf dem anderen lesbar. Dies bietet ein Verhalten, das analog zu einer Netzwerkverbindung ist, bei der die vom Client geschriebenen Daten vom Server gelesen werden können und umgekehrt.
Die Duplex-Streams sind symmetrisch; einer oder der andere kann ohne Unterschiede im Verhalten verwendet werden.
stream.finished(stream[, options], callback)#
stream<Stream>|<ReadableStream>|<WritableStream>Ein lesbarer und/oder schreibbarer Stream/Webstream.options<Object>error<boolean>Wenn auffalsegesetzt, wird ein Aufruf vonemit('error', err)nicht als beendet behandelt. Standard:true.readable<boolean>Wenn auffalsegesetzt, wird der Callback aufgerufen, wenn der Stream endet, auch wenn der Stream noch lesbar sein könnte. Standard:true.writable<boolean>Wenn auffalsegesetzt, wird der Callback aufgerufen, wenn der Stream endet, auch wenn der Stream noch beschreibbar sein könnte. Standard:true.signal<AbortSignal>ermöglicht das Abbrechen des Wartens auf das Stream-Ende. Der zugrunde liegende Stream wird nicht abgebrochen, wenn das Signal abgebrochen wird. Der Callback wird mit einemAbortErroraufgerufen. Alle durch diese Funktion hinzugefügten registrierten Listener werden ebenfalls entfernt.
callback<Function>Eine Callback-Funktion, die ein optionales Fehlerargument entgegennimmt.- Rückgabe:
<Function>Eine Bereinigungsfunktion, die alle registrierten Listener entfernt.
Eine Funktion, um benachrichtigt zu werden, wenn ein Stream nicht mehr lesbar oder beschreibbar ist oder einen Fehler oder ein vorzeitiges Close-Ereignis erfahren hat.
const { finished } = require('node:stream');
const fs = require('node:fs');
const rs = fs.createReadStream('archive.tar');
finished(rs, (err) => {
if (err) {
console.error('Stream failed.', err);
} else {
console.log('Stream is done reading.');
}
});
rs.resume(); // Drain the stream.
Besonders nützlich bei der Fehlerbehandlung, wenn ein Stream vorzeitig zerstört wurde (wie bei einer abgebrochenen HTTP-Anfrage) und kein 'end' oder 'finish' emittiert.
Die finished-API bietet eine Promise-Version.
stream.finished() hinterlässt hängende Ereignis-Listener (insbesondere 'error', 'end', 'finish' und 'close'), nachdem der callback aufgerufen wurde. Der Grund hierfür ist, dass unerwartete 'error'-Ereignisse (aufgrund fehlerhafter Stream-Implementierungen) keine unerwarteten Abstürze verursachen. Wenn dies unerwünscht ist, muss die zurückgegebene Bereinigungsfunktion im Callback aufgerufen werden.
const cleanup = finished(rs, (err) => {
cleanup();
// ...
});
stream.pipeline(source[, ...transforms], destination, callback)#
stream.pipeline(streams, callback)#
streams<Stream[]>|<Iterable[]>|<AsyncIterable[]>|<Function[]>|<ReadableStream[]>|<WritableStream[]>|<TransformStream[]>source<Stream>|<Iterable>|<AsyncIterable>|<Function>|<ReadableStream>- Rückgabe:
<Iterable>|<AsyncIterable>
- Rückgabe:
...transforms<Stream>|<Function>|<TransformStream>source<AsyncIterable>- Rückgabe:
<AsyncIterable>
destination<Stream>|<Function>|<WritableStream>source<AsyncIterable>- Rückgabe:
<AsyncIterable>|<Promise>
callback<Function>Wird aufgerufen, wenn die Pipeline vollständig abgeschlossen ist.err<Error>valAufgelöster Wert des durchdestinationzurückgegebenenPromise.
- Rückgabe:
<Stream>
Eine Modulmethode, um zwischen Streams und Generatoren zu pipen, Fehler weiterzuleiten, ordnungsgemäß aufzuräumen und einen Callback bereitzustellen, wenn die Pipeline vollständig ist.
const { pipeline } = require('node:stream');
const fs = require('node:fs');
const zlib = require('node:zlib');
// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.
// A pipeline to gzip a potentially huge tar file efficiently:
pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz'),
(err) => {
if (err) {
console.error('Pipeline failed.', err);
} else {
console.log('Pipeline succeeded.');
}
},
);
Die pipeline-API bietet eine Promise-Version.
stream.pipeline() ruft stream.destroy(err) bei allen Streams auf, außer bei
Readable-Streams, die'end'oder'close'emittiert haben.Writable-Streams, die'finish'oder'close'emittiert haben.
stream.pipeline() hinterlässt hängende Ereignis-Listener auf den Streams, nachdem der callback aufgerufen wurde. Bei Wiederverwendung von Streams nach einem Fehler kann dies zu Speicherlecks bei Ereignis-Listenern und verschluckten Fehlern führen. Wenn der letzte Stream lesbar ist, werden hängende Ereignis-Listener entfernt, damit der letzte Stream später konsumiert werden kann.
stream.pipeline() schließt alle Streams, wenn ein Fehler auftritt. Die Verwendung von IncomingRequest mit pipeline könnte zu unerwartetem Verhalten führen, da es den Socket zerstören würde, ohne die erwartete Antwort zu senden. Siehe das Beispiel unten.
const fs = require('node:fs');
const http = require('node:http');
const { pipeline } = require('node:stream');
const server = http.createServer((req, res) => {
const fileStream = fs.createReadStream('./fileNotExist.txt');
pipeline(fileStream, res, (err) => {
if (err) {
console.log(err); // No such file
// this message can't be sent once `pipeline` already destroyed the socket
return res.end('error!!!');
}
});
});
stream.compose(...streams)#
Stabilität: 1 - stream.compose ist experimentell.
streams<Stream[]>|<Iterable[]>|<AsyncIterable[]>|<Function[]>|<ReadableStream[]>|<WritableStream[]>|<TransformStream[]>|<Duplex[]>|<Function>- Rückgabe:
<stream.Duplex>
Kombiniert zwei oder mehr Streams zu einem Duplex-Stream, der in den ersten Stream schreibt und aus dem letzten liest. Jeder bereitgestellte Stream wird mithilfe von stream.pipeline in den nächsten gepipt. Wenn einer der Streams einen Fehler erzeugt, werden alle zerstört, einschließlich des äußeren Duplex-Streams.
Da stream.compose einen neuen Stream zurückgibt, der wiederum in andere Streams gepipt werden kann (und sollte), ermöglicht es Komposition. Im Gegensatz dazu ist bei der Übergabe von Streams an stream.pipeline normalerweise der erste Stream ein lesbarer Stream und der letzte ein beschreibbarer Stream, wodurch ein geschlossener Kreislauf entsteht.
Wenn eine Function übergeben wird, muss es sich um eine Factory-Methode handeln, die ein source-Iterable entgegennimmt.
import { compose, Transform } from 'node:stream';
const removeSpaces = new Transform({
transform(chunk, encoding, callback) {
callback(null, String(chunk).replace(' ', ''));
},
});
async function* toUpper(source) {
for await (const chunk of source) {
yield String(chunk).toUpperCase();
}
}
let res = '';
for await (const buf of compose(removeSpaces, toUpper).end('hello world')) {
res += buf;
}
console.log(res); // prints 'HELLOWORLD'
stream.compose kann verwendet werden, um asynchrone Iterables, Generatoren und Funktionen in Streams umzuwandeln.
AsyncIterablekonvertiert in ein lesbaresDuplex. Kann keinnullliefern.AsyncGeneratorFunctionkonvertiert in ein lesbares/beschreibbares Transform-Duplex. Muss ein Quell-AsyncIterableals ersten Parameter entgegennehmen. Kann keinnullliefern.AsyncFunctionkonvertiert in ein beschreibbaresDuplex. Muss entwedernulloderundefinedzurückgeben.
import { compose } from 'node:stream';
import { finished } from 'node:stream/promises';
// Convert AsyncIterable into readable Duplex.
const s1 = compose(async function*() {
yield 'Hello';
yield 'World';
}());
// Convert AsyncGenerator into transform Duplex.
const s2 = compose(async function*(source) {
for await (const chunk of source) {
yield String(chunk).toUpperCase();
}
});
let res = '';
// Convert AsyncFunction into writable Duplex.
const s3 = compose(async function(source) {
for await (const chunk of source) {
res += chunk;
}
});
await finished(compose(s1, s2, s3));
console.log(res); // prints 'HELLOWORLD'
Der Einfachheit halber ist die Methode readable.compose(stream) für <Readable> und <Duplex>-Streams als Wrapper für diese Funktion verfügbar.
stream.isErrored(stream)#
stream<Readable>|<Writable>|<Duplex>|<WritableStream>|<ReadableStream>- Rückgabewert:
<boolean>
Gibt zurück, ob der Stream auf einen Fehler gestoßen ist.
stream.isReadable(stream)#
stream<Readable>|<Duplex>|<ReadableStream>- Rückgabe:
<boolean>|<null>- Gibt nurnullzurück, wennstreamkein gültigerReadable,DuplexoderReadableStreamist.
Gibt zurück, ob der Stream lesbar ist.
stream.isWritable(stream)#
stream<Writable>|<Duplex>|<WritableStream>- Rückgabe:
<boolean>|<null>- Gibt nurnullzurück, wennstreamkein gültigerWritable,DuplexoderWritableStreamist.
Gibt zurück, ob der Stream beschreibbar ist.
stream.Readable.from(iterable[, options])#
iterable<Iterable>Objekt, das dasSymbol.asyncIteratoroderSymbol.iterator-Iterable-Protokoll implementiert. Emittiert ein 'error'-Ereignis, wenn ein null-Wert übergeben wird.options<Object>Optionen, die annew stream.Readable([options])bereitgestellt werden. Standardmäßig setztReadable.from()options.objectModeauftrue, es sei denn, dies wird explizit durch Setzen vonoptions.objectModeauffalseabgewählt.- Rückgabe:
<stream.Readable>
Eine Hilfsmethode zum Erstellen von lesbaren Streams aus Iteratoren.
const { Readable } = require('node:stream');
async function * generate() {
yield 'hello';
yield 'streams';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
});
Der Aufruf von Readable.from(string) oder Readable.from(buffer) führt aus Leistungsgründen nicht dazu, dass die Strings oder Puffer iteriert werden, um den Semantiken anderer Streams zu entsprechen.
Wenn ein Iterable-Objekt, das Promises enthält, als Argument übergeben wird, kann dies zu einer nicht behandelten Ablehnung (unhandled rejection) führen.
const { Readable } = require('node:stream');
Readable.from([
new Promise((resolve) => setTimeout(resolve('1'), 1500)),
new Promise((_, reject) => setTimeout(reject(new Error('2')), 1000)), // Unhandled rejection
]);
stream.Readable.fromWeb(readableStream[, options])#
readableStream<ReadableStream>options<Object>encoding<string>highWaterMark<number>objectMode<boolean>signal<AbortSignal>
- Rückgabe:
<stream.Readable>
stream.Readable.isDisturbed(stream)#
stream<stream.Readable>|<ReadableStream>- Rückgabe:
boolean
Gibt zurück, ob der Stream gelesen oder abgebrochen wurde.
stream.Readable.toWeb(streamReadable[, options])#
streamReadable<stream.Readable>options<Object>strategy<Object>highWaterMark<number>Die maximale interne Warteschlangengröße (des erstelltenReadableStream), bevor Gegendruck (Backpressure) beim Lesen aus dem angegebenenstream.Readableangewendet wird. Wenn kein Wert angegeben wird, wird er vom angegebenenstream.Readableübernommen.size<Function>Eine Funktion, die die Größe des angegebenen Datenpakets bestimmt. Wenn kein Wert bereitgestellt wird, ist die Größe für alle Pakete1.
type<string>Gibt den Typ des erstelltenReadableStreaman. Muss'bytes'oder undefiniert sein.
- Gibt zurück:
<ReadableStream>
stream.Writable.fromWeb(writableStream[, options])#
writableStream<WritableStream>options<Object>decodeStrings<boolean>highWaterMark<number>objectMode<boolean>signal<AbortSignal>
- Rückgabe:
<stream.Writable>
stream.Writable.toWeb(streamWritable)#
streamWritable<stream.Writable>- Rückgabe:
<WritableStream>
stream.Duplex.from(src)#
src<Stream>|<Blob>|<ArrayBuffer>|<string>|<Iterable>|<AsyncIterable>|<AsyncGeneratorFunction>|<AsyncFunction>|<Promise>|<Object>|<ReadableStream>|<WritableStream>
Eine Hilfsmethode zum Erstellen von Duplex-Streams.
Streamkonvertiert einen beschreibbaren Stream in ein beschreibbaresDuplexund einen lesbaren Stream in einDuplex.Blobkonvertiert in ein lesbaresDuplex.stringkonvertiert in ein lesbaresDuplex.ArrayBufferkonvertiert in ein lesbaresDuplex.AsyncIterablekonvertiert in ein lesbaresDuplex. Kann keinnullliefern.AsyncGeneratorFunctionkonvertiert in ein lesbares/beschreibbares Transform-Duplex. Muss ein Quell-AsyncIterableals ersten Parameter entgegennehmen. Kann keinnullliefern.AsyncFunctionkonvertiert in ein beschreibbaresDuplex. Muss entwedernulloderundefinedzurückgeben.Object ({ writable, readable })konvertiertreadableundwritableinStreamund kombiniert sie dann zu einemDuplex, wobei dasDuplexin daswritableschreibt und aus demreadableliest.Promisekonvertiert in ein lesbaresDuplex. Der Wertnullwird ignoriert.ReadableStreamkonvertiert in ein lesbaresDuplex.WritableStreamkonvertiert in ein beschreibbaresDuplex.- Rückgabe:
<stream.Duplex>
Wenn ein Iterable-Objekt, das Promises enthält, als Argument übergeben wird, kann dies zu einer nicht behandelten Ablehnung (unhandled rejection) führen.
const { Duplex } = require('node:stream');
Duplex.from([
new Promise((resolve) => setTimeout(resolve('1'), 1500)),
new Promise((_, reject) => setTimeout(reject(new Error('2')), 1000)), // Unhandled rejection
]);
stream.Duplex.fromWeb(pair[, options])#
pair<Object>readable<ReadableStream>writable<WritableStream>
options<Object>- Rückgabe:
<stream.Duplex>
import { Duplex } from 'node:stream'; import { ReadableStream, WritableStream, } from 'node:stream/web'; const readable = new ReadableStream({ start(controller) { controller.enqueue('world'); }, }); const writable = new WritableStream({ write(chunk) { console.log('writable', chunk); }, }); const pair = { readable, writable, }; const duplex = Duplex.fromWeb(pair, { encoding: 'utf8', objectMode: true }); duplex.write('hello'); for await (const chunk of duplex) { console.log('readable', chunk); }const { Duplex } = require('node:stream'); const { ReadableStream, WritableStream, } = require('node:stream/web'); const readable = new ReadableStream({ start(controller) { controller.enqueue('world'); }, }); const writable = new WritableStream({ write(chunk) { console.log('writable', chunk); }, }); const pair = { readable, writable, }; const duplex = Duplex.fromWeb(pair, { encoding: 'utf8', objectMode: true }); duplex.write('hello'); duplex.once('readable', () => console.log('readable', duplex.read()));
stream.Duplex.toWeb(streamDuplex[, options])#
streamDuplex<stream.Duplex>options<Object>readableType<string>Gibt den Typ derReadableStream-Hälfte des erstellten lesbar-beschreibbaren Paares an. Muss'bytes'oder undefiniert sein. (options.typeist ein veralteter Alias für diese Option.)
- Rückgabewert:
<Object>readable<ReadableStream>writable<WritableStream>
import { Duplex } from 'node:stream'; const duplex = Duplex({ objectMode: true, read() { this.push('world'); this.push(null); }, write(chunk, encoding, callback) { console.log('writable', chunk); callback(); }, }); const { readable, writable } = Duplex.toWeb(duplex); writable.getWriter().write('hello'); const { value } = await readable.getReader().read(); console.log('readable', value);const { Duplex } = require('node:stream'); const duplex = Duplex({ objectMode: true, read() { this.push('world'); this.push(null); }, write(chunk, encoding, callback) { console.log('writable', chunk); callback(); }, }); const { readable, writable } = Duplex.toWeb(duplex); writable.getWriter().write('hello'); readable.getReader().read().then((result) => { console.log('readable', result.value); });
stream.addAbortSignal(signal, stream)#
signal<AbortSignal>Ein Signal, das eine mögliche Stornierung darstellt.stream<Stream>|<ReadableStream>|<WritableStream>Ein Stream, an den ein Signal angehängt werden soll.
Hängt ein AbortSignal an einen lesbaren oder beschreibbaren Stream an. Dies ermöglicht es dem Code, die Stream-Zerstörung mithilfe eines AbortController zu steuern.
Der Aufruf von abort auf dem AbortController, der dem übergebenen AbortSignal entspricht, verhält sich genauso wie der Aufruf von .destroy(new AbortError()) auf dem Stream und controller.error(new AbortError()) für Webstreams.
const fs = require('node:fs');
const controller = new AbortController();
const read = addAbortSignal(
controller.signal,
fs.createReadStream(('object.json')),
);
// Later, abort the operation closing the stream
controller.abort();
Oder die Verwendung eines AbortSignal mit einem lesbaren Stream als asynchron iterierbares Objekt.
const controller = new AbortController();
setTimeout(() => controller.abort(), 10_000); // set a timeout
const stream = addAbortSignal(
controller.signal,
fs.createReadStream(('object.json')),
);
(async () => {
try {
for await (const chunk of stream) {
await process(chunk);
}
} catch (e) {
if (e.name === 'AbortError') {
// The operation was cancelled
} else {
throw e;
}
}
})();
Oder die Verwendung eines AbortSignal mit einem ReadableStream.
const controller = new AbortController();
const rs = new ReadableStream({
start(controller) {
controller.enqueue('hello');
controller.enqueue('world');
controller.close();
},
});
addAbortSignal(controller.signal, rs);
finished(rs, (err) => {
if (err) {
if (err.name === 'AbortError') {
// The operation was cancelled
}
}
});
const reader = rs.getReader();
reader.read().then(({ value, done }) => {
console.log(value); // hello
console.log(done); // false
controller.abort();
});
stream.getDefaultHighWaterMark(objectMode)#
Gibt den Standard-highWaterMark zurück, der von Streams verwendet wird. Standardwert ist 65536 (64 KiB) oder 16 für objectMode.
stream.setDefaultHighWaterMark(objectMode, value)#
Setzt den Standard-highWaterMark, der von Streams verwendet wird.
API für Stream-Implementierer#
Die node:stream-Modul-API wurde entwickelt, um es zu ermöglichen, Streams einfach mithilfe des prototypischen Vererbungsmodells von JavaScript zu implementieren.
Zuerst würde ein Stream-Entwickler eine neue JavaScript-Klasse deklarieren, die eine der vier grundlegenden Stream-Klassen erweitert (stream.Writable, stream.Readable, stream.Duplex oder stream.Transform), und sicherstellen, dass sie den entsprechenden Konstruktor der Elternklasse aufruft.
const { Writable } = require('node:stream');
class MyWritable extends Writable {
constructor({ highWaterMark, ...options }) {
super({ highWaterMark });
// ...
}
}
Berücksichtigen Sie beim Erweitern von Streams, welche Optionen der Benutzer bereitstellen kann und sollte, bevor Sie diese an den Basis-Konstruktor weiterleiten. Wenn die Implementierung beispielsweise Annahmen bezüglich der autoDestroy- und emitClose-Optionen trifft, erlauben Sie dem Benutzer nicht, diese zu überschreiben. Seien Sie explizit darüber, welche Optionen weitergeleitet werden, anstatt alle Optionen implizit weiterzuleiten.
Die neue Stream-Klasse muss dann eine oder mehrere spezifische Methoden implementieren, abhängig von der Art des zu erstellenden Streams, wie in der folgenden Tabelle detailliert beschrieben.
| Anwendungsfall | Klasse | Zu implementierende Methode(n) |
|---|---|---|
| Nur lesen | Readable |
_read() |
| Nur schreiben | Writable |
_write(), _writev(), _final() |
| Lesen und Schreiben | Duplex |
_read(), _write(), _writev(), _final() |
| Verarbeitung geschriebener Daten, dann Lesen des Ergebnisses | Transform |
_transform(), _flush(), _final() |
Der Implementierungscode für einen Stream sollte niemals die "öffentlichen" Methoden eines Streams aufrufen, die für den Gebrauch durch Konsumenten vorgesehen sind (wie im Abschnitt API für Stream-Konsumenten beschrieben). Dies kann zu unerwünschten Nebeneffekten im Anwendungscode führen, der den Stream konsumiert.
Vermeiden Sie das Überschreiben öffentlicher Methoden wie write(), end(), cork(), uncork(), read() und destroy() oder das Emittieren interner Ereignisse wie 'error', 'data', 'end', 'finish' und 'close' über .emit(). Dies kann aktuelle und zukünftige Stream-Invarianten brechen, was zu Verhaltens- und/oder Kompatibilitätsproblemen mit anderen Streams, Stream-Dienstprogrammen und Benutzererwartungen führen kann.
Vereinfachte Konstruktion#
Für viele einfache Fälle ist es möglich, einen Stream zu erstellen, ohne auf Vererbung angewiesen zu sein. Dies kann erreicht werden, indem Instanzen der Objekte stream.Writable, stream.Readable, stream.Duplex oder stream.Transform direkt erstellt und entsprechende Methoden als Konstruktor-Optionen übergeben werden.
const { Writable } = require('node:stream');
const myWritable = new Writable({
construct(callback) {
// Initialize state and load resources...
},
write(chunk, encoding, callback) {
// ...
},
destroy() {
// Free resources...
},
});
Implementierung eines beschreibbaren Streams#
Die Klasse stream.Writable wird erweitert, um einen Writable-Stream zu implementieren.
Benutzerdefinierte Writable-Streams müssen den new stream.Writable([options])-Konstruktor aufrufen und die Methode writable._write() und/oder writable._writev() implementieren.
new stream.Writable([options])#
options<Object>highWaterMark<number>Puffer-Level, ab demstream.write()beginnt,falsezurückzugeben. Standard:65536(64 KiB) oder16fürobjectMode-Streams.decodeStrings<boolean>Ob anstream.write()übergebenestrings (mit der im Aufruf angegebenen Kodierung) inBufferumgewandelt werden sollen, bevor sie anstream._write()übergeben werden. Andere Datentypen werden nicht konvertiert (d. h.Buffers werden nicht instrings dekodiert). Das Setzen auffalseverhindert die Konvertierung vonstrings. Standard:true.defaultEncoding<string>Die Standardkodierung, die verwendet wird, wenn kein Argument für die Kodierung anstream.write()übergeben wird. Standard:'utf8'.objectMode<boolean>Obstream.write(anyObj)eine gültige Operation ist. Wenn gesetzt, ist es möglich, andere JavaScript-Werte als Strings,<Buffer>,<TypedArray>oder<DataView>zu schreiben, sofern von der Stream-Implementierung unterstützt. Standard:false.emitClose<boolean>Ob der Stream'close'emittieren soll, nachdem er zerstört wurde. Standard:true.write<Function>Implementierung für die Methodestream._write().writev<Function>Implementierung für die Methodestream._writev().destroy<Function>Implementierung für die Methodestream._destroy().final<Function>Implementierung für die Methodestream._final().construct<Function>Implementierung für die Methodestream._construct().autoDestroy<boolean>Ob dieser Stream nach dem Beenden automatisch.destroy()auf sich selbst aufrufen soll. Standard:true.signal<AbortSignal>Ein Signal, das eine mögliche Stornierung darstellt.
const { Writable } = require('node:stream');
class MyWritable extends Writable {
constructor(options) {
// Calls the stream.Writable() constructor.
super(options);
// ...
}
}
Oder bei Verwendung von Konstruktoren im Vor-ES6-Stil
const { Writable } = require('node:stream');
const util = require('node:util');
function MyWritable(options) {
if (!(this instanceof MyWritable))
return new MyWritable(options);
Writable.call(this, options);
}
util.inherits(MyWritable, Writable);
Oder unter Verwendung des vereinfachten Konstruktoransatzes
const { Writable } = require('node:stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
},
});
Der Aufruf von abort auf dem AbortController, der dem übergebenen AbortSignal entspricht, verhält sich genauso wie der Aufruf von .destroy(new AbortError()) auf dem beschreibbaren Stream.
const { Writable } = require('node:stream');
const controller = new AbortController();
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
},
signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort();
writable._construct(callback)#
callback<Function>Rufen Sie diese Funktion auf (optional mit einem Fehlerargument), wenn der Stream die Initialisierung abgeschlossen hat.
Die Methode _construct() DARF NICHT direkt aufgerufen werden. Sie kann von untergeordneten Klassen implementiert werden und wird in diesem Fall nur von internen Writable-Klassenmethoden aufgerufen.
Diese optionale Funktion wird einen Tick nach dem Zurückkehren des Stream-Konstruktors aufgerufen, wodurch alle _write()-, _final()- und _destroy()-Aufrufe verzögert werden, bis callback aufgerufen wird. Dies ist nützlich, um den Zustand zu initialisieren oder Ressourcen asynchron zu initialisieren, bevor der Stream verwendet werden kann.
const { Writable } = require('node:stream');
const fs = require('node:fs');
class WriteStream extends Writable {
constructor(filename) {
super();
this.filename = filename;
this.fd = null;
}
_construct(callback) {
fs.open(this.filename, 'w', (err, fd) => {
if (err) {
callback(err);
} else {
this.fd = fd;
callback();
}
});
}
_write(chunk, encoding, callback) {
fs.write(this.fd, chunk, callback);
}
_destroy(err, callback) {
if (this.fd) {
fs.close(this.fd, (er) => callback(er || err));
} else {
callback(err);
}
}
}
writable._write(chunk, encoding, callback)#
chunk<Buffer>|<string>|<any>Der zu schreibendeBuffer, konvertiert aus dem anstream.write()übergebenenstring. Wenn diedecodeStrings-Option des Streamsfalseist oder der Stream im Objektmodus arbeitet, wird das Paket nicht konvertiert und ist das, was anstream.write()übergeben wurde.encoding<string>Wenn das Datenpaket ein String ist, dann istencodingdie Zeichenkodierung dieses Strings. Wenn das Paket einBufferist oder der Stream im Objektmodus arbeitet, kannencodingignoriert werden.callback<Function>Rufen Sie diese Funktion auf (optional mit einem Fehlerargument), wenn die Verarbeitung für das bereitgestellte Datenpaket abgeschlossen ist.
Alle Writable-Stream-Implementierungen müssen eine Methode writable._write() und/oder writable._writev() bereitstellen, um Daten an die zugrunde liegende Ressource zu senden.
Transform-Streams stellen ihre eigene Implementierung von writable._write() bereit.
Diese Funktion DARF NICHT direkt vom Anwendungscode aufgerufen werden. Sie sollte von untergeordneten Klassen implementiert und nur von internen Writable-Klassenmethoden aufgerufen werden.
Die callback-Funktion muss synchron innerhalb von writable._write() oder asynchron (d. h. ein anderer Tick) aufgerufen werden, um zu signalisieren, dass der Schreibvorgang erfolgreich war oder mit einem Fehler fehlgeschlagen ist. Das erste an den callback übergebene Argument muss das Error-Objekt sein, wenn der Aufruf fehlgeschlagen ist, oder null, wenn der Schreibvorgang erfolgreich war.
Alle Aufrufe von writable.write(), die zwischen dem Zeitpunkt des Aufrufs von writable._write() und dem Aufruf des callback stattfinden, führen dazu, dass die geschriebenen Daten gepuffert werden. Wenn der callback aufgerufen wird, emittiert der Stream möglicherweise ein 'drain'-Ereignis. Wenn eine Stream-Implementierung in der Lage ist, mehrere Datenpakete gleichzeitig zu verarbeiten, sollte die Methode writable._writev() implementiert werden.
Wenn die decodeStrings-Eigenschaft in den Konstruktoroptionen explizit auf false gesetzt ist, bleibt chunk dasselbe Objekt, das an .write() übergeben wurde, und kann ein String anstelle eines Buffer sein. Dies dient der Unterstützung von Implementierungen, die über eine optimierte Handhabung für bestimmte String-Datenkodierungen verfügen. In diesem Fall gibt das encoding-Argument die Zeichenkodierung des Strings an. Andernfalls kann das encoding-Argument sicher ignoriert werden.
Die Methode writable._write() ist mit einem Unterstrich versehen, da sie intern für die Klasse ist, die sie definiert, und niemals direkt von Benutzerprogrammen aufgerufen werden sollte.
writable._writev(chunks, callback)#
chunks<Object[]>Die zu schreibenden Daten. Der Wert ist ein Array von<Object>, die jeweils ein diskretes Datenpaket zum Schreiben darstellen. Die Eigenschaften dieser Objekte sind:chunk<Buffer>|<string>Eine Buffer-Instanz oder ein String, die/der die zu schreibenden Daten enthält. Derchunkist ein String, wenn derWritablemit der auffalsegesetztendecodeStrings-Option erstellt wurde und ein String anwrite()übergeben wurde.encoding<string>Die Zeichenkodierung deschunk. WennchunkeinBufferist, ist dasencoding'buffer'.
callback<Function>Eine Callback-Funktion (optional mit einem Fehlerargument), die aufgerufen wird, wenn die Verarbeitung für die bereitgestellten Pakete abgeschlossen ist.
Diese Funktion DARF NICHT direkt vom Anwendungscode aufgerufen werden. Sie sollte von untergeordneten Klassen implementiert und nur von internen Writable-Klassenmethoden aufgerufen werden.
Die Methode writable._writev() kann zusätzlich oder alternativ zu writable._write() in Stream-Implementierungen implementiert werden, die in der Lage sind, mehrere Datenpakete gleichzeitig zu verarbeiten. Wenn sie implementiert ist und gepufferte Daten aus früheren Schreibvorgängen vorhanden sind, wird _writev() anstelle von _write() aufgerufen.
Die Methode writable._writev() ist mit einem Unterstrich versehen, da sie intern für die Klasse ist, die sie definiert, und niemals direkt von Benutzerprogrammen aufgerufen werden sollte.
writable._destroy(err, callback)#
err<Error>Ein möglicher Fehler.callback<Function>Eine Callback-Funktion, die ein optionales Fehlerargument entgegennimmt.
Die Methode _destroy() wird von writable.destroy() aufgerufen. Sie kann von untergeordneten Klassen überschrieben werden, DARF aber nicht direkt aufgerufen werden.
writable._final(callback)#
callback<Function>Rufen Sie diese Funktion auf (optional mit einem Fehlerargument), wenn das Schreiben aller verbleibenden Daten abgeschlossen ist.
Die Methode _final() DARF nicht direkt aufgerufen werden. Sie kann von untergeordneten Klassen implementiert werden und wird in diesem Fall nur von internen Writable-Klassenmethoden aufgerufen.
Diese optionale Funktion wird vor dem Schließen des Streams aufgerufen, wodurch das 'finish'-Ereignis verzögert wird, bis callback aufgerufen wird. Dies ist nützlich, um Ressourcen zu schließen oder gepufferte Daten zu schreiben, bevor ein Stream endet.
Fehler während des Schreibens#
Fehler, die während der Verarbeitung der Methoden writable._write(), writable._writev() und writable._final() auftreten, müssen durch Aufrufen des Callbacks und Weitergeben des Fehlers als erstes Argument propagiert werden. Das Auslösen eines Error innerhalb dieser Methoden oder das manuelle Emittieren eines 'error'-Ereignisses führt zu undefiniertem Verhalten.
Wenn ein Readable-Stream in einen Writable-Stream piped und der Writable-Stream einen Fehler emittiert, wird der Readable-Stream entpiped.
const { Writable } = require('node:stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
},
});
Ein Beispiel für einen beschreibbaren Stream#
Das Folgende illustriert eine recht simplistische (und etwas sinnlose) benutzerdefinierte Writable-Stream-Implementierung. Obwohl diese spezifische Writable-Stream-Instanz keinen wirklichen Nutzen hat, illustriert das Beispiel jedes der erforderlichen Elemente einer benutzerdefinierten Writable-Stream-Instanz.
const { Writable } = require('node:stream');
class MyWritable extends Writable {
_write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
}
}
Dekodieren von Puffern in einem beschreibbaren Stream#
Das Dekodieren von Puffern ist eine häufige Aufgabe, zum Beispiel bei der Verwendung von Transformern, deren Eingabe ein String ist. Dies ist kein trivialer Prozess bei der Verwendung von Mehrbyte-Zeichenkodierungen wie UTF-8. Das folgende Beispiel zeigt, wie man Mehrbyte-Strings unter Verwendung von StringDecoder und Writable dekodiert.
const { Writable } = require('node:stream');
const { StringDecoder } = require('node:string_decoder');
class StringWritable extends Writable {
constructor(options) {
super(options);
this._decoder = new StringDecoder(options?.defaultEncoding);
this.data = '';
}
_write(chunk, encoding, callback) {
if (encoding === 'buffer') {
chunk = this._decoder.write(chunk);
}
this.data += chunk;
callback();
}
_final(callback) {
this.data += this._decoder.end();
callback();
}
}
const euro = [[0xE2, 0x82], [0xAC]].map(Buffer.from);
const w = new StringWritable();
w.write('currency: ');
w.write(euro[0]);
w.end(euro[1]);
console.log(w.data); // currency: €
Implementierung eines lesbaren Streams#
Die Klasse stream.Readable wird erweitert, um einen Readable-Stream zu implementieren.
Benutzerdefinierte Readable-Streams müssen den new stream.Readable([options])-Konstruktor aufrufen und die Methode readable._read() implementieren.
new stream.Readable([options])#
options<Object>highWaterMark<number>Die maximale Anzahl von Bytes, die im internen Puffer gespeichert werden, bevor das Lesen aus der zugrunde liegenden Ressource eingestellt wird. Standard:65536(64 KiB) oder16fürobjectMode-Streams.encoding<string>Falls angegeben, werden Puffer unter Verwendung der angegebenen Kodierung in Strings dekodiert. Standard:null.objectMode<boolean>Ob dieser Stream sich wie ein Stream von Objekten verhalten soll. Das bedeutet, dassstream.read(n)einen einzelnen Wert anstelle einesBufferder Größenzurückgibt. Standard:false.emitClose<boolean>Ob der Stream'close'emittieren soll, nachdem er zerstört wurde. Standard:true.read<Function>Implementierung für die Methodestream._read().destroy<Function>Implementierung für die Methodestream._destroy().construct<Function>Implementierung für die Methodestream._construct().autoDestroy<boolean>Ob dieser Stream nach dem Beenden automatisch.destroy()auf sich selbst aufrufen soll. Standard:true.signal<AbortSignal>Ein Signal, das eine mögliche Stornierung darstellt.
const { Readable } = require('node:stream');
class MyReadable extends Readable {
constructor(options) {
// Calls the stream.Readable(options) constructor.
super(options);
// ...
}
}
Oder bei Verwendung von Konstruktoren im Vor-ES6-Stil
const { Readable } = require('node:stream');
const util = require('node:util');
function MyReadable(options) {
if (!(this instanceof MyReadable))
return new MyReadable(options);
Readable.call(this, options);
}
util.inherits(MyReadable, Readable);
Oder unter Verwendung des vereinfachten Konstruktoransatzes
const { Readable } = require('node:stream');
const myReadable = new Readable({
read(size) {
// ...
},
});
Der Aufruf von abort auf dem AbortController, der dem übergebenen AbortSignal entspricht, verhält sich genauso wie der Aufruf von .destroy(new AbortError()) auf dem erstellten Readable.
const { Readable } = require('node:stream');
const controller = new AbortController();
const read = new Readable({
read(size) {
// ...
},
signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort();
readable._construct(callback)#
callback<Function>Rufen Sie diese Funktion auf (optional mit einem Fehlerargument), wenn der Stream die Initialisierung abgeschlossen hat.
Die Methode _construct() DARF NICHT direkt aufgerufen werden. Sie kann von untergeordneten Klassen implementiert werden und wird in diesem Fall nur von internen Readable-Klassenmethoden aufgerufen.
Diese optionale Funktion wird im nächsten Tick vom Stream-Konstruktor geplant, wodurch alle _read()- und _destroy()-Aufrufe verzögert werden, bis callback aufgerufen wird. Dies ist nützlich, um den Zustand zu initialisieren oder Ressourcen asynchron zu initialisieren, bevor der Stream verwendet werden kann.
const { Readable } = require('node:stream');
const fs = require('node:fs');
class ReadStream extends Readable {
constructor(filename) {
super();
this.filename = filename;
this.fd = null;
}
_construct(callback) {
fs.open(this.filename, (err, fd) => {
if (err) {
callback(err);
} else {
this.fd = fd;
callback();
}
});
}
_read(n) {
const buf = Buffer.alloc(n);
fs.read(this.fd, buf, 0, n, null, (err, bytesRead) => {
if (err) {
this.destroy(err);
} else {
this.push(bytesRead > 0 ? buf.slice(0, bytesRead) : null);
}
});
}
_destroy(err, callback) {
if (this.fd) {
fs.close(this.fd, (er) => callback(er || err));
} else {
callback(err);
}
}
}
readable._read(size)#
size<number>Anzahl der asynchron zu lesenden Bytes.
Diese Funktion DARF NICHT direkt vom Anwendungscode aufgerufen werden. Sie sollte von untergeordneten Klassen implementiert und nur von internen Readable-Klassenmethoden aufgerufen werden.
Alle Readable-Stream-Implementierungen müssen eine Implementierung der Methode readable._read() bereitstellen, um Daten aus der zugrunde liegenden Ressource abzurufen.
Wenn readable._read() aufgerufen wird und Daten von der Ressource verfügbar sind, sollte die Implementierung beginnen, diese Daten mithilfe der Methode this.push(dataChunk) in die Lesewarteschlange zu pushen. _read() wird nach jedem Aufruf von this.push(dataChunk) erneut aufgerufen, sobald der Stream bereit ist, mehr Daten zu akzeptieren. _read() kann so lange weiter aus der Ressource lesen und Daten pushen, bis readable.push() false zurückgibt. Erst wenn _read() erneut aufgerufen wird, nachdem es gestoppt hat, sollte es das Pushen weiterer Daten in die Warteschlange fortsetzen.
Sobald die Methode readable._read() aufgerufen wurde, wird sie erst wieder aufgerufen, wenn über die Methode readable.push() weitere Daten gepusht wurden. Leere Daten wie leere Puffer und Strings führen nicht dazu, dass readable._read() aufgerufen wird.
Das size-Argument ist beratend. Für Implementierungen, bei denen ein "Lesen" eine einzelne Operation ist, die Daten zurückgibt, kann das size-Argument verwendet werden, um zu bestimmen, wie viele Daten abgerufen werden sollen. Andere Implementierungen können dieses Argument ignorieren und einfach Daten bereitstellen, wann immer sie verfügbar werden. Es besteht keine Notwendigkeit zu "warten", bis size Bytes verfügbar sind, bevor stream.push(chunk) aufgerufen wird.
Die Methode readable._read() ist mit einem Unterstrich versehen, da sie intern für die Klasse ist, die sie definiert, und niemals direkt von Benutzerprogrammen aufgerufen werden sollte.
readable._destroy(err, callback)#
err<Error>Ein möglicher Fehler.callback<Function>Eine Callback-Funktion, die ein optionales Fehlerargument entgegennimmt.
Die Methode _destroy() wird von readable.destroy() aufgerufen. Sie kann von untergeordneten Klassen überschrieben werden, DARF aber nicht direkt aufgerufen werden.
readable.push(chunk[, encoding])#
chunk<Buffer>|<TypedArray>|<DataView>|<string>|<null>|<any>Datenpaket, das in die Lesewarteschlange gepusht werden soll. Bei Streams, die nicht im Objektmodus arbeiten, musschunkein<string>,<Buffer>,<TypedArray>oder<DataView>sein. Für Streams im Objektmodus kannchunkjeder JavaScript-Wert sein.encoding<string>Kodierung von String-Datenstücken. Muss eine gültigeBuffer-Kodierung sein, wie z. B.'utf8'oder'ascii'.- Rückgabe:
<boolean>true, wenn weiterhin zusätzliche Datenpakete gepusht werden dürfen; andernfallsfalse.
Wenn chunk ein <Buffer>, <TypedArray>, <DataView> oder <string> ist, wird das Datenpaket zur internen Warteschlange hinzugefügt, damit Benutzer des Streams es konsumieren können. Das Übergeben von chunk als null signalisiert das Ende des Streams (EOF), nach dem keine Daten mehr geschrieben werden können.
Wenn der Readable im Pausenmodus arbeitet, können die mit readable.push() hinzugefügten Daten durch Aufrufen der Methode readable.read() gelesen werden, wenn das 'readable'-Ereignis emittiert wird.
Wenn der Readable im fließenden Modus (flowing mode) arbeitet, werden die mit readable.push() hinzugefügten Daten durch Emittieren eines 'data'-Ereignisses geliefert.
Die Methode readable.push() ist so flexibel wie möglich gestaltet. Wenn man beispielsweise eine Low-Level-Quelle umhüllt, die irgendeine Form von Pause/Wiederaufnahme-Mechanismus und einen Daten-Callback bietet, kann die Low-Level-Quelle von der benutzerdefinierten Readable-Instanz umhüllt werden.
// `_source` is an object with readStop() and readStart() methods,
// and an `ondata` member that gets called when it has data, and
// an `onend` member that gets called when the data is over.
class SourceWrapper extends Readable {
constructor(options) {
super(options);
this._source = getLowLevelSourceObject();
// Every time there's data, push it into the internal buffer.
this._source.ondata = (chunk) => {
// If push() returns false, then stop reading from source.
if (!this.push(chunk))
this._source.readStop();
};
// When the source ends, push the EOF-signaling `null` chunk.
this._source.onend = () => {
this.push(null);
};
}
// _read() will be called when the stream wants to pull more data in.
// The advisory size argument is ignored in this case.
_read(size) {
this._source.readStart();
}
}
Die Methode readable.push() wird verwendet, um den Inhalt in den internen Puffer zu pushen. Sie kann durch die Methode readable._read() gesteuert werden.
Bei Streams, die nicht im Objektmodus arbeiten: Wenn der chunk-Parameter von readable.push() undefined ist, wird er als leerer String oder Puffer behandelt. Siehe readable.push('') für weitere Informationen.
Fehler während des Lesens#
Fehler, die während der Verarbeitung von readable._read() auftreten, müssen durch die Methode readable.destroy(err) propagiert werden. Das Auslösen eines Error innerhalb von readable._read() oder das manuelle Emittieren eines 'error'-Ereignisses führt zu undefiniertem Verhalten.
const { Readable } = require('node:stream');
const myReadable = new Readable({
read(size) {
const err = checkSomeErrorCondition();
if (err) {
this.destroy(err);
} else {
// Do some work.
}
},
});
Ein Beispiel für einen Zähl-Stream#
Das Folgende ist ein grundlegendes Beispiel für einen Readable-Stream, der die Zahlen von 1 bis 1.000.000 in aufsteigender Reihenfolge emittiert und dann endet.
const { Readable } = require('node:stream');
class Counter extends Readable {
constructor(opt) {
super(opt);
this._max = 1000000;
this._index = 1;
}
_read() {
const i = this._index++;
if (i > this._max)
this.push(null);
else {
const str = String(i);
const buf = Buffer.from(str, 'ascii');
this.push(buf);
}
}
}
Implementierung eines Duplex-Streams#
Ein Duplex-Stream implementiert sowohl Readable als auch Writable, wie etwa eine TCP-Socket-Verbindung.
Da JavaScript keine Unterstützung für Mehrfachvererbung hat, wird die Klasse stream.Duplex erweitert, um einen Duplex-Stream zu implementieren (im Gegensatz zur Erweiterung der Klassen stream.Readable und stream.Writable).
Die Klasse stream.Duplex erbt prototypisch von stream.Readable und parasitär von stream.Writable, aber instanceof funktioniert ordnungsgemäß für beide Basisklassen aufgrund des Überschreibens von Symbol.hasInstance auf stream.Writable.
Benutzerdefinierte Duplex-Streams müssen den new stream.Duplex([options])-Konstruktor aufrufen und beide Methoden readable._read() und writable._write() implementieren.
new stream.Duplex(options)#
options<Object>Wird an beideWritable- undReadable-Konstruktoren übergeben. Hat auch die folgenden Felder:allowHalfOpen<boolean>Wenn auffalsegesetzt, beendet der Stream die Writable-Seite automatisch, wenn die Readable-Seite endet. Standard:true.readable<boolean>Legt fest, ob derDuplexlesbar sein soll. Standard:true.writable<boolean>Legt fest, ob derDuplexbeschreibbar sein soll. Standard:true.readableObjectMode<boolean>SetztobjectModefür die lesbare Seite des Streams. Hat keine Auswirkung, wennobjectModetrueist. Standard:false.writableObjectMode<boolean>SetztobjectModefür die beschreibbare Seite des Streams. Hat keine Auswirkung, wennobjectModetrueist. Standard:false.readableHighWaterMark<number>SetzthighWaterMarkfür die lesbare Seite des Streams. Hat keine Auswirkung, wennhighWaterMarkbereitgestellt wird.writableHighWaterMark<number>SetzthighWaterMarkfür die beschreibbare Seite des Streams. Hat keine Auswirkung, wennhighWaterMarkbereitgestellt wird.
const { Duplex } = require('node:stream');
class MyDuplex extends Duplex {
constructor(options) {
super(options);
// ...
}
}
Oder bei Verwendung von Konstruktoren im Vor-ES6-Stil
const { Duplex } = require('node:stream');
const util = require('node:util');
function MyDuplex(options) {
if (!(this instanceof MyDuplex))
return new MyDuplex(options);
Duplex.call(this, options);
}
util.inherits(MyDuplex, Duplex);
Oder unter Verwendung des vereinfachten Konstruktoransatzes
const { Duplex } = require('node:stream');
const myDuplex = new Duplex({
read(size) {
// ...
},
write(chunk, encoding, callback) {
// ...
},
});
Bei Verwendung von Pipeline
const { Transform, pipeline } = require('node:stream');
const fs = require('node:fs');
pipeline(
fs.createReadStream('object.json')
.setEncoding('utf8'),
new Transform({
decodeStrings: false, // Accept string input rather than Buffers
construct(callback) {
this.data = '';
callback();
},
transform(chunk, encoding, callback) {
this.data += chunk;
callback();
},
flush(callback) {
try {
// Make sure is valid json.
JSON.parse(this.data);
this.push(this.data);
callback();
} catch (err) {
callback(err);
}
},
}),
fs.createWriteStream('valid-object.json'),
(err) => {
if (err) {
console.error('failed', err);
} else {
console.log('completed');
}
},
);
Ein Beispiel für einen Duplex-Stream#
Das Folgende illustriert ein einfaches Beispiel für einen Duplex-Stream, der ein hypothetisches Low-Level-Quellobjekt umhüllt, in das Daten geschrieben und aus dem gelesen werden können, wenn auch unter Verwendung einer API, die nicht mit Node.js-Streams kompatibel ist. Das Beispiel zeigt einen einfachen Duplex-Stream, der eingehende geschriebene Daten über die Writable-Schnittstelle puffert, die dann über die Readable-Schnittstelle wieder ausgelesen werden.
const { Duplex } = require('node:stream');
const kSource = Symbol('source');
class MyDuplex extends Duplex {
constructor(source, options) {
super(options);
this[kSource] = source;
}
_write(chunk, encoding, callback) {
// The underlying source only deals with strings.
if (Buffer.isBuffer(chunk))
chunk = chunk.toString();
this[kSource].writeSomeData(chunk);
callback();
}
_read(size) {
this[kSource].fetchSomeData(size, (data, encoding) => {
this.push(Buffer.from(data, encoding));
});
}
}
Der wichtigste Aspekt eines Duplex-Streams ist, dass die Readable- und Writable-Seiten unabhängig voneinander arbeiten, obwohl sie in einer einzigen Objektinstanz koexistieren.
Duplex-Streams im Objektmodus#
Für Duplex-Streams kann objectMode mithilfe der Optionen readableObjectMode und writableObjectMode ausschließlich für die Readable- oder Writable-Seite gesetzt werden.
Im folgenden Beispiel wird beispielsweise ein neuer Transform-Stream (der ein Typ von Duplex-Stream ist) erstellt, der eine Writable-Seite im Objektmodus hat, die JavaScript-Zahlen akzeptiert, welche auf der Readable-Seite in hexadezimale Strings konvertiert werden.
const { Transform } = require('node:stream');
// All Transform streams are also Duplex Streams.
const myTransform = new Transform({
writableObjectMode: true,
transform(chunk, encoding, callback) {
// Coerce the chunk to a number if necessary.
chunk |= 0;
// Transform the chunk into something else.
const data = chunk.toString(16);
// Push the data onto the readable queue.
callback(null, '0'.repeat(data.length % 2) + data);
},
});
myTransform.setEncoding('ascii');
myTransform.on('data', (chunk) => console.log(chunk));
myTransform.write(1);
// Prints: 01
myTransform.write(10);
// Prints: 0a
myTransform.write(100);
// Prints: 64
Implementierung eines Transform-Streams#
Ein Transform-Stream ist ein Duplex-Stream, bei dem die Ausgabe in irgendeiner Weise aus der Eingabe berechnet wird. Beispiele sind zlib-Streams oder crypto-Streams, die Daten komprimieren, verschlüsseln oder entschlüsseln.
Es gibt keine Anforderung, dass die Ausgabe die gleiche Größe wie die Eingabe, die gleiche Anzahl von Paketen hat oder zur gleichen Zeit ankommt. Zum Beispiel wird ein Hash-Stream immer nur ein einzelnes Ausgabepaket haben, das bereitgestellt wird, wenn die Eingabe endet. Ein zlib-Stream erzeugt eine Ausgabe, die entweder viel kleiner oder viel größer als seine Eingabe ist.
Die Klasse stream.Transform wird erweitert, um einen Transform-Stream zu implementieren.
Die Klasse stream.Transform erbt prototypisch von stream.Duplex und implementiert ihre eigenen Versionen der Methoden writable._write() und readable._read(). Benutzerdefinierte Transform-Implementierungen müssen die Methode transform._transform() implementieren und können auch die Methode transform._flush() implementieren.
Bei der Verwendung von Transform-Streams ist Vorsicht geboten, da in den Stream geschriebene Daten dazu führen können, dass die Writable-Seite des Streams pausiert wird, wenn die Ausgabe auf der Readable-Seite nicht konsumiert wird.
new stream.Transform([options])#
options<Object>Wird an beideWritable- undReadable-Konstruktoren übergeben. Hat auch die folgenden Felder:transform<Function>Implementierung für die Methodestream._transform().flush<Function>Implementierung für die Methodestream._flush().
const { Transform } = require('node:stream');
class MyTransform extends Transform {
constructor(options) {
super(options);
// ...
}
}
Oder bei Verwendung von Konstruktoren im Vor-ES6-Stil
const { Transform } = require('node:stream');
const util = require('node:util');
function MyTransform(options) {
if (!(this instanceof MyTransform))
return new MyTransform(options);
Transform.call(this, options);
}
util.inherits(MyTransform, Transform);
Oder unter Verwendung des vereinfachten Konstruktoransatzes
const { Transform } = require('node:stream');
const myTransform = new Transform({
transform(chunk, encoding, callback) {
// ...
},
});
Ereignis: 'end'#
Das 'end'-Ereignis stammt von der Klasse stream.Readable. Das 'end'-Ereignis wird emittiert, nachdem alle Daten ausgegeben wurden, was nach dem Aufruf des Callbacks in transform._flush() erfolgt. Im Fehlerfall sollte 'end' nicht emittiert werden.
Ereignis: 'finish'#
Das 'finish'-Ereignis stammt von der Klasse stream.Writable. Das 'finish'-Ereignis wird emittiert, nachdem stream.end() aufgerufen wurde und alle Datenpakete von stream._transform() verarbeitet wurden. Im Fehlerfall sollte 'finish' nicht emittiert werden.
transform._flush(callback)#
callback<Function>Eine Callback-Funktion (optional mit einem Fehlerargument und Daten), die aufgerufen wird, wenn verbleibende Daten geflusht wurden.
Diese Funktion DARF NICHT direkt vom Anwendungscode aufgerufen werden. Sie sollte von untergeordneten Klassen implementiert und nur von internen Readable-Klassenmethoden aufgerufen werden.
In einigen Fällen muss eine Transformationsoperation am Ende des Streams ein zusätzliches Stück Daten emittieren. Zum Beispiel speichert ein zlib-Komprimierungsstream eine Menge internen Zustands, der zur optimalen Komprimierung der Ausgabe verwendet wird. Wenn der Stream endet, müssen diese zusätzlichen Daten jedoch geflusht werden, damit die komprimierten Daten vollständig sind.
Benutzerdefinierte Transform-Implementierungen können die Methode transform._flush() implementieren. Diese wird aufgerufen, wenn keine weiteren Daten zum Konsumieren geschrieben wurden, aber bevor das 'end'-Ereignis emittiert wird, das das Ende des Readable-Streams signalisiert.
Innerhalb der transform._flush()-Implementierung kann die Methode transform.push() null oder mehrmals aufgerufen werden, wie angemessen. Die callback-Funktion muss aufgerufen werden, wenn die Flush-Operation abgeschlossen ist.
Die Methode transform._flush() ist mit einem Unterstrich versehen, da sie intern für die Klasse ist, die sie definiert, und niemals direkt von Benutzerprogrammen aufgerufen werden sollte.
transform._transform(chunk, encoding, callback)#
chunk<Buffer>|<string>|<any>Der zu transformierendeBuffer, konvertiert aus dem anstream.write()übergebenenstring. Wenn diedecodeStrings-Option des Streamsfalseist oder der Stream im Objektmodus arbeitet, wird das Paket nicht konvertiert und ist das, was anstream.write()übergeben wurde.encoding<string>Wenn das Paket ein String ist, ist dies der Kodierungstyp. Wenn das Paket ein Buffer ist, ist dies der spezielle Wert'buffer'. Ignorieren Sie es in diesem Fall.callback<Function>Eine Callback-Funktion (optional mit einem Fehlerargument und Daten), die aufgerufen wird, nachdem das bereitgestellte Datenpaketchunkverarbeitet wurde.
Diese Funktion DARF NICHT direkt vom Anwendungscode aufgerufen werden. Sie sollte von untergeordneten Klassen implementiert und nur von internen Readable-Klassenmethoden aufgerufen werden.
Alle Transform-Stream-Implementierungen müssen eine Methode _transform() bereitstellen, um Eingaben zu akzeptieren und Ausgaben zu erzeugen. Die transform._transform()-Implementierung handhabt die zu schreibenden Bytes, berechnet eine Ausgabe und leitet diese dann unter Verwendung der Methode transform.push() an den lesbaren Teil weiter.
Die Methode transform.push() kann null- oder mehrmals aufgerufen werden, um aus einem einzelnen Eingabeblock (Chunk) eine Ausgabe zu generieren, je nachdem, wie viel als Ergebnis des Blocks ausgegeben werden soll.
Es ist möglich, dass für einen bestimmten Block von Eingabedaten keine Ausgabe generiert wird.
Die callback-Funktion darf nur aufgerufen werden, wenn der aktuelle Block vollständig verarbeitet wurde. Das erste Argument, das an den callback übergeben wird, muss ein Error-Objekt sein, falls während der Verarbeitung der Eingabe ein Fehler aufgetreten ist, andernfalls null. Wenn ein zweites Argument an den callback übergeben wird, wird dieses an die Methode transform.push() weitergeleitet, jedoch nur, wenn das erste Argument "falsy" (falsch) ist. Mit anderen Worten: Folgendes ist äquivalent:
transform.prototype._transform = function(data, encoding, callback) {
this.push(data);
callback();
};
transform.prototype._transform = function(data, encoding, callback) {
callback(null, data);
};
Die Methode transform._transform() ist mit einem Unterstrich versehen, da sie intern für die Klasse ist, die sie definiert, und niemals direkt von Benutzerprogrammen aufgerufen werden sollte.
transform._transform() wird niemals parallel aufgerufen; Streams implementieren einen Warteschlangenmechanismus, und um den nächsten Block zu erhalten, muss der callback aufgerufen werden, entweder synchron oder asynchron.
Klasse: stream.PassThrough#
Die Klasse stream.PassThrough ist eine triviale Implementierung eines Transform-Streams, der die Eingabebytes einfach an die Ausgabe weiterreicht. Ihr Zweck dient primär Beispielen und Tests, aber es gibt einige Anwendungsfälle, in denen stream.PassThrough als Baustein für neuartige Arten von Streams nützlich ist.
Zusätzliche Hinweise#
Kompatibilität von Streams mit asynchronen Generatoren und asynchronen Iteratoren#
Mit der Unterstützung von asynchronen Generatoren und Iteratoren in JavaScript sind asynchrone Generatoren zu diesem Zeitpunkt effektiv ein Stream-Konstrukt erster Klasse auf Sprachebene.
Einige gängige Interop-Fälle für die Verwendung von Node.js-Streams mit asynchronen Generatoren und asynchronen Iteratoren werden unten aufgeführt.
Verarbeitung von lesbaren Streams mit asynchronen Iteratoren#
(async function() {
for await (const chunk of readable) {
console.log(chunk);
}
})();
Asynchrone Iteratoren registrieren einen permanenten Fehler-Handler für den Stream, um unvorhergesehene Fehler nach der Zerstörung (Post-Destroy-Errors) zu verhindern.
Erstellung von lesbaren Streams mit asynchronen Generatoren#
Ein lesbarer Node.js-Stream kann mithilfe der Hilfsmethode Readable.from() aus einem asynchronen Generator erstellt werden.
const { Readable } = require('node:stream');
const ac = new AbortController();
const signal = ac.signal;
async function * generate() {
yield 'a';
await someLongRunningFn({ signal });
yield 'b';
yield 'c';
}
const readable = Readable.from(generate());
readable.on('close', () => {
ac.abort();
});
readable.on('data', (chunk) => {
console.log(chunk);
});
Weiterleitung an schreibbare Streams von asynchronen Iteratoren#
Beim Schreiben auf einen schreibbaren Stream von einem asynchronen Iterator aus ist auf die korrekte Handhabung von Backpressure und Fehlern zu achten. stream.pipeline() abstrahiert die Handhabung von Backpressure und Backpressure-bezogenen Fehlern.
const fs = require('node:fs');
const { pipeline } = require('node:stream');
const { pipeline: pipelinePromise } = require('node:stream/promises');
const writable = fs.createWriteStream('./file');
const ac = new AbortController();
const signal = ac.signal;
const iterator = createIterator({ signal });
// Callback Pattern
pipeline(iterator, writable, (err, value) => {
if (err) {
console.error(err);
} else {
console.log(value, 'value returned');
}
}).on('close', () => {
ac.abort();
});
// Promise Pattern
pipelinePromise(iterator, writable)
.then((value) => {
console.log(value, 'value returned');
})
.catch((err) => {
console.error(err);
ac.abort();
});
Kompatibilität mit älteren Node.js-Versionen#
Vor Node.js 0.10 war die Readable-Stream-Schnittstelle einfacher, aber auch weniger leistungsfähig und nützlich.
- Anstatt auf Aufrufe der Methode
stream.read()zu warten, würden'data'-Ereignisse sofort ausgelöst. Anwendungen, die einen gewissen Arbeitsaufwand betreiben mussten, um zu entscheiden, wie Daten zu behandeln sind, mussten gelesene Daten in Puffern speichern, damit die Daten nicht verloren gingen. - Die Methode
stream.pause()war eher beratend als garantiert. Dies bedeutete, dass es weiterhin notwendig war, darauf vorbereitet zu sein,'data'-Ereignisse zu empfangen, selbst wenn sich der Stream in einem pausierten Zustand befand.
In Node.js 0.10 wurde die Klasse Readable hinzugefügt. Zur Abwärtskompatibilität mit älteren Node.js-Programmen wechseln Readable-Streams in den „Flowing Mode“, wenn ein 'data'-Ereignis-Handler hinzugefügt wird oder wenn die Methode stream.resume() aufgerufen wird. Der Effekt ist, dass man sich selbst bei Nichtverwendung der neuen Methode stream.read() und des 'readable'-Ereignisses keine Sorgen mehr über den Verlust von 'data'-Blöcken machen muss.
Während die meisten Anwendungen weiterhin normal funktionieren, führt dies in den folgenden Bedingungen zu einem Grenzfall:
- Es wird kein
'data'-Ereignis-Listener hinzugefügt. - Die Methode
stream.resume()wird niemals aufgerufen. - Der Stream wird an kein schreibbares Ziel weitergeleitet.
Betrachten Sie zum Beispiel den folgenden Code:
// WARNING! BROKEN!
net.createServer((socket) => {
// We add an 'end' listener, but never consume the data.
socket.on('end', () => {
// It will never get here.
socket.end('The message was received but was not processed.\n');
});
}).listen(1337);
Vor Node.js 0.10 wurden die eingehenden Nachrichtendaten einfach verworfen. In Node.js 0.10 und neuer bleibt der Socket jedoch für immer pausiert.
Die Problemumgehung in dieser Situation besteht darin, die Methode stream.resume() aufzurufen, um den Datenfluss zu starten.
// Workaround.
net.createServer((socket) => {
socket.on('end', () => {
socket.end('The message was received but was not processed.\n');
});
// Start the flow of data, discarding it.
socket.resume();
}).listen(1337);
Zusätzlich dazu, dass neue Readable-Streams in den Flowing-Modus wechseln, können Streams im Stil vor 0.10 mithilfe der Methode readable.wrap() in eine Readable-Klasse gewrappt werden.
readable.read(0)#
Es gibt einige Fälle, in denen es notwendig ist, eine Aktualisierung der zugrunde liegenden lesbaren Stream-Mechanismen auszulösen, ohne tatsächlich Daten zu konsumieren. In solchen Fällen ist es möglich, readable.read(0) aufzurufen, was immer null zurückgibt.
Wenn der interne Lesepuffer unter dem highWaterMark liegt und der Stream derzeit nicht liest, löst der Aufruf von stream.read(0) einen Low-Level-Aufruf von stream._read() aus.
Während die meisten Anwendungen dies fast nie tun müssen, gibt es innerhalb von Node.js Situationen, in denen dies getan wird, insbesondere in den Interna der Readable-Stream-Klasse.
readable.push('')#
Die Verwendung von readable.push('') wird nicht empfohlen.
Das Pushen eines Null-Byte <string>, <Buffer>, <TypedArray> oder <DataView> an einen Stream, der sich nicht im Objektmodus befindet, hat einen interessanten Nebeneffekt. Da es ein Aufruf an readable.push() ist, beendet der Aufruf den Lesevorgang. Da das Argument jedoch ein leerer String ist, werden dem lesbaren Puffer keine Daten hinzugefügt, sodass es nichts gibt, was ein Benutzer konsumieren könnte.
highWaterMark-Diskrepanz nach dem Aufruf von readable.setEncoding()#
Die Verwendung von readable.setEncoding() ändert das Verhalten, wie der highWaterMark im Nicht-Objektmodus arbeitet.
Typischerweise wird die Größe des aktuellen Puffers anhand des highWaterMark in Bytes gemessen. Nachdem jedoch setEncoding() aufgerufen wurde, beginnt die Vergleichsfunktion, die Größe des Puffers in Zeichen zu messen.
Dies ist in üblichen Fällen mit latin1 oder ascii kein Problem. Es wird jedoch empfohlen, bei der Arbeit mit Strings, die Multi-Byte-Zeichen enthalten könnten, auf dieses Verhalten zu achten.