Поток[src]
Исходный код: lib/stream.js
Поток — это абстрактный интерфейс для работы с потоковыми данными в Node.js. Модуль node:stream предоставляет API для реализации интерфейса потоков.
Node.js предоставляет множество потоковых объектов. Например, запрос к HTTP-серверу и process.stdout являются экземплярами потоков.
Потоки могут быть читаемыми, записываемыми или обоими сразу. Все потоки являются экземплярами EventEmitter.
Чтобы получить доступ к модулю node:stream:
const stream = require('node:stream'); copy Модуль node:stream полезен для создания новых типов экземпляров потоков. Обычно нет необходимости использовать модуль node:stream для чтения потоков.
Структура этого документа
Этот документ состоит из двух основных разделов и третьего раздела с примечаниями. В первом разделе объясняется, как использовать существующие потоки в приложении. Во втором разделе объясняется, как создавать новые типы потоков.
Типы потоков
В Node.js есть четыре основных типа потоков:
-
Writable: потоки, в которые можно записывать данные (например,fs.createWriteStream()). -
Readable: потоки, из которых можно читать данные (например,fs.createReadStream()). -
Duplex: потоки, которые одновременно являютсяReadableиWritable(например,net.Socket). -
Transform: потокиDuplex, которые могут изменять или преобразовывать данные при их записи и чтении (например,zlib.createDeflate()).
Кроме того, этот модуль включает вспомогательные функции stream.duplexPair(), stream.pipeline(), stream.finished(), stream.Readable.from() и stream.addAbortSignal().
API потоков на основе промисов
API stream/promises предоставляет альтернативный набор асинхронных вспомогательных функций для потоков, которые возвращают объекты Promise вместо использования обратных вызовов. Доступ к API можно получить через require('node:stream/promises') или require('node:stream').promises.
stream.pipeline(streams[, options])
stream.pipeline(source[, ...transforms], destination[, options])
-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> -
source<Stream> | <Iterable> | <AsyncIterable> | <Function>- Возвращает: <Promise> | <AsyncIterable>
-
...transforms<Stream> | <Function>-
source<AsyncIterable> - Возвращает: <Promise> | <AsyncIterable>
-
-
destination<Stream> | <Function>-
source<AsyncIterable> - Возвращает: <Promise> | <AsyncIterable>
-
-
options<Object> Параметры конвейера-
signal<AbortSignal> -
end<boolean> Завершать поток назначения при завершении исходного потока. Потоки преобразования завершаются всегда, даже если это значение равноfalse. По умолчанию:true.
-
- Возвращает: <Promise> Выполняется после завершения конвейера.
CommonJS
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);Модули JavaScript
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.');Чтобы использовать AbortSignal, передайте его в объекте параметров последним аргументом. При отмене сигнала для базового конвейера будет вызван destroy с объектом AbortError.
CommonJS
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); // AbortErrorМодули JavaScript
import { 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
}API pipeline также поддерживает асинхронные генераторы:
CommonJS
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);Модули JavaScript
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.');Не забудьте обработать аргумент signal, переданный асинхронному генератору. Это особенно важно, если асинхронный генератор является источником конвейера (то есть первым аргументом), иначе конвейер никогда не завершится.
CommonJS
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);Модули JavaScript
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.');API pipeline также предоставляет версию с обратным вызовом:
stream.finished(stream[, options])
-
stream<Stream> | <ReadableStream> | <WritableStream> Читаемый и/или записываемый поток/веб-поток. -
options<Object>-
error<boolean> | <undefined> -
readable<boolean> | <undefined> -
writable<boolean> | <undefined> -
signal<AbortSignal> | <undefined> -
cleanup<boolean> | <undefined> Еслиtrue, удаляет зарегистрированные этой функцией обработчики событий до выполнения промиса. По умолчанию:false.
-
- Возвращает: <Promise> Выполняется, когда поток становится нечитаемым или недоступным для записи.
CommonJS
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.Модули JavaScript
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.API finished также предоставляет версию с обратным вызовом.
stream.finished() оставляет висящие обработчики событий (в частности, 'error', 'end', 'finish' и 'close') после выполнения или отклонения возвращённого промиса. Это сделано для того, чтобы неожиданные события 'error' (из-за неправильной реализации потоков) не приводили к неожиданным сбоям. Если такое поведение нежелательно, следует установить для options.cleanup значение true:
await finished(rs, { cleanup: true }); copy Режим объектов
Все потоки, создаваемые API Node.js, работают исключительно со строками, объектами <Buffer>, <TypedArray> и <DataView>:
-
StringsиBuffers— наиболее распространённые типы, используемые в потоках. -
TypedArrayиDataViewпозволяют работать с двоичными данными в таких типах, какInt32ArrayилиUint8Array. Когда в поток записывается TypedArray или DataView, Node.js обрабатывает необработанные байты.
Однако реализации потоков могут работать и с другими типами значений JavaScript (за исключением null, который выполняет в потоках особую функцию). Такие потоки считаются работающими в «режиме объектов».
Экземпляры потоков переводятся в режим объектов с помощью параметра objectMode при создании потока. Переводить существующий поток в режим объектов небезопасно.
Буферизация
Как потоки Writable, так и потоки Readable хранят данные во внутреннем буфере.
Объём потенциально буферизуемых данных зависит от параметра highWaterMark, переданного конструктору потока. Для обычных потоков параметр highWaterMark задаёт общее количество байтов. Для потоков, работающих в режиме объектов, highWaterMark задаёт общее количество объектов. Для потоков, работающих со строками (но не декодирующих их), highWaterMark задаёт общее количество кодовых единиц UTF-16.
Данные буферизуются в потоках Readable, когда реализация вызывает stream.push(chunk). Если потребитель потока не вызывает stream.read(), данные остаются во внутренней очереди до их потребления.
Когда общий размер внутреннего буфера чтения достигает порогового значения, заданного параметром highWaterMark, поток временно прекращает чтение данных из базового ресурса, пока буферизованные данные не будут потреблены (то есть поток перестаёт вызывать внутренний метод readable._read(), используемый для заполнения буфера чтения).
Данные буферизуются в потоках Writable при многократном вызове метода writable.write(chunk). Пока общий размер внутреннего буфера записи меньше порогового значения, заданного параметром highWaterMark, вызовы writable.write() возвращают true. Когда размер внутреннего буфера достигает или превышает highWaterMark, возвращается false.
Одна из ключевых целей API stream, особенно метода stream.pipe(), — ограничить буферизацию данных приемлемым уровнем, чтобы источники и получатели с разной скоростью не перегружали доступную память.
Параметр highWaterMark задаёт порог, а не ограничение: он определяет объём данных, который поток буферизует, прежде чем перестать запрашивать новые данные. Как правило, он не устанавливает строгого ограничения на использование памяти. Конкретные реализации потоков могут устанавливать более строгие ограничения, но это необязательно.
Поскольку потоки Duplex и Transform одновременно являются Readable и Writable, каждый из них поддерживает два отдельных внутренних буфера для чтения и записи. Это позволяет каждой стороне работать независимо от другой, сохраняя подходящий и эффективный поток данных. Например, экземпляры net.Socket — это потоки Duplex, чья сторона Readable позволяет потреблять данные, полученные из сокета, а сторона Writable позволяет записывать данные в сокет. Поскольку скорость записи данных в сокет может быть выше или ниже скорости их получения, каждая сторона должна работать (и буферизовать данные) независимо от другой.
Механизмы внутренней буферизации являются деталями внутренней реализации и могут быть изменены в любое время. Однако в некоторых сложных реализациях внутренние буферы можно получить с помощью writable.writableBuffer или readable.readableBuffer. Использовать эти недокументированные свойства не рекомендуется.
API для потребителей потоков
Практически все приложения Node.js, какими бы простыми они ни были, так или иначе используют потоки. Ниже приведён пример использования потоков в приложении Node.js, реализующем HTTP-сервер:
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 copy Потоки Writable (например, res в примере) предоставляют такие методы, как write() и end(), которые используются для записи данных в поток.
Потоки Readable используют API EventEmitter, чтобы уведомлять код приложения о доступности данных для чтения из потока. Эти данные можно считывать из потока несколькими способами.
Потоки Writable и Readable используют API EventEmitter различными способами для передачи текущего состояния потока.
Потоки Duplex и Transform являются одновременно потоками Writable и Readable.
Приложениям, которые записывают данные в поток или считывают их из потока, не требуется напрямую реализовывать интерфейсы потоков, и, как правило, у них нет причин вызывать require('node:stream').
Разработчикам, желающим реализовать новые типы потоков, следует обратиться к разделу API для разработчиков потоков.
Записываемые потоки
Записываемые потоки представляют собой абстракцию назначения, в которое записываются данные.
Примеры потоков Writable:
- HTTP-запросы на стороне клиента
- HTTP-ответы на стороне сервера
- потоки записи fs
- потоки zlib
- потоки crypto
- TCP-сокеты
- стандартный ввод дочернего процесса
-
process.stdout,process.stderr
Некоторые из этих примеров на самом деле являются потоками Duplex, реализующими интерфейс Writable.
Все потоки Writable реализуют интерфейс, определённый классом stream.Writable.
Хотя конкретные экземпляры потоков Writable могут различаться в разных аспектах, все потоки Writable следуют одному и тому же основному шаблону использования, показанному в примере ниже:
const myStream = getWritableStreamSomehow();
myStream.write('some data');
myStream.write('some more data');
myStream.end('done writing data'); copy Класс: stream.Writable
Событие: 'close'
Событие 'close' генерируется, когда поток и все его базовые ресурсы (например, файловый дескриптор) закрыты. Это событие указывает, что больше не будет генерироваться никаких событий и дальнейшие вычисления выполняться не будут.
Поток Writable всегда сгенерирует событие 'close', если он создан с параметром emitClose.
Событие: 'drain'
Если вызов stream.write(chunk) возвращает false, событие 'drain' будет сгенерировано, когда можно будет возобновить запись данных в поток.
// 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);
}
}
} copy Событие: 'error'
- Тип: <Error>
Событие 'error' генерируется, если при записи или передаче данных через поток произошла ошибка. При вызове обработчику передаётся единственный аргумент Error.
При генерации события 'error' поток закрывается, если при его создании параметру autoDestroy не было присвоено значение false.
После 'error' не должно генерироваться никаких дальнейших событий, кроме 'close' (включая события 'error').
Событие: 'finish'
Событие 'finish' генерируется после вызова метода stream.end() и сброса всех данных в базовую систему.
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'); copy Событие: 'pipe'
-
srcисходный поток <stream.Readable>, передающий данные в этот записываемый поток
Событие 'pipe' генерируется, когда для читаемого потока вызывается метод stream.pipe(), добавляющий этот записываемый поток в список назначений.
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); copy Событие: 'unpipe'
-
srcисходный поток <stream.Readable>, который отключил передачу в этот записываемый поток
Событие 'unpipe' генерируется, когда для потока Readable вызывается метод stream.unpipe(), удаляющий этот поток Writable из списка назначений.
Это событие также генерируется, если поток Writable выдаёт ошибку при передаче в него данных из потока Readable.
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); copy
writable.cork()
Метод writable.cork() принудительно буферизует все записанные данные в памяти. Буферизованные данные будут сброшены при вызове метода stream.uncork() или stream.end().
Основная цель writable.cork() — обеспечить обработку ситуации, когда несколько небольших фрагментов данных быстро записываются в поток один за другим. Вместо немедленной передачи в базовое назначение метод writable.cork() буферизует все фрагменты до вызова writable.uncork(), который передаст их все в writable._writev(), если он реализован. Это предотвращает блокировку в начале очереди, при которой данные буферизуются в ожидании обработки первого небольшого фрагмента. Однако использование writable.cork() без реализации writable._writev() может отрицательно сказаться на пропускной способности.
См. также: writable.uncork(), writable._writev().
writable.destroy([error])
Уничтожает поток. При необходимости генерирует событие 'error' и событие 'close' (если только для emitClose не задано значение false). После этого вызова записываемый поток завершает работу, а последующие вызовы write() или end() приведут к ошибке ERR_STREAM_DESTROYED. Это разрушительный и немедленный способ уничтожить поток. Данные, передаваемые предыдущими вызовами write(), могли ещё не быть сброшены, что может привести к ошибке ERR_STREAM_DESTROYED. Если перед закрытием необходимо сбросить данные, вместо destroy используйте end() или дождитесь события 'drain', прежде чем уничтожать поток.
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 copy const { Writable } = require('node:stream');
const myStream = new Writable();
myStream.destroy();
myStream.on('error', function wontHappen() {}); copy const { Writable } = require('node:stream');
const myStream = new Writable();
myStream.destroy();
myStream.write('foo', (error) => console.error(error.code));
// ERR_STREAM_DESTROYED copy После вызова destroy() все дальнейшие вызовы ничего не делают, и никакие дальнейшие ошибки, кроме ошибок из _destroy(), не будут генерироваться как 'error'.
Разработчикам не следует переопределять этот метод; вместо этого следует реализовать writable._destroy().
writable.destroyed
- Тип: <boolean>
Равно true после вызова writable.destroy().
const { Writable } = require('node:stream');
const myStream = new Writable();
console.log(myStream.destroyed); // false
myStream.destroy();
console.log(myStream.destroyed); // true copy
writable.end([chunk[, encoding]][, callback])
-
chunkНеобязательные данные для записи. Для потоков, работающих не в объектном режиме,chunkдолжен быть значением типа <string>, <Buffer>, <TypedArray> или <DataView>. Для потоков в объектном режимеchunkможет быть любым значением JavaScript, кромеnull. <string> | <Buffer> | <TypedArray> | <DataView> | <any> -
encodingКодировка, еслиchunk— строка <string> -
callbackОбратный вызов по завершении работы потока. <Function> - Возвращает: <this>
Вызов метода writable.end() сообщает, что в Writable больше не будут записываться данные. Необязательные аргументы chunk и encoding позволяют записать последний дополнительный фрагмент данных непосредственно перед закрытием потока.
Вызов метода stream.write() после вызова stream.end() приведёт к ошибке.
// 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! copy
writable.setDefaultEncoding(encoding)
Метод writable.setDefaultEncoding() задаёт кодировку encoding по умолчанию для потока Writable.
writable.uncork()
Метод writable.uncork() сбрасывает все данные, буферизованные с момента вызова stream.cork().
При использовании writable.cork() и writable.uncork() для управления буферизацией записей в поток откладывайте вызовы writable.uncork() с помощью process.nextTick(). Это позволяет объединять все вызовы writable.write(), происходящие в течение одной фазы цикла событий Node.js.
stream.cork();
stream.write('some ');
stream.write('data ');
process.nextTick(() => stream.uncork()); copy Если метод writable.cork() вызывается для потока несколько раз, для сброса буферизованных данных необходимо вызвать writable.uncork() столько же раз.
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();
}); copy См. также: writable.cork().
writable.writable
- Тип: <boolean>
Равно true, если вызов writable.write() безопасен, то есть поток не был уничтожен, в нём не произошла ошибка и его работа не завершилась.
writable.writableAborted
- Тип: <boolean>
Возвращает признак того, что поток был уничтожен или в нём произошла ошибка до генерации 'finish'.
writable.writableEnded
- Тип: <boolean>
Равно true после вызова writable.end(). Это свойство не указывает, были ли сброшены данные; для этого используйте writable.writableFinished.
writable.writableCorked
- Тип: <integer>
Количество вызовов writable.uncork(), необходимых для полной разблокировки потока.
writable.writableFinished
- Тип: <boolean>
Устанавливается в true непосредственно перед генерацией события 'finish'.
writable.writableHighWaterMark
- Тип: <number>
Возвращает значение highWaterMark, переданное при создании этого Writable.
writable.writableLength
- Тип: <number>
Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение позволяет узнать состояние highWaterMark.
writable.writableNeedDrain
- Тип: <boolean>
Равно true, если буфер потока был заполнен и поток сгенерирует 'drain'.
writable[Symbol.asyncDispose]()
Вызывает writable.destroy() с объектом AbortError и возвращает промис, который выполняется после завершения потока.
writable.write(chunk[, encoding][, callback])
-
chunkНеобязательные данные для записи. Для потоков, работающих не в объектном режиме,chunkдолжен быть значением типа <string>, <Buffer>, <TypedArray> или <DataView>. Для потоков в объектном режимеchunkможет быть любым значением JavaScript, кромеnull. <string> | <Buffer> | <TypedArray> | <DataView> | <any> -
encodingКодировка, еслиchunk— строка. По умолчанию:'utf8'<string> | <null> -
callbackОбратный вызов после сброса этого фрагмента данных. <Function> - Возвращает: <boolean>
false, если поток ожидает, что вызывающий код дождётся генерации события'drain', прежде чем продолжать запись дополнительных данных; в противном случае —true.
Метод writable.write() записывает данные в поток и вызывает переданный callback после полной обработки данных. Если возникает ошибка, callback будет вызван с этой ошибкой в качестве первого аргумента. Вызов callback выполняется асинхронно и до генерации 'error'.
Возвращаемое значение равно true, если после принятия chunk внутренний буфер меньше значения highWaterMark, заданного при создании потока. Если возвращено false, дальнейшую запись данных в поток следует прекратить до генерации события 'drain'.
Пока поток не освобождает место в буфере, вызовы write() будут буферизовать chunk и возвращать false. После передачи всех буферизованных на данный момент фрагментов (принятия их операционной системой к доставке) будет сгенерировано событие 'drain'. После того как write() вернёт false, не записывайте новые фрагменты до генерации события 'drain'. Хотя вызовы write() для потока, не освобождающего место в буфере, разрешены, Node.js будет буферизовать все записываемые фрагменты, пока не будет исчерпан максимум используемой памяти, после чего процесс будет безусловно прерван. Ещё до прерывания высокое потребление памяти ухудшает работу сборщика мусора и увеличивает RSS (который обычно не освобождается для системы, даже когда память больше не нужна). Поскольку TCP-сокеты могут никогда не освободить буфер, если удалённый узел не считывает данные, запись в такой сокет может привести к уязвимости, которой можно воспользоваться удалённо.
Запись данных, когда поток не освобождает место в буфере, особенно проблематична для потока Transform, поскольку потоки Transform по умолчанию приостановлены до тех пор, пока для них не начнётся передача данных или не будет добавлен обработчик события 'data' либо 'readable'.
Если данные для записи можно создавать или получать по требованию, рекомендуется инкапсулировать эту логику в поток Readable и использовать stream.pipe(). Однако если предпочтительно вызывать write(), можно соблюдать обратное давление и избегать проблем с памятью, используя событие 'drain':
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.');
}); copy Поток Writable в объектном режиме всегда игнорирует аргумент encoding.
Потоки для чтения
Потоки для чтения представляют собой абстракцию источника, из которого потребляются данные.
Примеры потоков Readable включают:
- HTTP-ответы на стороне клиента
- HTTP-запросы на стороне сервера
- потоки чтения fs
- потоки zlib
- потоки crypto
- TCP-сокеты
- stdout и stderr дочернего процесса
process.stdin
Все потоки Readable реализуют интерфейс, определенный классом stream.Readable.
Два режима чтения
Потоки Readable фактически работают в одном из двух режимов: активном и приостановленном. Эти режимы не связаны с режимом объектов. Поток Readable может находиться в режиме объектов или нет, независимо от того, работает ли он в активном или приостановленном режиме.
-
В активном режиме данные автоматически считываются из базовой системы и передаются приложению как можно быстрее с помощью событий через интерфейс
EventEmitter. -
В приостановленном режиме для чтения фрагментов данных из потока необходимо явно вызвать метод
stream.read().
Все потоки Readable начинают работу в приостановленном режиме, но могут быть переключены в активный режим одним из следующих способов:
- Добавить обработчик события
'data'. - Вызвать метод
stream.resume(). - Вызвать метод
stream.pipe(), чтобы передать данные потокуWritable.
Readable может вернуться в приостановленный режим одним из следующих способов:
- Если нет назначений для передачи данных, вызвав метод
stream.pause(). - Если есть назначения для передачи данных, удалив все такие назначения. Несколько назначений для передачи данных можно удалить вызовом метода
stream.unpipe().
Важно помнить, что Readable не будет генерировать данные, пока не будет предоставлен механизм для их потребления или игнорирования. Если механизм потребления отключен или удален, Readable попытается прекратить генерацию данных.
Из соображений обратной совместимости удаление обработчиков события 'data' не приостанавливает поток автоматически. Кроме того, если имеются назначения для передачи данных, вызов stream.pause() не гарантирует, что поток останется приостановленным после того, как эти назначения обработают данные и запросят новые.
Если поток Readable переключен в активный режим и нет потребителей, способных обработать данные, эти данные будут потеряны. Это может произойти, например, если вызвать метод readable.resume() без обработчика события 'data' или удалить из потока обработчик события 'data'.
Добавление обработчика события 'readable' автоматически останавливает поток, и данные необходимо потреблять с помощью readable.read(). Если обработчик события 'readable' удален, поток снова начнет работать в активном режиме, если у него есть обработчик события 'data'.
Три состояния
«Два режима» работы потока Readable — это упрощенная абстракция более сложного управления внутренним состоянием, происходящего в реализации потока Readable.
В частности, в каждый момент времени каждый Readable находится в одном из трех возможных состояний:
readable.readableFlowing === nullreadable.readableFlowing === falsereadable.readableFlowing === true
Когда readable.readableFlowing имеет значение null, механизм потребления данных потока не предоставлен. Поэтому поток не будет генерировать данные. В этом состоянии добавление обработчика события 'data', вызов метода readable.pipe() или вызов метода readable.resume() переключит readable.readableFlowing в состояние true, в результате чего Readable начнет активно генерировать события по мере поступления данных.
Вызов readable.pause(), readable.unpipe() или получение сигнала обратного давления приведет к тому, что readable.readableFlowing перейдет в состояние false, временно приостановив генерацию событий, но не генерацию данных. В этом состоянии добавление обработчика события 'data' не переключит readable.readableFlowing в состояние 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. copy Пока readable.readableFlowing имеет значение false, данные могут накапливаться во внутреннем буфере потока.
Выберите один способ работы с API
API потока Readable развивался в разных версиях Node.js и предоставляет несколько способов потребления данных потока. В целом разработчикам следует выбрать один способ потребления данных и не следует использовать несколько способов для потребления данных из одного потока. В частности, сочетание on('data'), on('readable'), pipe() или асинхронных итераторов может привести к неожиданному поведению.
Класс: stream.Readable
Событие: 'close'
Событие 'close' испускается, когда поток и все его базовые ресурсы (например, файловый дескриптор) закрыты. Это событие указывает, что больше не будет испускаться никаких событий и дальнейшие вычисления выполняться не будут.
Поток Readable всегда испускает событие 'close', если он создан с параметром emitClose.
Событие: 'data'
-
chunk<Buffer> | <string> | <any> Фрагмент данных. Для потоков, работающих не в объектном режиме, фрагментом будет строка илиBuffer. Для потоков в объектном режиме фрагментом может быть любое значение JavaScript, кромеnull.
Событие 'data' испускается каждый раз, когда поток передаёт потребителю владение фрагментом данных. Это может происходить при переключении потока в режим передачи данных вызовом readable.pipe(), readable.resume() или добавлением функции-обработчика события 'data'. Событие 'data' также испускается при каждом вызове метода readable.read(), если доступен фрагмент данных для возврата.
Добавление обработчика события 'data' к потоку, который явно не был приостановлен, переключит поток в режим передачи данных. Данные будут передаваться сразу после того, как станут доступны.
Функции-обработчику передаётся фрагмент данных в виде строки, если для потока была задана кодировка по умолчанию с помощью метода readable.setEncoding(); в противном случае данные будут переданы как Buffer.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
}); copy Событие: 'end'
Событие 'end' испускается, когда в потоке больше нет данных для чтения.
Событие 'end' не будет испущено, пока данные не будут прочитаны полностью. Для этого можно переключить поток в режим передачи данных или многократно вызывать stream.read(), пока все данные не будут прочитаны.
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.');
}); copy Событие: 'error'
- Тип: <Error>
Событие 'error' может быть испущено реализацией Readable в любой момент. Обычно это происходит, если базовый поток не может создать данные из-за внутреннего сбоя или если реализация потока пытается передать недопустимый фрагмент данных.
Функции-обработчику будет передан один объект Error.
Событие: 'pause'
Событие 'pause' испускается при вызове stream.pause(), если readableFlowing не равен false.
Событие: 'readable'
Событие 'readable' испускается, когда в потоке есть данные для чтения, объём которых не превышает настроенный предел буфера (state.highWaterMark). Фактически оно указывает на то, что в буфере потока появилась новая информация. Если в буфере есть данные, их можно получить вызовом stream.read(). Кроме того, событие 'readable' может испускаться и при достижении конца потока.
const readable = getReadableStreamSomehow();
readable.on('readable', function() {
// There is some data to read now.
let data;
while ((data = this.read()) !== null) {
console.log(data);
}
}); copy Если достигнут конец потока, вызов stream.read() вернёт null и вызовет событие 'end'. Это также верно, если данных для чтения не было вовсе. Например, в следующем примере файл foo.txt пуст:
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');
}); copy При запуске этого скрипта вывод будет следующим:
$ node test.js readable: null end copy
В некоторых случаях добавление обработчика события 'readable' приводит к чтению некоторого объёма данных во внутренний буфер.
В целом механизмы событий readable.pipe() и 'data' проще для понимания, чем событие 'readable'. Однако обработка 'readable' может повысить пропускную способность.
Если одновременно используются 'readable' и 'data', 'readable' имеет приоритет при управлении потоком, то есть 'data' будет испускаться только при вызове stream.read(). Свойство readableFlowing получит значение false. Если при удалении 'readable' останутся обработчики 'data', поток начнёт передавать данные, то есть события 'data' будут испускаться без вызова .resume().
Событие: 'resume'
Событие 'resume' испускается при вызове stream.resume(), если readableFlowing не равен true.
readable.destroy([error])
Уничтожает поток. При необходимости испускает событие 'error' и событие 'close' (если только emitClose не установлено в false). После этого вызова читаемый поток освобождает все внутренние ресурсы, а последующие вызовы push() игнорируются.
После вызова destroy() все дальнейшие вызовы ничего не делают; дальнейшие ошибки, кроме _destroy(), могут быть испущены как 'error'.
Разработчикам не следует переопределять этот метод; вместо этого следует реализовать readable._destroy().
readable.isPaused()
- Возвращает: <boolean>
Метод readable.isPaused() возвращает текущее рабочее состояние Readable. Он используется главным образом механизмом, лежащим в основе метода readable.pipe(). В большинстве обычных случаев нет необходимости вызывать этот метод напрямую.
const readable = new stream.Readable(); readable.isPaused(); // === false readable.pause(); readable.isPaused(); // === true readable.resume(); readable.isPaused(); // === false copy
readable.pause()
- Возвращает: <this>
Метод readable.pause() останавливает испускание событий 'data' потоком, работающим в режиме передачи данных, и переключает его из этого режима. Все поступающие данные остаются во внутреннем буфере.
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);
}); copy Метод readable.pause() не действует, если имеется обработчик события 'readable'.
readable.pipe(destination[, options])
-
destination<stream.Writable> Назначение для записи данных -
options<Object> Параметры канала-
end<boolean> Завершать записывающий поток после завершения читающего. По умолчанию:true.
-
- Возвращает: <stream.Writable> Назначение, позволяющее выстраивать цепочки каналов, если это поток
DuplexилиTransform
Метод readable.pipe() подключает поток Writable к readable, автоматически переводя его в режим передачи данных и перенаправляя все данные в подключённый Writable. Поток данных будет регулироваться автоматически, чтобы более быстрый поток Readable не перегружал поток назначения Writable.
В следующем примере все данные из readable направляются в файл с именем 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); copy К одному потоку Readable можно подключить несколько потоков Writable.
Метод readable.pipe() возвращает ссылку на поток назначения, позволяя создавать цепочки потоков, соединённых каналами:
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); copy По умолчанию для потока назначения Writable вызывается stream.end(), когда исходный поток Readable испускает событие 'end', после чего в назначение больше нельзя записывать данные. Чтобы отключить такое поведение по умолчанию и оставить поток назначения открытым, параметр end можно передать со значением false:
reader.pipe(writer, { end: false });
reader.on('end', () => {
writer.end('Goodbye\n');
}); copy Важно учитывать, что если во время обработки поток Readable испустит ошибку, поток назначения Writable не будет автоматически закрыт. При возникновении ошибки каждый поток необходимо вручную закрыть, чтобы предотвратить утечки памяти.
Потоки Writable process.stderr и process.stdout не закрываются до завершения процесса Node.js независимо от указанных параметров.
readable.read([size])
-
size<number> Необязательный аргумент, задающий объём данных для чтения. - Возвращает: <string> | <Buffer> | <null> | <any>
Метод readable.read() считывает данные из внутреннего буфера и возвращает их. Если данных для чтения нет, возвращается null. По умолчанию данные возвращаются как объект Buffer, если только с помощью метода readable.setEncoding() не была задана кодировка или поток не работает в объектном режиме.
Необязательный аргумент size задаёт точное количество байтов для чтения. Если для чтения доступно меньше size байт, будет возвращено null, если только поток не завершился; в этом случае будут возвращены все данные, оставшиеся во внутреннем буфере.
Если аргумент size не указан, будут возвращены все данные из внутреннего буфера.
Значение аргумента size должно быть не больше 1 ГиБ.
Метод readable.read() следует вызывать только для потоков Readable, работающих в режиме паузы. В режиме передачи данных readable.read() вызывается автоматически, пока внутренний буфер не будет полностью опустошён.
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.');
}); copy Каждый вызов readable.read() возвращает фрагмент данных или null, означающий, что в данный момент данных для чтения больше нет. Эти фрагменты не объединяются автоматически. Поскольку один вызов read() не возвращает все данные, может потребоваться цикл while, который будет непрерывно считывать фрагменты до получения всех данных. При чтении большого файла .read() может временно вернуть null, указывая, что все данные из буфера считаны, но в него могут поступить новые. В таких случаях событие 'readable' испускается, как только в буфере появляются новые данные, а событие 'end' сигнализирует о завершении передачи данных.
Поэтому, чтобы прочитать всё содержимое файла из readable, необходимо собирать фрагменты при нескольких событиях 'readable':
const chunks = [];
readable.on('readable', () => {
let chunk;
while (null !== (chunk = readable.read())) {
chunks.push(chunk);
}
});
readable.on('end', () => {
const content = chunks.join('');
}); copy Поток Readable в объектном режиме всегда возвращает один элемент при вызове readable.read(size) независимо от значения аргумента size.
Если метод readable.read() возвращает фрагмент данных, также будет испущено событие 'data'.
Вызов stream.read([size]) после испускания события 'end' вернёт null. Ошибка времени выполнения при этом не возникнет.
readable.readable
- Тип: <boolean>
Равно true, если безопасно вызвать readable.read(); это означает, что поток не был уничтожен и не испустил события 'error' или 'end'.
readable.readableAborted
- Тип: <boolean>
Возвращает, был ли поток уничтожен или в нём произошла ошибка до испускания 'end'.
readable.readableEncoding
Геттер свойства encoding указанного потока Readable. Свойство encoding можно задать с помощью метода readable.setEncoding().
readable.readableFlowing
- Тип: <boolean>
Это свойство отражает текущее состояние потока Readable, описанное в разделе Три состояния.
readable.readableHighWaterMark
- Тип: <number>
Возвращает значение highWaterMark, переданное при создании этого Readable.
readable.readableLength
- Тип: <number>
Это свойство содержит количество байтов (или объектов) в очереди, готовых к чтению. Значение позволяет получить сведения о состоянии highWaterMark.
readable.resume()
- Возвращает: <this>
Метод readable.resume() возобновляет испускание событий 'data' явно приостановленным потоком Readable и переключает его в режим передачи данных.
Метод readable.resume() можно использовать для полного чтения данных из потока без их фактической обработки:
getReadableStreamSomehow()
.resume()
.on('end', () => {
console.log('Reached the end, but did not read anything.');
}); copy Метод readable.resume() не действует, если имеется обработчик события 'readable'.
readable.setEncoding(encoding)
Метод readable.setEncoding() задаёт кодировку символов для данных, считываемых из потока Readable.
По умолчанию кодировка не задана, и данные потока возвращаются как объекты Buffer. Если задать кодировку, данные потока будут возвращаться в виде строк в указанной кодировке, а не объектов Buffer. Например, вызов readable.setEncoding('utf8') приведёт к интерпретации выходных данных как данных UTF-8 и передаче их в виде строк. Вызов readable.setEncoding('hex') приведёт к кодированию данных в виде шестнадцатеричной строки.
Поток Readable корректно обрабатывает многобайтовые символы, передаваемые через поток, которые в противном случае были бы декодированы некорректно при считывании из потока в виде объектов Buffer.
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);
}); copy
readable.unpipe([destination])
-
destination<stream.Writable> Необязательный поток, который требуется отсоединить от канала - Возвращает: <this>
Метод readable.unpipe() отсоединяет поток Writable, ранее подключённый с помощью метода stream.pipe().
Если параметр destination не указан, отсоединяются все каналы.
Если параметр destination указан, но для него не настроен канал, метод ничего не делает.
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); copy
readable.unshift(chunk[, encoding])
-
chunk<Buffer> | <TypedArray> | <DataView> | <string> | <null> | <any> Фрагмент данных, добавляемый в начало очереди чтения. Для потоков, работающих не в объектном режиме,chunkдолжен быть значением типа <string>, <Buffer>, <TypedArray>, <DataView> илиnull. Для потоков в объектном режимеchunkможет быть любым значением JavaScript. -
encoding<string> Кодировка фрагментов-строк. Должна быть допустимой кодировкойBuffer, например'utf8'или'ascii'.
Передача chunk в качестве null сигнализирует о конце потока (EOF) и действует так же, как readable.push(null); после этого записывать данные больше нельзя. Сигнал EOF помещается в конец буфера, а все буферизованные данные будут сброшены.
Метод readable.unshift() помещает фрагмент данных обратно во внутренний буфер. Это полезно в некоторых ситуациях, когда поток считывается кодом, которому нужно «отменить чтение» части данных, предварительно полученных из источника, чтобы передать их другому потребителю.
Метод stream.unshift(chunk) нельзя вызывать после испускания события 'end', иначе будет выброшена ошибка времени выполнения.
Разработчикам, использующим stream.unshift(), часто следует рассмотреть возможность перехода на поток Transform. Дополнительные сведения см. в разделе API для разработчиков потоков.
// 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;
}
}
} copy В отличие от stream.push(chunk), stream.unshift(chunk) не завершает процесс чтения сбросом внутреннего состояния чтения потока. Это может привести к неожиданным результатам, если readable.unshift() вызван во время чтения (например, из реализации stream._read() пользовательского потока). Немедленный вызов stream.push('') после readable.unshift() правильно сбросит состояние чтения, однако лучше вообще не вызывать readable.unshift() во время чтения.
readable.wrap(stream)
До Node.js 0.10 потоки не реализовывали весь API модуля node:stream в его нынешнем виде. (Дополнительные сведения см. в разделе Совместимость.)
При использовании старой библиотеки Node.js, испускающей события 'data' и имеющей лишь рекомендательный метод stream.pause(), метод readable.wrap() позволяет создать поток Readable, использующий старый поток в качестве источника данных.
Метод readable.wrap() требуется редко, но он предоставлен для удобства взаимодействия со старыми приложениями и библиотеками Node.js.
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.
}); copy
readable[Symbol.asyncIterator]()
- Возвращает: <AsyncIterator> для полного чтения потока.
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); copy Если цикл завершается с помощью break, return или throw, поток будет уничтожен. Иными словами, итерация по потоку приводит к его полному чтению. Поток будет считываться фрагментами размером, равным параметру highWaterMark. В приведённом выше примере кода данные будут считаны одним фрагментом, если размер файла меньше 64 КиБ, поскольку для fs.createReadStream() не указан параметр highWaterMark.
readable[Symbol.asyncDispose]()
Вызывает readable.destroy() с AbortError и возвращает промис, который выполняется после завершения потока.
readable.compose(stream[, options])
-
stream<Stream> | <Iterable> | <AsyncIterable> | <Function> -
options<Object>-
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Duplex> поток, объединённый с потоком
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(['this is', 'compose as operator']).compose(splitToWords);
const words = await wordsStream.toArray();
console.log(words); // prints ['this', 'is', 'compose', 'as', 'operator'] copy Дополнительные сведения см. в разделе stream.compose.
readable.iterator([options])
-
options<Object>-
destroyOnReturn<boolean> Если задано значениеfalse, вызовreturnдля асинхронного итератора или выход из итерацииfor await...ofс помощьюbreak,returnилиthrowне приведёт к уничтожению потока. По умолчанию:true.
-
- Возвращает: <AsyncIterator> для чтения потока.
Итератор, созданный этим методом, позволяет отменить уничтожение потока, если цикл for await...of завершается с помощью return, break или throw, либо указать, что итератор должен уничтожить поток, если во время итерации в потоке возникла ошибка.
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(); copy
readable.map(fn[, options])
-
fn<Function> | <AsyncFunction> функция, применяемая к каждому фрагменту потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
highWaterMark<number> количество элементов, буферизуемых в ожидании обработки сопоставленных элементов пользователем. По умолчанию:concurrency * 2 - 1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Readable> поток, преобразованный функцией
fn.
Этот метод позволяет преобразовать поток. Функция fn вызывается для каждого фрагмента потока. Если функция fn возвращает промис, этот промис будет обработан с помощью await перед передачей результата в выходной поток.
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.
} copy
readable.filter(fn[, options])
-
fn<Function> | <AsyncFunction> функция для фильтрации фрагментов потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
highWaterMark<number> количество элементов, буферизуемых в ожидании обработки отфильтрованных элементов пользователем. По умолчанию:concurrency * 2 - 1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Readable> поток, отфильтрованный предикатом
fn.
Этот метод позволяет фильтровать поток. Для каждого фрагмента потока вызывается функция fn; если она возвращает истинное значение, фрагмент передаётся в выходной поток. Если функция fn возвращает промис, этот промис будет обработан с помощью await.
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);
} copy
readable.forEach(fn[, options])
-
fn<Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Promise> промис, выполняющийся после завершения потока.
Этот метод позволяет перебирать поток. Для каждого фрагмента потока вызывается функция fn. Если функция fn возвращает промис, этот промис будет обработан с помощью await.
Этот метод отличается от циклов for await...of тем, что позволяет при необходимости обрабатывать фрагменты параллельно. Кроме того, итерацию forEach можно остановить только посредством передачи параметра signal и прерывания связанного с ним AbortController, тогда как for await...of можно остановить с помощью break или return. В обоих случаях поток будет уничтожен.
Этот метод отличается от прослушивания события 'data' тем, что во внутренних механизмах он использует событие readable и позволяет ограничить количество одновременных вызовов fn.
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 copy
readable.toArray([options])
-
options<Object>-
signal<AbortSignal> позволяет отменить операцию toArray при прерывании сигнала.
-
- Возвращает: <Promise> промис, содержащий массив с содержимым потока.
Этот метод позволяет легко получить содержимое потока.
Поскольку этот метод считывает весь поток в память, он сводит на нет преимущества потоков. Он предназначен для совместимости и удобства, а не для основного способа чтения потоков.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
await Readable.from([1, 2, 3, 4]).toArray(); // [1, 2, 3, 4]
// 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(); copy
readable.some(fn[, options])
-
fn<Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Promise> промис, результатом которого будет
true, еслиfnвернуло истинное значение хотя бы для одного из фрагментов.
Этот метод похож на Array.prototype.some и вызывает fn для каждого фрагмента потока, пока ожидаемое возвращаемое значение не станет true (или любым другим истинным значением). Как только ожидаемое возвращаемое значение вызова fn для фрагмента становится истинным, поток уничтожается, а промис выполняется со значением true. Если ни один из вызовов fn для фрагментов не возвращает истинное значение, промис выполняется со значением false.
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 copy
readable.find(fn[, options])
-
fn<Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Promise> промис, результатом которого будет первый фрагмент, для которого
fnвернуло истинное значение, либоundefined, если элемент не найден.
Этот метод похож на Array.prototype.find и вызывает fn для каждого фрагмента потока, чтобы найти фрагмент, для которого fn имеет истинное значение. Как только ожидаемое возвращаемое значение вызова fn становится истинным, поток уничтожается, а промис выполняется со значением, для которого fn вернуло истинное значение. Если все вызовы fn для фрагментов возвращают ложное значение, промис выполняется со значением undefined.
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 copy
readable.every(fn[, options])
-
fn<Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Promise> промис, результатом которого будет
true, еслиfnвернуло истинное значение для всех фрагментов.
Этот метод похож на Array.prototype.every и вызывает fn для каждого фрагмента потока, чтобы проверить, являются ли все ожидаемые возвращаемые значения истинными для fn. Как только ожидаемое возвращаемое значение вызова fn для фрагмента становится ложным, поток уничтожается, а промис выполняется со значением false. Если все вызовы fn для фрагментов возвращают истинное значение, промис выполняется со значением true.
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 copy
readable.flatMap(fn[, options])
-
fn<Function> | <AsyncGeneratorFunction> | <AsyncFunction> функция, применяемая к каждому фрагменту потока.-
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
options<Object>-
concurrency<number> максимальное число одновременных вызововfnдля обработки потока. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Readable> поток, преобразованный функцией
fnс разворачиванием результата.
Этот метод возвращает новый поток, применяя заданную функцию обратного вызова к каждому фрагменту исходного потока и затем разворачивая результат.
Функция fn может возвращать поток, другой итерируемый объект или асинхронный итерируемый объект; результирующие потоки будут объединены (развёрнуты) в возвращаемый поток.
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);
} copy
readable.drop(limit[, options])
-
limit<number> количество фрагментов, отбрасываемых из читаемого потока. -
options<Object>-
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Readable> поток, из которого удалены
limitфрагментов.
Этот метод возвращает новый поток, отбросив первые limit фрагментов.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).drop(2).toArray(); // [3, 4] copy
readable.take(limit[, options])
-
limit<number> количество фрагментов, извлекаемых из читаемого потока. -
options<Object>-
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Readable> поток, содержащий
limitизвлечённых фрагментов.
Этот метод возвращает новый поток с первыми limit фрагментами.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).take(2).toArray(); // [1, 2] copy
readable.reduce(fn[, initial[, options]])
-
fn<Function> | <AsyncFunction> функция-редуктор, вызываемая для каждого фрагмента потока.-
previous<any> значение, полученное при последнем вызовеfn, или значениеinitial, если оно указано, либо первый фрагмент потока. -
data<any> фрагмент данных из потока. -
options<Object>-
signal<AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызовfn.
-
-
-
initial<any> начальное значение для редукции. -
options<Object>-
signal<AbortSignal> позволяет уничтожить поток при прерывании сигнала.
-
- Возвращает: <Promise> промис с итоговым значением редукции.
Этот метод последовательно вызывает fn для каждого фрагмента потока, передавая результат вычисления для предыдущего элемента. Он возвращает промис с итоговым значением редукции.
Если начальное значение initial не задано, в качестве начального значения используется первый фрагмент потока. Если поток пуст, промис отклоняется с ошибкой TypeError, у которой свойство кода имеет значение ERR_INVALID_ARGS.
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); copy Функция-редуктор перебирает элементы потока по одному, поэтому параметр concurrency и параллельная обработка отсутствуют. Чтобы выполнить reduce параллельно, можно вынести асинхронную функцию в метод readable.map.
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); copy Двунаправленные потоки и потоки преобразования
Класс: stream.Duplex
Двунаправленные потоки — это потоки, реализующие интерфейсы Readable и Writable.
Примеры потоков Duplex:
duplex.allowHalfOpen
- Тип: <boolean>
Если false, то поток автоматически завершит записываемую сторону при завершении читаемой стороны. Изначально значение задаётся параметром конструктора allowHalfOpen, по умолчанию равным true.
Это значение можно изменить вручную, чтобы изменить поведение полузакрытого режима существующего экземпляра потока Duplex, но сделать это необходимо до вызова события 'end'.
Класс: stream.Transform
Потоки преобразования — это потоки Duplex, в которых выходные данные некоторым образом связаны с входными. Как и все потоки Duplex, потоки Transform реализуют интерфейсы Readable и Writable.
Примеры потоков Transform:
transform.destroy([error])
Уничтожает поток и при необходимости генерирует событие 'error'. После этого вызова поток освободит все внутренние ресурсы. Реализующим классам не следует переопределять этот метод, вместо этого нужно реализовать readable._destroy(). Реализация _destroy() по умолчанию для Transform также генерирует 'close', если для emitClose не задано значение false.
После вызова destroy() все последующие вызовы ничего не делают, и дальнейшие ошибки, кроме _destroy(), могут быть сгенерированы как 'error'.
stream.duplexPair([options])
-
options<Object> Значение, передаваемое обоим конструкторамDuplexдля задания таких параметров, как буферизация. - Возвращает: <Array> из двух экземпляров
Duplex.
Вспомогательная функция duplexPair возвращает массив из двух элементов, каждый из которых является потоком Duplex, соединённым с другой стороной:
const [ sideA, sideB ] = duplexPair(); copy
Всё, что записывается в один поток, становится доступным для чтения из другого. Это поведение аналогично сетевому соединению, в котором данные, записанные клиентом, становятся доступными для чтения сервером, и наоборот.
Двунаправленные потоки симметричны: любой из них можно использовать без различий в поведении.
stream.finished(stream[, options], callback)
-
stream<Stream> | <ReadableStream> | <WritableStream> Читаемый и/или записываемый поток/веб-поток. -
options<Object>-
error<boolean> Если задано значениеfalse, вызовemit('error', err)не считается завершением. По умолчанию:true. -
readable<boolean> Если задано значениеfalse, обратный вызов будет вызван при завершении потока, даже если поток всё ещё может быть доступен для чтения. По умолчанию:true. -
writable<boolean> Если задано значениеfalse, обратный вызов будет вызван при завершении потока, даже если в поток всё ещё можно записывать. По умолчанию:true. -
signal<AbortSignal> позволяет прервать ожидание завершения потока. Базовый поток не будет прерван, если сигнал будет прерван. Обратный вызов будет вызван с ошибкойAbortError. Все зарегистрированные этой функцией обработчики также будут удалены.
-
-
callback<Function> Функция обратного вызова, принимающая необязательный аргумент ошибки. - Возвращает: <Function> Функция очистки, удаляющая все зарегистрированные обработчики.
Функция для получения уведомления о том, что поток больше не доступен для чтения или записи, либо в нём произошла ошибка или преждевременное закрытие.
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. copy Особенно полезна при обработке ошибок, когда поток уничтожается преждевременно (например, при прерванном HTTP-запросе) и не генерирует 'end' или 'finish'.
API finished предоставляет версию на основе промисов.
stream.finished() оставляет необработанные слушатели событий (в частности, 'error', 'end', 'finish' и 'close') после вызова callback. Это сделано для того, чтобы неожиданные события 'error' (из-за некорректной реализации потока) не приводили к неожиданным сбоям. Если такое поведение нежелательно, в обратном вызове необходимо вызвать возвращённую функцию очистки:
const cleanup = finished(rs, (err) => {
cleanup();
// ...
}); copy
stream.pipeline(source[, ...transforms], destination, callback)
stream.pipeline(streams, callback)
-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> -
source<Stream> | <Iterable> | <AsyncIterable> | <Function> | <ReadableStream>- Возвращает: <Iterable> | <AsyncIterable>
-
...transforms<Stream> | <Function> | <TransformStream>-
source<AsyncIterable> - Возвращает: <AsyncIterable>
-
-
destination<Stream> | <Function> | <WritableStream>-
source<AsyncIterable> - Возвращает: <AsyncIterable> | <Promise>
-
-
callback<Function> Вызывается после полного завершения конвейера.-
err<Error> -
valРазрешённое значениеPromise, возвращённого функциейdestination.
-
- Возвращает: <Stream>
Метод модуля для передачи данных между потоками и генераторами с перенаправлением ошибок, корректной очисткой и вызовом обратного вызова после завершения конвейера.
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.');
}
},
); copy API pipeline предоставляет версию на основе промисов.
stream.pipeline() вызовет stream.destroy(err) для всех потоков, кроме:
-
Readableпотоков, которые сгенерировали'end'или'close'. -
Writableпотоков, которые сгенерировали'finish'или'close'.
stream.pipeline() оставляет необработанные слушатели событий в потоках после вызова callback. При повторном использовании потоков после сбоя это может привести к утечке слушателей событий и подавлению ошибок. Если последний поток доступен для чтения, необработанные слушатели событий будут удалены, чтобы этот поток можно было использовать позднее.
stream.pipeline() закрывает все потоки при возникновении ошибки. Использование IncomingRequest с pipeline может привести к неожиданному поведению: сокет будет уничтожен без отправки ожидаемого ответа. См. пример ниже:
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!!!');
}
});
}); copy
stream.compose(...streams)
stream.compose является экспериментальным.-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> | <Duplex[]> | <Function> - Возвращает: <stream.Duplex>
Объединяет два или более потоков в поток Duplex, который записывает данные в первый поток и читает их из последнего. Каждый переданный поток соединяется со следующим с помощью stream.pipeline. При возникновении ошибки в любом из потоков уничтожаются все потоки, включая внешний поток Duplex.
Поскольку stream.compose возвращает новый поток, который, в свою очередь, можно (и следует) передавать в другие потоки, эта функция позволяет создавать композиции. В отличие от этого, при передаче потоков в stream.pipeline первый поток обычно является читаемым, а последний — записываемым, образуя замкнутую цепь.
Если передан Function, он должен быть фабричным методом, принимающим source Iterable.
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' copy stream.compose можно использовать для преобразования асинхронных итерируемых объектов, генераторов и функций в потоки.
-
AsyncIterableпреобразует объект в читаемыйDuplex. Не может выдаватьnull. -
AsyncGeneratorFunctionпреобразует объект в читаемый/записываемый поток преобразованияDuplex. Первым параметром должен принимать исходныйAsyncIterable. Не может выдаватьnull. -
AsyncFunctionпреобразует объект в записываемыйDuplex. Должен возвращать либоnull, либоundefined.
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' copy См. readable.compose(stream), чтобы использовать stream.compose в качестве оператора.
stream.isErrored(stream)
-
stream<Readable> | <Writable> | <Duplex> | <WritableStream> | <ReadableStream> - Возвращает: <boolean>
Возвращает значение, указывающее, произошла ли в потоке ошибка.
stream.isReadable(stream)
-
stream<Readable> | <Duplex> | <ReadableStream> - Возвращает: <boolean> | <null> — возвращает
nullтолько в том случае, еслиstreamне является допустимымReadable,DuplexилиReadableStream.
Возвращает значение, указывающее, доступен ли поток для чтения.
stream.isWritable(stream)
-
stream<Writable> | <Duplex> | <WritableStream> - Возвращает: <boolean> | <null> — возвращает
nullтолько в том случае, еслиstreamне является допустимымWritable,DuplexилиWritableStream.
Возвращает значение, указывающее, доступен ли поток для записи.
stream.Readable.from(iterable[, options])
-
iterable<Iterable> Объект, реализующий протокол перебораSymbol.asyncIteratorилиSymbol.iterator. Генерирует событие 'error', если передано значение null. -
options<Object> Параметры, передаваемые вnew stream.Readable([options]). По умолчаниюReadable.from()устанавливаетoptions.objectModeв значениеtrue, если это явно не отменено установкойoptions.objectModeв значениеfalse. - Возвращает: <stream.Readable>
Вспомогательный метод для создания читаемых потоков из итераторов.
const { Readable } = require('node:stream');
async function * generate() {
yield 'hello';
yield 'streams';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
}); copy В целях повышения производительности вызов Readable.from(string) или Readable.from(buffer) не приводит к перебору строк или буферов в соответствии с семантикой других потоков.
Если в качестве аргумента передан объект Iterable, содержащий промисы, это может привести к необработанному отклонению промиса.
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
]); copy
stream.Readable.fromWeb(readableStream[, options])
-
readableStream<ReadableStream> -
options<Object>-
encoding<string> -
highWaterMark<number> -
objectMode<boolean> -
signal<AbortSignal>
-
- Возвращает: <stream.Readable>
stream.Readable.isDisturbed(stream)
-
stream<stream.Readable> | <ReadableStream> - Возвращает:
boolean
Возвращает значение, указывающее, считывались ли данные из потока или был ли он отменён.
stream.Readable.toWeb(streamReadable[, options])
-
streamReadable<stream.Readable> -
options<Object>-
strategy<Object>-
highWaterMark<number> Максимальный размер внутренней очереди (созданногоReadableStream), после которого при чтении из указанногоstream.Readableначинает применяться обратное давление. Если значение не задано, оно будет взято из указанногоstream.Readable. -
size<Function> Функция, вычисляющая размер указанного фрагмента данных. Если значение не задано, размер всех фрагментов будет1.
-
-
- Возвращает: <ReadableStream>
stream.Writable.fromWeb(writableStream[, options])
-
writableStream<WritableStream> -
options<Object>-
decodeStrings<boolean> -
highWaterMark<number> -
objectMode<boolean> -
signal<AbortSignal>
-
- Возвращает: <stream.Writable>
stream.Writable.toWeb(streamWritable)
-
streamWritable<stream.Writable> - Возвращает: <WritableStream>
stream.Duplex.from(src)
-
src<Stream> | <Blob> | <ArrayBuffer> | <string> | <Iterable> | <AsyncIterable> | <AsyncGeneratorFunction> | <AsyncFunction> | <Promise> | <Object> | <ReadableStream> | <WritableStream>
Вспомогательный метод для создания двунаправленных потоков.
-
Streamпреобразует поток для записи в записываемыйDuplex, а поток для чтения — вDuplex. -
Blobпреобразует в читаемыйDuplex. -
stringпреобразует в читаемыйDuplex. -
ArrayBufferпреобразует в читаемыйDuplex. -
AsyncIterableпреобразует в читаемыйDuplex. Не может выдаватьnull. -
AsyncGeneratorFunctionпреобразует в преобразующийDuplexдля чтения и записи. Первым параметром должен принимать исходныйAsyncIterable. Не может выдаватьnull. -
AsyncFunctionпреобразует в записываемыйDuplex. Должен возвращать либоnull, либоundefined -
Object ({ writable, readable })преобразуетreadableиwritableвStream, а затем объединяет их вDuplex, гдеDuplexбудет записывать вwritableи читать изreadable. -
Promiseпреобразует в читаемыйDuplex. Значениеnullигнорируется. -
ReadableStreamпреобразует в читаемыйDuplex. -
WritableStreamпреобразует в записываемыйDuplex. - Возвращает: <stream.Duplex>
Если в качестве аргумента передан объект Iterable, содержащий промисы, это может привести к необработанному отклонению промиса.
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
]); copy
stream.Duplex.fromWeb(pair[, options])
-
pair<Object>-
readable<ReadableStream> -
writable<WritableStream>
-
-
options<Object> - Возвращает: <stream.Duplex>
Модули JavaScript
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);
}CommonJS
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)
-
streamDuplex<stream.Duplex> - Возвращает: <Object>
-
readable<ReadableStream> -
writable<WritableStream>
-
Модули JavaScript
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);CommonJS
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> Сигнал, обозначающий возможность отмены -
stream<Stream> | <ReadableStream> | <WritableStream> Поток, к которому нужно привязать сигнал.
Привязывает AbortSignal к потоку для чтения или записи. Это позволяет коду управлять уничтожением потока с помощью AbortController.
Вызов abort для AbortController, соответствующего переданному AbortSignal, будет работать так же, как вызов .destroy(new AbortError()) для потока, а для веб-потоков — controller.error(new AbortError()).
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(); copy Или используя AbortSignal с читаемым потоком в качестве асинхронного итерируемого объекта:
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;
}
}
})(); copy Или используя AbortSignal с 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();
}); copy
stream.getDefaultHighWaterMark(objectMode)
Возвращает значение highWaterMark по умолчанию, используемое потоками. По умолчанию это 65536 (64 КиБ) или 16 для objectMode.
API для разработчиков потоков
API модуля node:stream разработан так, чтобы потоки можно было легко реализовывать с помощью модели прототипного наследования JavaScript.
Сначала разработчик потока объявляет новый класс JavaScript, расширяющий один из четырёх основных классов потоков (stream.Writable, stream.Readable, stream.Duplex или stream.Transform), и вызывает конструктор соответствующего родительского класса:
const { Writable } = require('node:stream');
class MyWritable extends Writable {
constructor({ highWaterMark, ...options }) {
super({ highWaterMark });
// ...
}
} copy При расширении потоков учитывайте, какие параметры пользователь может и должен передавать, прежде чем передавать их базовому конструктору. Например, если реализация предполагает определённые значения параметров autoDestroy и emitClose, не позволяйте пользователю переопределять их. Явно указывайте, какие параметры передаются, вместо неявной передачи всех параметров.
Затем новый класс потока должен реализовать один или несколько определённых методов в зависимости от создаваемого типа потока, как показано в таблице ниже:
| Сценарий использования | Класс | Методы для реализации |
|---|---|---|
| Только чтение | Readable |
_read() |
| Только запись | Writable |
_write(), _writev(), _final()
|
| Чтение и запись | Duplex |
_read(), _write(), _writev(), _final()
|
| Обработка записанных данных с последующим чтением результата | Transform |
_transform(), _flush(), _final()
|
Код реализации потока никогда не должен вызывать «публичные» методы потока, предназначенные для потребителей (как описано в разделе API для потребителей потоков). Это может привести к нежелательным побочным эффектам в коде приложения, использующего поток.
Не переопределяйте публичные методы, такие как write(), end(), cork(), uncork(), read() и destroy(), и не генерируйте внутренние события, такие как 'error', 'data', 'end', 'finish' и 'close' через .emit(). Это может нарушить текущие и будущие инварианты потоков, что приведёт к проблемам с поведением и/или совместимостью с другими потоками, утилитами потоков и ожиданиями пользователей.
Упрощённое создание
Во многих простых случаях поток можно создать без наследования. Для этого достаточно напрямую создать экземпляры объектов stream.Writable, stream.Readable, stream.Duplex или stream.Transform и передать соответствующие методы в качестве параметров конструктора.
const { Writable } = require('node:stream');
const myWritable = new Writable({
construct(callback) {
// Initialize state and load resources...
},
write(chunk, encoding, callback) {
// ...
},
destroy() {
// Free resources...
},
}); copy Реализация записываемого потока
Класс stream.Writable расширяется для реализации потока Writable.
Пользовательские потоки Writable должны вызывать конструктор new stream.Writable([options]) и реализовывать метод writable._write() и/или writable._writev().
new stream.Writable([options])
-
options<Object>-
highWaterMark<number> Размер буфера, при которомstream.write()начинает возвращатьfalse. По умолчанию:65536(64 КиБ) или16для потоковobjectMode. -
decodeStrings<boolean> Определяет, следует ли кодировать переданные вstream.write()значенияstringв значенияBuffer(с кодировкой, указанной при вызовеstream.write()) перед передачей вstream._write(). Другие типы данных не преобразуются (то есть значенияBufferне декодируются в значенияstring). Значение false предотвращает преобразование значенийstring. По умолчанию:true. -
defaultEncoding<string> Кодировка по умолчанию, используемая, если при вызовеstream.write()не указана кодировка. По умолчанию:'utf8'. -
objectMode<boolean> Определяет, является лиstream.write(anyObj)допустимой операцией. Если параметр задан, становится возможным записывать значения JavaScript, отличные от строк, <Buffer>, <TypedArray> или <DataView>, если это поддерживается реализацией потока. По умолчанию:false. -
emitClose<boolean> Определяет, должен ли поток генерировать'close'после уничтожения. По умолчанию:true. -
write<Function> Реализация методаstream._write(). -
writev<Function> Реализация методаstream._writev(). -
destroy<Function> Реализация методаstream._destroy(). -
final<Function> Реализация методаstream._final(). -
construct<Function> Реализация методаstream._construct(). -
autoDestroy<boolean> Определяет, должен ли поток автоматически вызывать.destroy()после завершения. По умолчанию:true. -
signal<AbortSignal> Сигнал, обозначающий возможность отмены.
-
const { Writable } = require('node:stream');
class MyWritable extends Writable {
constructor(options) {
// Calls the stream.Writable() constructor.
super(options);
// ...
}
} copy Или при использовании конструкторов в стиле до ES6:
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); copy Или с использованием упрощённого конструктора:
const { Writable } = require('node:stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
},
}); copy Вызов abort для AbortController, соответствующего переданному AbortSignal, будет работать так же, как вызов .destroy(new AbortError()) для записываемого потока.
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(); copy
writable._construct(callback)
-
callback<Function> Вызовите эту функцию (необязательно с аргументом ошибки), когда инициализация потока завершится.
Метод _construct() НЕЛЬЗЯ вызывать напрямую. Его могут реализовать дочерние классы; в таком случае он будет вызываться только внутренними методами класса Writable.
Эта необязательная функция будет вызвана в следующем тике после возврата конструктора потока, откладывая вызовы _write(), _final() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной подготовки ресурсов до начала использования потока.
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);
}
}
} copy
writable._write(chunk, encoding, callback)
-
chunk<Buffer> | <string> | <any> ЗаписываемыйBuffer, преобразованный изstring, переданного вstream.write(). Если параметрdecodeStringsпотока имеет значениеfalseили поток работает в режиме объектов, фрагмент не преобразуется и будет иметь исходное значение, переданное вstream.write(). -
encoding<string> Если фрагмент является строкой, тоencoding— это кодировка символов этой строки. Если фрагмент являетсяBufferили поток работает в режиме объектов,encodingможно игнорировать. -
callback<Function> Вызовите эту функцию (необязательно с аргументом ошибки), когда обработка переданного фрагмента завершится.
Все реализации потоков Writable должны предоставлять метод writable._write() и/или writable._writev() для отправки данных в нижележащий ресурс.
Потоки Transform предоставляют собственную реализацию writable._write().
Эту функцию НЕЛЬЗЯ вызывать напрямую из кода приложения. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Writable.
Функция callback должна быть вызвана синхронно внутри writable._write() или асинхронно (то есть в другом тике), чтобы сообщить об успешном завершении записи или об ошибке. Первым аргументом, передаваемым в callback, должен быть объект Error, если вызов завершился ошибкой, или null, если запись прошла успешно.
Все вызовы writable.write(), происходящие между вызовом writable._write() и вызовом callback, приводят к буферизации записываемых данных. При вызове callback поток может сгенерировать событие 'drain'. Если реализация потока может обрабатывать несколько фрагментов данных одновременно, следует реализовать метод writable._writev().
Если свойству decodeStrings явно присвоено значение false в параметрах конструктора, то chunk останется тем же объектом, который передан в .write(), и может быть строкой, а не Buffer. Это необходимо для поддержки реализаций, оптимизированно обрабатывающих определённые кодировки строковых данных. В этом случае аргумент encoding указывает кодировку символов строки. В противном случае аргумент encoding можно безопасно игнорировать.
Перед именем метода writable._write() стоит символ подчёркивания, поскольку он является внутренним для определяющего его класса, и пользовательские программы не должны вызывать его напрямую.
writable._writev(chunks, callback)
-
chunks<Object[]> Записываемые данные. Это массив объектов <Object>, каждый из которых представляет отдельный фрагмент данных для записи. Свойства этих объектов:-
chunk<Buffer> | <string> Экземпляр буфера или строка, содержащая записываемые данные. Значениеchunkбудет строкой, еслиWritableсоздан с параметромdecodeStrings, равнымfalse, и вwrite()передана строка. -
encoding<string> Кодировка символовchunk. ЕслиchunkявляетсяBuffer, значениеencodingбудет'buffer'.
-
-
callback<Function> Функция обратного вызова (необязательно с аргументом ошибки), которая будет вызвана после завершения обработки переданных фрагментов.
Эту функцию НЕЛЬЗЯ вызывать напрямую из кода приложения. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Writable.
Метод writable._writev() может быть реализован в дополнение к writable._write() или вместо него в реализациях потоков, способных обрабатывать несколько фрагментов данных одновременно. Если этот метод реализован и имеются данные, буферизованные при предыдущих операциях записи, вместо _write() будет вызван _writev().
Перед именем метода writable._writev() стоит символ подчёркивания, поскольку он является внутренним для определяющего его класса, и пользовательские программы не должны вызывать его напрямую.
writable._destroy(err, callback)
-
err<Error> Возможная ошибка. -
callback<Function> Функция обратного вызова, принимающая необязательный аргумент ошибки.
Метод _destroy() вызывается методом writable.destroy(). Его можно переопределить в дочерних классах, но его нельзя вызывать напрямую.
writable._final(callback)
-
callback<Function> Вызовите эту функцию (необязательно с аргументом ошибки), когда запись оставшихся данных будет завершена.
Метод _final() нельзя вызывать напрямую. Его могут реализовать дочерние классы; в таком случае он будет вызываться только внутренними методами класса Writable.
Эта необязательная функция будет вызвана перед закрытием потока, откладывая событие 'finish' до вызова callback. Это полезно для закрытия ресурсов или записи буферизованных данных перед завершением потока.
Ошибки при записи
Об ошибках, возникающих при выполнении методов writable._write(), writable._writev() и writable._final(), необходимо сообщать, вызывая функцию обратного вызова и передавая ошибку в качестве первого аргумента. Выбрасывание Error внутри этих методов или ручная генерация события 'error' приводят к неопределённому поведению.
Если поток Readable передаёт данные в поток Writable, то при возникновении ошибки в Writable поток Readable будет отключён от конвейера.
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();
}
},
}); copy Пример записываемого потока
Ниже показана довольно простая (и в некоторой степени бесполезная) реализация пользовательского потока Writable. Хотя этот конкретный экземпляр потока Writable не имеет особой практической пользы, пример демонстрирует все необходимые элементы пользовательского экземпляра потока Writable:
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();
}
}
} copy Декодирование буферов в записываемом потоке
Декодирование буферов — распространённая задача, например при использовании преобразователей со строковым входом. При работе с многобайтными кодировками символов, такими как UTF-8, этот процесс не так прост. В следующем примере показано, как декодировать многобайтные строки с помощью StringDecoder и Writable.
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: € copy Реализация потока для чтения
Класс stream.Readable расширяется для реализации потока Readable.
Пользовательские потоки Readable должны вызывать конструктор new stream.Readable([options]) и реализовывать метод readable._read().
new stream.Readable([options])
-
options<Object>-
highWaterMark<number> Максимальное количество байтов, которое можно хранить во внутреннем буфере до прекращения чтения из нижележащего ресурса. По умолчанию:65536(64 КиБ) или16для потоковobjectMode. -
encoding<string> Если задано, буферы будут декодироваться в строки с использованием указанной кодировки. По умолчанию:null. -
objectMode<boolean> Должен ли этот поток работать как поток объектов. Это означает, чтоstream.read(n)возвращает одно значение, а неBufferразмеромn. По умолчанию:false. -
emitClose<boolean> Должен ли поток генерировать событие'close'после уничтожения. По умолчанию:true. -
read<Function> Реализация методаstream._read(). -
destroy<Function> Реализация методаstream._destroy(). -
construct<Function> Реализация методаstream._construct(). -
autoDestroy<boolean> Должен ли этот поток автоматически вызывать.destroy()для себя после завершения. По умолчанию:true. -
signal<AbortSignal> Сигнал, представляющий возможную отмену.
-
const { Readable } = require('node:stream');
class MyReadable extends Readable {
constructor(options) {
// Calls the stream.Readable(options) constructor.
super(options);
// ...
}
} copy Или при использовании конструкторов в стиле до ES6:
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); copy Или при использовании упрощенного подхода к созданию конструктора:
const { Readable } = require('node:stream');
const myReadable = new Readable({
read(size) {
// ...
},
}); copy Вызов abort для AbortController, соответствующего переданному AbortSignal, будет вести себя так же, как вызов .destroy(new AbortError()) для созданного потока для чтения.
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(); copy
readable._construct(callback)
-
callback<Function> Вызовите эту функцию (при необходимости с аргументом ошибки), когда инициализация потока завершена.
Метод _construct() НЕ ДОЛЖЕН вызываться напрямую. Он может быть реализован дочерними классами; в этом случае его будут вызывать только внутренние методы класса Readable.
Эта необязательная функция будет запланирована конструктором потока на следующий такт, откладывая вызовы _read() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной инициализации ресурсов до того, как поток можно будет использовать.
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);
}
}
} copy
readable._read(size)
-
size<number> Количество байтов для асинхронного чтения
Эта функция НЕ ДОЛЖНА вызываться непосредственно кодом приложения. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Readable.
Все реализации потока Readable должны предоставлять реализацию метода readable._read() для получения данных из нижележащего ресурса.
При вызове readable._read(), если из ресурса доступны данные, реализация должна начать помещать их в очередь чтения с помощью метода this.push(dataChunk). _read() будет вызван снова после каждого вызова this.push(dataChunk), как только поток будет готов принять новые данные. _read() может продолжать чтение из ресурса и помещение данных в очередь, пока readable.push() не вернет false. Возобновлять помещение дополнительных данных в очередь следует только после повторного вызова _read(), когда оно было остановлено.
После вызова метода readable._read() он не будет вызван снова, пока через метод readable.push() не будут переданы дополнительные данные. Пустые данные, например пустые буферы и строки, не приводят к вызову readable._read().
Аргумент size носит рекомендательный характер. В реализациях, где «чтение» — это одна операция, возвращающая данные, аргумент size можно использовать для определения объема запрашиваемых данных. Другие реализации могут игнорировать этот аргумент и просто предоставлять данные по мере их появления. Не нужно «ждать», пока станет доступно size байт, прежде чем вызывать stream.push(chunk).
Метод readable._read() начинается с символа подчеркивания, поскольку он является внутренним для определяющего его класса и никогда не должен вызываться напрямую пользовательскими программами.
readable._destroy(err, callback)
-
err<Error> Возможная ошибка. -
callback<Function> Функция обратного вызова, принимающая необязательный аргумент ошибки.
Метод _destroy() вызывается методом readable.destroy(). Его можно переопределить в дочерних классах, но его нельзя вызывать напрямую.
readable.push(chunk[, encoding])
-
chunk<Buffer> | <TypedArray> | <DataView> | <string> | <null> | <any> Фрагмент данных для помещения в очередь чтения. Для потоков, работающих не в режиме объектов,chunkдолжен иметь тип <string>, <Buffer>, <TypedArray> или <DataView>. Для потоков в режиме объектовchunkможет иметь любое значение JavaScript. -
encoding<string> Кодировка строковых фрагментов. Должна быть допустимой кодировкойBuffer, например'utf8'или'ascii'. - Возвращает: <boolean>
true, если можно продолжать помещать дополнительные фрагменты данных; в противном случае —false.
Если chunk имеет тип <Buffer>, <TypedArray>, <DataView> или <string>, chunk данных будет добавлено во внутреннюю очередь, откуда пользователи потока смогут их получить. Передача chunk в качестве null сигнализирует о конце потока (EOF), после которого записывать данные больше нельзя.
Когда Readable работает в приостановленном режиме, данные, добавленные с помощью readable.push(), можно прочитать, вызвав метод readable.read() при генерации события 'readable'.
Когда Readable работает в режиме передачи данных, данные, добавленные с помощью readable.push(), будут переданы путем генерации события 'data'.
Метод readable.push() разработан с максимальной гибкостью. Например, низкоуровневый источник, предоставляющий некоторый механизм приостановки/возобновления и функцию обратного вызова для данных, можно обернуть пользовательским экземпляром Readable:
// `_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();
}
} copy Метод readable.push() используется для помещения содержимого во внутренний буфер. Его можно вызывать из метода readable._read().
Для потоков, работающих не в режиме объектов, если параметр chunk метода readable.push() равен undefined, он будет считаться пустой строкой или буфером. Дополнительные сведения см. в разделе readable.push('').
Ошибки при чтении
Ошибки, возникающие при обработке readable._read(), должны передаваться через метод readable.destroy(err). Выбрасывание Error изнутри readable._read() или ручная генерация события 'error' приводит к неопределенному поведению.
const { Readable } = require('node:stream');
const myReadable = new Readable({
read(size) {
const err = checkSomeErrorCondition();
if (err) {
this.destroy(err);
} else {
// Do some work.
}
},
}); copy Пример потока-счетчика
Ниже приведен простой пример потока Readable, который последовательно выдает числа от 1 до 1 000 000, а затем завершается.
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);
}
}
} copy Реализация дуплексного потока
Duplex — это поток, реализующий одновременно Readable и Writable, например TCP-соединение через сокет.
Поскольку JavaScript не поддерживает множественное наследование, для реализации потока Duplex расширяется класс stream.Duplex (вместо расширения классов stream.Readable и stream.Writable).
Класс stream.Duplex наследуется по прототипу от stream.Readable и паразитически — от stream.Writable, однако instanceof корректно работает с обоими базовыми классами благодаря переопределению Symbol.hasInstance в stream.Writable.
Пользовательские потоки Duplex должны вызывать конструктор new stream.Duplex([options]) и реализовывать оба метода: readable._read() и writable._write().
new stream.Duplex(options)
-
options<Object> Передается конструкторамWritableиReadable. Также содержит следующие поля:-
allowHalfOpen<boolean> Если задано значениеfalse, сторона для записи потока автоматически завершится при завершении стороны для чтения. По умолчанию:true. -
readable<boolean> Определяет, должна ли сторонаDuplexбыть доступна для чтения. По умолчанию:true. -
writable<boolean> Определяет, должна ли сторонаDuplexбыть доступна для записи. По умолчанию:true. -
readableObjectMode<boolean> ЗадаетobjectModeдля стороны потока, доступной для чтения. Не влияет на поток, еслиobjectModeимеет значениеtrue. По умолчанию:false. -
writableObjectMode<boolean> ЗадаетobjectModeдля стороны потока, доступной для записи. Не влияет на поток, еслиobjectModeимеет значениеtrue. По умолчанию:false. -
readableHighWaterMark<number> ЗадаетhighWaterMarkдля стороны потока, доступной для чтения. Не влияет на поток, если указанhighWaterMark. -
writableHighWaterMark<number> ЗадаетhighWaterMarkдля стороны потока, доступной для записи. Не влияет на поток, если указанhighWaterMark.
-
const { Duplex } = require('node:stream');
class MyDuplex extends Duplex {
constructor(options) {
super(options);
// ...
}
} copy Или при использовании конструкторов в стиле до ES6:
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); copy Или при использовании упрощенного подхода к созданию конструктора:
const { Duplex } = require('node:stream');
const myDuplex = new Duplex({
read(size) {
// ...
},
write(chunk, encoding, callback) {
// ...
},
}); copy При использовании конвейера:
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');
}
},
); copy Пример дуплексного потока
Ниже показан простой пример потока Duplex, оборачивающего гипотетический низкоуровневый объект-источник, в который можно записывать данные и из которого можно читать данные, хотя его API несовместим с потоками Node.js. Также ниже показан простой пример потока Duplex, который буферизует входящие записанные данные через интерфейс Writable, а затем возвращает их через интерфейс Readable.
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));
});
}
} copy Важнейшая особенность потока Duplex состоит в том, что стороны Readable и Writable работают независимо друг от друга, несмотря на то что сосуществуют в одном экземпляре объекта.
Дуплексные потоки в режиме объектов
Для потоков Duplex параметр objectMode можно задать только для стороны Readable или Writable, используя соответственно параметры readableObjectMode и writableObjectMode.
Например, в следующем примере создается новый поток Transform (разновидность потока Duplex), у которого сторона Writable работает в режиме объектов и принимает числа JavaScript, преобразуемые в шестнадцатеричные строки на стороне Readable.
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 copy Реализация преобразующего потока
Поток Transform — это поток Duplex, в котором выходные данные каким-либо образом вычисляются на основе входных. К таким потокам относятся потоки zlib или crypto, которые сжимают, шифруют или расшифровывают данные.
Выходные данные не обязательно должны иметь тот же размер, то же количество фрагментов или поступать в то же время, что и входные. Например, у потока Hash будет только один выходной фрагмент, который предоставляется после завершения ввода. Поток zlib будет выдавать данные, объем которых может быть значительно меньше или больше объема входных данных.
Для реализации потока Transform расширяется класс stream.Transform.
Класс stream.Transform наследуется по прототипу от stream.Duplex и реализует собственные версии методов writable._write() и readable._read(). Пользовательские реализации Transform должны реализовывать метод transform._transform() и могут также реализовывать метод transform._flush().
При использовании потоков Transform следует учитывать, что записываемые в поток данные могут приостановить сторону Writable, если выходные данные на стороне Readable не считываются.
new stream.Transform([options])
-
options<Object> Передается конструкторамWritableиReadable. Также содержит следующие поля:-
transform<Function> Реализация методаstream._transform(). -
flush<Function> Реализация методаstream._flush().
-
const { Transform } = require('node:stream');
class MyTransform extends Transform {
constructor(options) {
super(options);
// ...
}
} copy Или при использовании конструкторов в стиле до ES6:
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); copy Или при использовании упрощенного подхода к созданию конструктора:
const { Transform } = require('node:stream');
const myTransform = new Transform({
transform(chunk, encoding, callback) {
// ...
},
}); copy Событие: 'end'
Событие 'end' относится к классу stream.Readable. Событие 'end' генерируется после вывода всех данных, то есть после вызова функции обратного вызова в transform._flush(). В случае ошибки событие 'end' генерироваться не должно.
Событие: 'finish'
Событие 'finish' относится к классу stream.Writable. Событие 'finish' генерируется после вызова stream.end() и обработки всех фрагментов методом stream._transform(). В случае ошибки событие 'finish' генерироваться не должно.
transform._flush(callback)
-
callback<Function> Функция обратного вызова (при необходимости с аргументом ошибки и данными), вызываемая после сброса оставшихся данных.
Эта функция НЕ ДОЛЖНА вызываться кодом приложения. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Readable.
В некоторых случаях при преобразовании в конце потока может потребоваться выдать дополнительные данные. Например, поток сжатия zlib хранит определенный объем внутреннего состояния, необходимого для оптимального сжатия выходных данных. Однако при завершении потока эти дополнительные данные нужно сбросить, чтобы сжатые данные были полными.
Пользовательские реализации Transform могут реализовывать метод transform._flush(). Он будет вызван, когда больше не останется записанных данных для обработки, но до генерации события 'end', сигнализирующего о завершении потока Readable.
В реализации transform._flush() метод transform.push() может вызываться ноль или более раз, по мере необходимости. Функцию callback следует вызвать по завершении операции сброса.
Метод transform._flush() начинается с символа подчеркивания, поскольку он является внутренним для определяющего его класса и никогда не должен вызываться напрямую пользовательскими программами.
transform._transform(chunk, encoding, callback)
-
chunk<Buffer> | <string> | <any> ФрагментBufferдля преобразования, полученный изstring, переданного вstream.write(). Если параметр потокаdecodeStringsимеет значениеfalseили поток работает в режиме объектов, фрагмент не преобразуется и будет иметь то же значение, что и переданное вstream.write(). -
encoding<string> Если фрагмент является строкой, это тип кодировки. Если фрагмент является буфером, это специальное значение'buffer'. В этом случае игнорируйте его. -
callback<Function> Функция обратного вызова (при необходимости с аргументом ошибки и данными), вызываемая после обработки переданногоchunk.
Эта функция НЕ ДОЛЖНА вызываться кодом приложения. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Readable.
Все реализации потока Transform должны предоставлять метод _transform() для приема входных данных и создания выходных. Реализация transform._transform() обрабатывает записываемые байты, вычисляет выходные данные, а затем передает их читаемой части с помощью метода transform.push().
Метод transform.push() может вызываться ноль или более раз для создания выходных данных из одного входного фрагмента — в зависимости от того, сколько данных нужно выдать в результате обработки этого фрагмента.
Для любого заданного фрагмента входных данных выходные данные могут не создаваться.
Функцию callback следует вызывать только после полной обработки текущего фрагмента. Первым аргументом, передаваемым в callback, должен быть объект Error, если при обработке входных данных произошла ошибка, и null в противном случае. Если в callback передается второй аргумент, он будет передан методу transform.push(), но только если первый аргумент имеет ложное значение. Иными словами, следующие варианты эквивалентны:
transform.prototype._transform = function(data, encoding, callback) {
this.push(data);
callback();
};
transform.prototype._transform = function(data, encoding, callback) {
callback(null, data);
}; copy Метод transform._transform() начинается с символа подчеркивания, поскольку он является внутренним для определяющего его класса и никогда не должен вызываться напрямую пользовательскими программами.
transform._transform() никогда не вызывается параллельно: потоки реализуют механизм очереди, и для получения следующего фрагмента необходимо вызвать callback — синхронно или асинхронно.
Класс: stream.PassThrough
Класс stream.PassThrough — это простая реализация потока Transform, которая просто передает входные байты на выход. Он предназначен главным образом для примеров и тестирования, однако существуют сценарии, в которых stream.PassThrough служит полезным строительным блоком для новых типов потоков.
Дополнительные примечания
Совместимость потоков с асинхронными генераторами и асинхронными итераторами
Благодаря поддержке асинхронных генераторов и итераторов в JavaScript асинхронные генераторы теперь фактически являются потоковой конструкцией языка первого класса.
Ниже приведены распространённые примеры взаимодействия потоков Node.js с асинхронными генераторами и асинхронными итераторами.
Чтение потоков для чтения с помощью асинхронных итераторов
(async function() {
for await (const chunk of readable) {
console.log(chunk);
}
})(); copy Асинхронные итераторы регистрируют в потоке постоянный обработчик ошибок, чтобы предотвратить необработанные ошибки, возникающие после уничтожения потока.
Создание потоков для чтения с помощью асинхронных генераторов
Поток для чтения Node.js можно создать из асинхронного генератора с помощью вспомогательного метода Readable.from():
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);
}); copy Передача данных в потоки для записи из асинхронных итераторов
При записи в поток для записи из асинхронного итератора обеспечьте правильную обработку обратного давления и ошибок. stream.pipeline() берёт на себя обработку обратного давления и связанных с ним ошибок:
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();
}); copy Совместимость со старыми версиями Node.js
До Node.js 0.10 интерфейс потоков Readable был проще, но при этом менее мощным и полезным.
- Вместо ожидания вызовов метода
stream.read()события'data'начинали генерироваться немедленно. Приложениям, которым требовалось выполнить некоторую работу, чтобы решить, как обработать данные, приходилось сохранять прочитанные данные в буферах, чтобы не потерять их. - Метод
stream.pause()носил рекомендательный характер, а не гарантированный. Это означало, что необходимо было быть готовым получать события'data'даже когда поток находился в приостановленном состоянии.
В Node.js 0.10 был добавлен класс Readable. Для обратной совместимости со старыми программами Node.js потоки Readable переключаются в «режим передачи данных», когда добавляется обработчик события 'data' или вызывается метод stream.resume(). В результате даже при отказе от нового метода stream.read() и события 'readable' больше не нужно беспокоиться о потере фрагментов 'data'.
Хотя большинство приложений продолжат работать как обычно, в следующих условиях может возникнуть особый случай:
- Не добавлен обработчик события
'data'. - Метод
stream.resume()никогда не вызывается. - Поток не передаёт данные ни в один поток для записи.
Например, рассмотрим следующий код:
// 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); copy До Node.js 0.10 данные входящего сообщения просто отбрасывались. Однако в Node.js 0.10 и более поздних версиях сокет остаётся приостановленным навсегда.
В этой ситуации можно обойти проблему, вызвав метод stream.resume(), чтобы начать передачу данных:
// 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); copy Помимо переключения новых потоков Readable в режим передачи данных, потоки старого формата (до версии 0.10) можно обернуть в класс Readable с помощью метода readable.wrap().
readable.read(0)
В некоторых случаях необходимо инициировать обновление механизмов базового потока для чтения, не потребляя при этом никаких данных. В таких случаях можно вызвать readable.read(0), который всегда возвращает null.
Если внутренний буфер чтения меньше highWaterMark и поток в данный момент не выполняет чтение, вызов stream.read(0) инициирует низкоуровневый вызов stream._read().
Хотя большинству приложений почти никогда не потребуется это делать, в Node.js есть ситуации, когда это необходимо, особенно во внутренних механизмах класса потоков Readable.
readable.push('')
Использовать readable.push('') не рекомендуется.
Передача строки <string>, объекта <Buffer>, объекта <TypedArray> или объекта <DataView> нулевой длины в поток, работающий не в объектном режиме, имеет интересный побочный эффект. Поскольку это вызов readable.push(), он завершает процесс чтения. Однако, поскольку аргументом является пустая строка, в буфер чтения не добавляются данные, поэтому пользователю нечего считывать.
Несоответствие highWaterMark после вызова readable.setEncoding()
Использование readable.setEncoding() изменит поведение highWaterMark в не объектном режиме.
Обычно размер текущего буфера сравнивается с highWaterMark в байтах. Однако после вызова setEncoding() функция сравнения начинает измерять размер буфера в символах.
В распространённых случаях с latin1 или ascii это не создаёт проблем. Но при работе со строками, содержащими многобайтовые символы, рекомендуется учитывать такое поведение.
© Joyent, Inc. and other Node contributors
Licensed under the MIT License.
Node.js is a trademark of Joyent, Inc. and is used with its permission.
We are not endorsed by or affiliated with Joyent.
https://nodejs.org/dist/latest-v22.x/docs/api/stream.html