Поток[src]
Исходный код: lib/stream.js
Поток — это абстрактный интерфейс для работы с потоковыми данными в Node.js. Модуль stream предоставляет API для реализации интерфейса потока.
Node.js предоставляет множество объектов потоков. Например, запрос к HTTP-серверу и process.stdout — оба являются экземплярами потоков.
Потоки могут быть читаемыми, записываемыми или обоими. Все потоки являются экземплярами EventEmitter.
Для доступа к модулю stream:
const stream = require('stream'); Модуль stream полезен для создания новых типов экземпляров потоков. Обычно нет необходимости использовать модуль stream для потребления потоков.
Структура этого документа
Этот документ содержит две основные секции и третью секцию для заметок. В первой секции объясняется, как использовать существующие потоки в приложении. Во второй секции объясняется, как создавать новые типы потоков.
Типы потоков
В Node.js существует четыре основных типа потоков:
-
Writable: потоки, в которые можно записывать данные (например,fs.createWriteStream()). -
Readable: потоки, из которых можно читать данные (например,fs.createReadStream()). -
Duplex: потоки, которые являются какReadable, так иWritable(например,net.Socket). -
Transform:Duplexпотоки, которые могут изменять или преобразовывать данные при записи и чтении (например,zlib.createDeflate()).
Кроме того, этот модуль включает служебные функции stream.pipeline(), stream.finished(), stream.Readable.from() и stream.addAbortSignal().
API потоков с обещаниями
API stream/promises предоставляет альтернатный набор асинхронных служебных функций для потоков, которые возвращают объекты Promise вместо использования обратных вызовов. К API можно получить доступ через require('stream/promises') или require('stream').promises.
Режим объектов
Все потоки, созданные API Node.js, работают исключительно со строками и Buffer (или Uint8Array) объектами. Однако реализации потоков могут работать с другими типами JavaScript-значений (за исключением null, который служит специальной цели в потоках). Такие потоки считаются работающими в «режиме объектов».
Экземпляры потоков переключаются в режим объектов с помощью параметра objectMode при создании потока. Попытка переключить существующий поток в режим объектов небезопасна.
Буферизация
Как Writable, так и Readable потоки будут хранить данные во внутренней буферной памяти.
Количество данных, потенциально буферизуемых, зависит от параметра highWaterMark , переданного в конструктор потока. Для обычных потоков параметр highWaterMark определяет общее количество байтов. Для потоков, работающих в режиме объектов, параметр highWaterMark определяет общее количество объектов.
Данные буферизуются в 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('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 in JSON at position 1 Writable потоки (например, res в примере) предоставляют методы, такие как write() и end(), которые используются для записи данных в поток.
Readable потоки используют API EventEmitter для уведомления кода приложения о том, что данные доступны для чтения из потока. Эти доступные данные можно прочитать из потока несколькими способами.
Оба Writable и Readable потока используют API EventEmitter различными способами для передачи текущего состояния потока.
Duplex и Transform потоки являются как Writable, так и Readable.
Приложения, которые либо записывают данные в поток, либо потребляют данные из потока, не обязаны реализовывать интерфейсы потоков напрямую и обычно не имеют причин вызывать require('stream').
Разработчики, желающие реализовать новые типы потоков, должны обратиться к разделу API для разработчиков потоков.
Потоки записи
Потоки записи — это абстракция назначения, в которое записываются данные.
Примеры Writable потоков включают:
- HTTP-запросы, на стороне клиента
- HTTP-ответы, на стороне сервера
- Потоки записи в fs
- Потоки zlib
- Потоки криптографии
- 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'); Класс: 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);
}
}
} Событие: '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'); Событие: '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); Событие: 'unpipe'
-
srcИсточник-поток <stream.Readable>, который отсоединил этот поток записи
Событие 'unpipe' генерируется, когда метод stream.unpipe() вызывается на потоке Readable, удаляя этот поток 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);
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<Ошибка> Необязательно, ошибка, которую нужно излучить со событием'error'. - Возвращает: <this>
Уничтожить поток. Необязательно излучить событие 'error', и излучить событие 'close' (если emitClose не установлено в false). После этого вызова поток записи завершен, и последующие вызовы write() или end() приведут к ошибке ERR_STREAM_DESTROYED. Это деструктивный и немедленный способ уничтожить поток. Предыдущие вызовы write() могут не завершиться, и могут вызвать ошибку ERR_STREAM_DESTROYED. Используйте end() вместо destroy, если данные должны быть отправлены прежде, чем произойдёт закрытие, или подождите события 'drain' перед уничтожением потока.
const { Writable } = require('stream');
const myStream = new Writable();
const fooErr = new Error('foo error');
myStream.destroy(fooErr);
myStream.on('error', (fooErr) => console.error(fooErr.message)); // foo error const { Writable } = require('stream');
const myStream = new Writable();
myStream.destroy();
myStream.on('error', function wontHappen() {}); const { Writable } = require('stream');
const myStream = new Writable();
myStream.destroy();
myStream.write('foo', (error) => console.error(error.code));
// ERR_STREAM_DESTROYED После того, как destroy() был вызван, любые дальнейшие вызовы будут недействительными, и больше никаких ошибок, кроме _destroy(), не могут быть излучены в виде 'error'.
Реализаторы не должны переопределять этот метод, а вместо этого реализовывать writable._destroy().
writable.destroyed
Является ли поток true после вызова writable.destroy().
const { Writable } = require('stream');
const myStream = new Writable();
console.log(myStream.destroyed); // false
myStream.destroy();
console.log(myStream.destroyed); // true
writable.end([chunk[, encoding]][, callback])
-
chunk<string> | <Buffer> | <Uint8Array> | <any> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжно быть строкой,BufferилиUint8Array. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript, кромеnull. -
encoding<string> Кодировка, еслиchunkявляется строкой -
callback<Function> Обратный вызов, когда поток завершен. - Возвращает: <this>
Вызов метода writable.end() сигнализирует о том, что больше данных не будет записано в Writable. Необязательные аргументы chunk и encoding позволяют записать последний фрагмент данных непосредственно перед закрытием потока.
Вызов метода stream.write() после вызова stream.end() вызовет ошибку.
// Write 'hello, ' and then end with 'world!'.
const fs = require('fs');
const file = fs.createWriteStream('example.txt');
file.write('hello, ');
file.end('world!');
// Writing more now is not allowed!
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()); Если метод 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();
}); См. также: writable.cork().
writable.writable
Истинно, если безопасно вызывать writable.write(), что означает, что поток не был уничтожен, не получил ошибку или не был завершен.
writable.writableEnded
Истинно после вызова writable.end(). Это свойство не указывает, были ли данные очищены, для этого используйте writable.writableFinished.
writable.writableCorked
Количество вызовов writable.uncork(), необходимых для полной разблокировки потока.
writable.writableFinished
Устанавливается в true непосредственно перед тем, как будет отправлено событие 'finish'.
writable.writableHighWaterMark
Возвращает значение highWaterMark заданное при создании этого Writable.
writable.writableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение предоставляет данные для интроспекции состояния highWaterMark.
writable.writableNeedDrain
Истинно, если буфер потока заполнен, и поток выпустит событие 'drain'.
writable.writableObjectMode
Получение свойства objectMode заданного потока Writable.
writable.write(chunk[, encoding][, callback])
-
chunk<string> | <Buffer> | <Uint8Array> | <any> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжно быть строкой,BufferилиUint8Array. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript, кромеnull. -
encoding<string> | <null> Кодировка, еслиchunkявляется строкой. По умолчанию:'utf8' -
callback<Function> Обратный вызов, когда этот фрагмент данных очищен. - Возвращает: <boolean>
falseесли поток хочет, чтобы вызывающий код ждал события'drain'перед продолжением записи дополнительных данных; иначеtrue.
Метод writable.write() записывает данные в поток и вызывает переданный callback после того, как данные будут полностью обработаны. При возникновении ошибки, callback вызывается с ошибкой в качестве первого аргумента. callback вызывается асинхронно и до того, как будет отправлено событие 'error'.
Возвращаемое значение true если внутренний буфер меньше highWaterMark заданного при создании потока после принятия chunk. Если возвращено false, дальнейшие попытки записи данных в поток должны быть приостановлены до тех пор, пока не будет отправлено событие 'drain'.
Пока поток не очищен, вызовы write() будут буферизовать chunk, и возвращать false. После того, как все буферизованные фрагменты будут очищены (приняты для доставки операционной системой), будет отправлено событие 'drain' . Рекомендуется, что после того, как write() вернет false, больше фрагментов не следует записывать до тех пор, пока не будет отправлено событие 'drain'. Вызов write() для потока, который не очищен, разрешен, но Node.js будет буферизовать все записанные фрагменты, пока не произойдет максимальное использование памяти, после чего он прервется безусловно. Даже до прерывания, высокое использование памяти вызовет плохую работу сборщика мусора и высокое значение RSS (которое обычно не возвращается системе, даже после того, как память больше не нужна). Так как TCP-сокеты могут никогда не очиститься, если удалённый узел не читает данные, запись в сокет, который не очищен, может привести к удалённо эксплуатируемой уязвимости.
Запись данных, пока поток не очищен, особенно проблематична для Transform, потому что потоки Transform приостановлены по умолчанию до тех пор, пока они не будут направлены или не будет добавлен обработчик события 'data' или 'readable'.
Если данные для записи могут быть сгенерированы или получены по требованию, рекомендуется инкапсулировать логику в Readable и использовать stream.pipe(). Однако, если вы предпочитаете вызывать write(), можно соблюдать обратную загрузку (backpressure) и избежать проблем с памятью, используя событие '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.');
}); Поток Writable в режиме объектов всегда будет игнорировать аргумент encoding.
Потоки чтения
Потоки чтения — это абстракция для источника, из которого потребляются данные.
Примеры потоков Readable включают:
- HTTP-ответы на стороне клиента
- HTTP-запросы на стороне сервера
- Потоки чтения fs
- Потоки zlib
- Потоки crypto
- TCP-сокеты
- Выходные потоки child process (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(), или получение обратной загрузки (backpressure) приведут к тому, что состояние потока readable.readableFlowing будет установлено как false, временно приостановив поток событий, но не останавливая генерацию данных. В этом состоянии подключение обработчика события 'data' не переключит readable.readableFlowing в состояние true.
const { PassThrough, Writable } = require('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()); });
pass.write('ok'); // Will not emit 'data'.
pass.resume(); // Must be called to make stream emit 'data'. Пока readable.readableFlowing находится в состоянии false, данные могут накапливаться в внутренему буфере потока.
Выбор одного стиля API
API потоков Readable эволюционировал на протяжении нескольких версий Node.js и предоставляет несколько методов потребления данных потока. В целом, разработчики должны выбрать один метод потребления данных и никогда не использовать несколько методов для потребления данных из одного потока. В частности, использование комбинации on('data'), on('readable'), pipe(), или асинхронных итераторов может привести к неинтуитивному поведению.
Для большинства пользователей рекомендуется использование метода readable.pipe(), поскольку он реализован для обеспечения наилучшего способа потребления данных потока. Разработчики, которым требуется более тонкий контроль над передачей и генерацией данных, могут использовать EventEmitter и readable.on('readable')/readable.read() или API readable.pause()/readable.resume().
Класс: stream.Readable
Событие: 'close'
Событие 'close' испускается, когда поток и все его базовые ресурсы (например, дескриптор файла) были закрыты. Событие указывает, что больше событий не будет испускаться, и дальнейшие вычисления не будут производиться.
Поток Readable всегда будет испускать событие 'close' если он был создан с параметром emitClose.
Событие: 'data'
-
chunk<Буфер> | <строка> | <любое> Блок данных. Для потоков, не работающих в режиме объектов, блок будет строкой или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.`);
}); Событие: '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.');
}); Событие: 'error'
Событие 'error' может быть испущено реализацией потока Readable в любое время. Обычно это происходит, если базовый поток не может генерировать данные из-за внутренней ошибки или при попытке потока отправить неверный блок данных.
Обработчик события получит объект Error.
Событие: 'pause'
Событие 'pause' испускается, когда вызывается stream.pause(), и readableFlowing не находится в состоянии false.
Событие: 'readable'
Событие 'readable' генерируется, когда данные доступны для чтения из потока или когда достигнут конец потока. По сути, событие 'readable' указывает на то, что у потока есть новые данные. Если данные доступны, stream.read() вернёт эти данные.
const readable = getReadableStreamSomehow();
readable.on('readable', function() {
// There is some data to read now.
let data;
while (data = this.read()) {
console.log(data);
}
}); Если достигнут конец потока, вызов stream.read() вернёт null и вызовет событие 'end'. Это также верно, если данных для чтения никогда не было. Например, в следующем примере foo.txt — это пустой файл:
const fs = require('fs');
const rr = fs.createReadStream('foo.txt');
rr.on('readable', () => {
console.log(`readable: ${rr.read()}`);
});
rr.on('end', () => {
console.log('end');
}); Вывод выполнения этого скрипта:
$ node test.js readable: null end
В некоторых случаях подключение обработчика события '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<Ошибка> Ошибка, которая будет передана в качестве полезной нагрузки в событии'error' - Возвращает: <this>
Уничтожить поток. При необходимости сгенерировать событие 'error' и событие 'close' (если emitClose не равно false). После этого вызова читабельный поток освободит все внутренние ресурсы, а последующие вызовы push() будут игнорироваться.
После вызова destroy() все последующие вызовы будут бездействовать, и не будут генерироваться ошибки, кроме тех, что возникают из _destroy(), как 'error'.
Реализаторы не должны переопределять этот метод, но вместо этого реализовать readable._destroy().
readable.destroyed
Была ли true после вызова 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
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);
}); Метод readable.pause() не имеет эффекта, если есть обработчик события 'readable'.
readable.pipe(destination[, options])
-
destination<stream.Writable> Цель для записи данных -
options<Объект> Параметры конвейера-
end<boolean> Завершить писатель, когда читатель завершит. По умолчанию:true.
-
- Возвращает: <stream.Writable> Назначение, позволяющее создавать цепочки конвейеров, если это
DuplexилиTransformпоток
Метод readable.pipe() подключает поток Writable к readable, вызывая автоматическое переключение в режим потока и отправку всех данных в подключенный Writable. Поток данных будет управляться автоматически, чтобы поток назначения Writable не перегружался более быстрым потоком Readable.
Следующий пример перенаправляет все данные из readable в файл с именем file.txt:
const fs = require('fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt'.
readable.pipe(writable); Возможно подключение нескольких потоков Writable к одному потоку Readable.
Метод readable.pipe() возвращает ссылку на целевой поток, что позволяет создавать цепочки соединённых потоков:
const fs = require('fs');
const r = fs.createReadStream('file.txt');
const z = zlib.createGzip();
const w = fs.createWriteStream('file.txt.gz');
r.pipe(z).pipe(w); По умолчанию, stream.end() вызывается в целевом Writable потоке, когда исходный Readable поток генерирует событие 'end', так что целевой поток больше не доступен для записи. Чтобы отключить это поведение по умолчанию, можно передать опцию end как false, что сохраняет открытым целевой поток:
reader.pipe(writer, { end: false });
reader.on('end', () => {
writer.end('Goodbye\n');
}); Важно отметить, что если поток Readable генерирует ошибку во время обработки, целевой поток Writable не закрывается автоматически. Если произошла ошибка, необходимо вручную закрыть каждый поток, чтобы избежать утечек памяти.
Потоки process.stderr и process.stdout Writable никогда не закрываются до выхода процесса Node.js, независимо от указанных параметров.
readable.read([size])
-
size<число> Необязательный аргумент для указания количества считываемых данных. - Возвращает: <строка> | <Буфер> | <null> | <любой>
Метод 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.');
}); Каждый вызов readable.read() возвращает фрагмент данных или null. Фрагменты не конкатенируются. Для получения всех данных в буфере необходим цикл 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('');
}); Поток Readable в режиме объектов всегда возвращает один элемент из вызова readable.read(size), независимо от значения аргумента size.
Если метод readable.read() возвращает фрагмент данных, то также будет выведено событие 'data'.
Вызов stream.read([size]) после того, как было выведено событие 'end', вернёт null. Ошибка во время выполнения не будет поднята.
readable.readable
Имеет значение true , если безопасно вызвать readable.read(), что означает, что поток не был уничтожен и не выведено событие 'error' или 'end'.
readable.readableAborted
Возвращает, был ли поток уничтожен или произошла ошибка до вывода события 'end'.
readable.readableDidRead
Возвращает, было ли выведено событие 'data'.
readable.readableEncoding
Получатель свойства encoding заданного потока Readable. Свойство encoding может быть установлено с помощью метода readable.setEncoding().
readable.readableEnded
Приобретает значение true при выводе события 'end'.
readable.readableFlowing
Это свойство отражает текущее состояние потока Readable, как описано в разделе Три состояния.
readable.readableHighWaterMark
Возвращает значение highWaterMark , переданное при создании этого потока Readable.
readable.readableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к чтению. Значение предоставляет данные о состоянии highWaterMark.
readable.readableObjectMode
Получатель свойства objectMode заданного потока Readable.
readable.resume()
- Возвращает: <this>
Метод readable.resume() заставляет явно приостановленный поток Readable возобновить выведение событий 'data', переключая поток в режим потоковой передачи.
Метод readable.resume() может быть использован для полного потребления данных из потока, не обрабатывая эти данные:
getReadableStreamSomehow()
.resume()
.on('end', () => {
console.log('Reached the end, but did not read anything.');
}); Метод 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);
});
readable.unpipe([destination])
-
destination<stream.Writable> Необязательный конкретный поток для отключения - Возвращает: <this>
Метод readable.unpipe() отсоединяет поток Writable , ранее подключенный с помощью метода stream.pipe().
Если destination не указан, то все подключения отсоединяются.
Если destination указан, но для него нет подключения, то метод ничего не делает.
const fs = require('fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt',
// but only for the first second.
readable.pipe(writable);
setTimeout(() => {
console.log('Stop writing to file.txt.');
readable.unpipe(writable);
console.log('Manually close the file stream.');
writable.end();
}, 1000);
readable.unshift(chunk[, encoding])
-
chunk<Buffer> | <Uint8Array> | <string> | <null> | <any> Фрагмент данных для добавления в начало очереди чтения. Для потоков, не работающих в режиме объектов,chunkдолжен быть строкой,Buffer,Uint8Arrayили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('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.match(/\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);
} else {
// Still reading the header.
header += str;
}
}
}
} В отличие от stream.push(chunk), stream.unshift(chunk) не завершит процесс чтения, сбросив внутреннее состояние чтения потока. Это может привести к неожиданным результатам, если readable.unshift() вызывается во время чтения (например, из реализации stream._read() в пользовательском потоке). После вызова readable.unshift() с немедленным stream.push('') состояние чтения будет сброшено надлежащим образом, однако лучше просто избегать вызова readable.unshift() во время выполнения чтения.
readable.wrap(stream)
До версии Node.js 0.10 потоки не реализовывали весь API модуля stream, как он определён сейчас. (См. Совместимость для получения дополнительной информации.)
При использовании старой библиотеки Node.js, которая испускает события 'data' и имеет метод stream.pause(), который является только рекомендательным, можно использовать метод readable.wrap() для создания потока Readable, использующего старый поток в качестве источника данных.
Использовать readable.wrap() будет редко необходимо, но метод был предоставлен для удобства взаимодействия со старыми приложениями и библиотеками Node.js.
const { OldReader } = require('./old-api-module.js');
const { Readable } = require('stream');
const oreader = new OldReader();
const myReader = new Readable().wrap(oreader);
myReader.on('readable', () => {
myReader.read(); // etc.
});
readable[Symbol.asyncIterator]()
- Возвращает: <AsyncIterator> для полного потребления потока.
const fs = require('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); Если цикл завершается с break, return, или throw, поток будет уничтожен. Другими словами, перебор потока полностью потребляет его. Поток будет читаться кусками размером, равным значению параметра highWaterMark. В приведённом выше примере данные будут в одном куске, если файл имеет меньше 64 КБ данных, так как параметр highWaterMark не предоставлен для fs.createReadStream().
readable.iterator([options])
-
options<Объект>-
destroyOnReturn<логическое значение> Если установлено вfalse, вызовreturnдля async iterator или выход из цикла итерацииfor await...ofс помощьюbreak,return, илиthrowне уничтожит поток. По умолчанию:true. -
destroyOnError<логическое значение> Если установлено вfalse, если поток испускает ошибку во время итерации, итератор не уничтожит поток. По умолчанию:true.
-
- Возвращает: <AsyncIterator> для потребления потока.
Итератор, созданный этим методом, даёт пользователям возможность отменить уничтожение потока, если цикл for await...of завершается с return, break, или throw, или если итератор должен уничтожить поток, если поток испустил ошибку во время итерации.
const { Readable } = require('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(); Потоки duplex и transform
Класс: stream.Duplex
Потоки duplex — это потоки, которые реализуют как интерфейс Readable, так и Writable.
Примеры потоков Duplex включают:
duplex.allowHalfOpen
Если false , поток автоматически завершит сторону записи при завершении стороны чтения. Изначально установлено в конструкторе allowHalfOpen, по умолчанию false.
Это можно изменить вручную, чтобы изменить поведение полуоткрытых потоков Duplex , но это изменение должно быть выполнено до испускания события 'end'.
Класс: stream.Transform
Потоки transform — это потоки Duplex, где выход каким-то образом связан с входом. Как и все потоки Duplex, потоки Transform реализуют как интерфейс Readable, так и Writable.
Примеры потоков Transform включают:
transform.destroy([error])
Уничтожить поток и, по желанию, испустить событие 'error'. После этого вызова поток transform освободит все внутренние ресурсы. Реализаторы не должны переопределять этот метод, но вместо этого должны реализовать readable._destroy(). По умолчанию реализация _destroy() для Transform также испускает 'close' , если emitClose не установлено в false.
После того, как destroy() был вызван, все последующие вызовы будут ничтожны, и дальнейшие ошибки, кроме ошибок _destroy() , не могут быть испущены как 'error'.
stream.finished(stream[, options], callback)
-
stream<Поток> Поток для чтения и/или записи. -
options<Объект>-
error<логическое значение> Если установлено вfalse, вызовemit('error', err)не рассматривается как завершение. По умолчанию:true. -
readable<логическое значение> Если установлено вfalse, колбэк будет вызван при завершении потока, даже если поток всё ещё может читать. По умолчанию:true. -
writable<логическое значение> Если установлено вfalse, колбэк будет вызван при завершении потока, даже если поток всё ещё может писать. По умолчанию:true. -
signal<AbortSignal> позволяет прервать ожидание завершения потока. Сам поток не будет прерван, если сигнал прерван. Колбэк будет вызван сAbortError. Все зарегистрированные обработчики, добавленные этой функцией, также будут удалены.
-
-
callback<Функция> Функция колбэка, принимающая необязательный аргумент ошибки. - Возвращает: <Функция> Функция очистки, которая удаляет все зарегистрированные обработчики.
Функция для получения уведомления, когда поток больше не является читаемым, записываемым или столкнулся с ошибкой или преждевременным закрытием.
const { finished } = require('stream');
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. Особо полезна в сценариях обработки ошибок, когда поток разрушается преждевременно (например, при прерывании HTTP-запроса) и не испускает 'end' или 'finish'.
API finished предоставляет версию с promise:
const { finished } = require('stream/promises');
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. stream.finished() оставляет висящие обработчики событий (в частности 'error', 'end', 'finish' и 'close') после вызова callback. Причина в том, что непредвиденные 'error' события (из-за неправильной реализации потоков) не вызывают непредвиденных сбоев. Если это поведение нежелательно, то возвращённая функция очистки должна быть вызвана в обратном вызове:
const cleanup = finished(rs, (err) => {
cleanup();
// ...
});
stream.pipeline(source[, ...transforms], destination, callback)
stream.pipeline(streams, callback)
-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> -
source<Stream> | <Iterable> | <AsyncIterable> | <Function>- Возвращает: <Iterable> | <AsyncIterable>
-
...transforms<Stream> | <Function>-
source<AsyncIterable> - Возвращает: <AsyncIterable>
-
-
destination<Stream> | <Function>-
source<AsyncIterable> - Возвращает: <AsyncIterable> | <Promise>
-
-
callback<Function> Вызывается, когда конвейер полностью завершён.-
err<Error> -
valРезультат выполненияPromiseвозвращённыйdestination.
-
- Возвращает: <Stream>
Метод модуля для перенаправления данных между потоками и генераторами, перенаправляя ошибки и должным образом очищая всё и предоставляя обратный вызов, когда конвейер завершён.
const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('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.');
}
}
); pipeline API предоставляет версию с промисом, которая также может принимать аргумент options в качестве последнего параметра со свойством signal <AbortSignal>. Когда сигнал прерывается, destroy будет вызван на базовом конвейере с AbortError.
const { pipeline } = require('stream/promises');
async function run() {
await pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz')
);
console.log('Pipeline succeeded.');
}
run().catch(console.error); Для использования AbortSignal, передайте его внутри объекта options, как последний аргумент:
const { pipeline } = require('stream/promises');
async function run() {
const ac = new AbortController();
const signal = ac.signal;
setTimeout(() => ac.abort(), 1);
await pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz'),
{ signal },
);
}
run().catch(console.error); // AbortError pipeline API также поддерживает асинхронные генераторы:
const { pipeline } = require('stream/promises');
const fs = require('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); Не забудьте обработать аргумент signal , переданный в асинхронный генератор. Особенно в том случае, когда асинхронный генератор является источником для конвейера (т.е. первый аргумент), или конвейер никогда не завершится.
const { pipeline } = require('stream/promises');
const fs = require('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); stream.pipeline() вызовет stream.destroy(err) для всех потоков, кроме:
-
Readableпотоков, которые испустили'end'или'close'. -
Writableпотоков, которые испустили'finish'или'close'.
stream.pipeline() оставляет висящие обработчики событий в потоках после вызова callback. В случае повторного использования потоков после сбоя это может привести к утечкам обработчиков событий и необработанным ошибкам.
stream.compose(...streams)
stream.compose экспериментальная.-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> - Возвращает: <stream.Duplex>
Объединяет два или более потоков в поток Duplex , который записывает в первый поток и считывает из последнего. Каждый предоставленный поток направляется в следующий с использованием stream.pipeline. Если какой-либо из потоков даёт ошибку, то все потоки уничтожаются, включая внешний поток Duplex.
Поскольку stream.compose возвращает новый поток, который, в свою очередь, может (и должен) быть перенаправлен в другие потоки, он позволяет композицию. В отличие от передачи потоков в stream.pipeline, обычно первый поток является потоком чтения, а последний — потоком записи, образуя замкнутый цикл.
Если передан Function, то он должен быть фабричным методом, принимающим source Iterable.
import { compose, Transform } from 'stream';
const removeSpaces = new Transform({
transform(chunk, encoding, callback) {
callback(null, String(chunk).replace(' ', ''));
}
});
async function* toUpper(source) {
for await (const chunk of source) {
yield String(chunk).toUpperCase();
}
}
let res = '';
for await (const buf of compose(removeSpaces, toUpper).end('hello world')) {
res += buf;
}
console.log(res); // prints 'HELLOWORLD' stream.compose может быть использован для преобразования асинхронных итерируемых объектов, генераторов и функций в потоки.
-
AsyncIterableпреобразует в читаемыйDuplex. Не может отдаватьnull. -
AsyncGeneratorFunctionпреобразует в читаемый/записываемый трансформаторDuplex. Должен принять исходныйAsyncIterableв качестве первого параметра. Не может отдаватьnull. -
AsyncFunctionпреобразует в записываемыйDuplex. Должен вернуть либоnullилиundefined.
import { compose } from 'stream';
import { finished } from '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'
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('stream');
async function * generate() {
yield 'hello';
yield 'streams';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
}); Вызов Readable.from(string) или Readable.from(buffer) не будет итерировать строки или буферы для соответствия семантике других потоков по причинам производительности.
stream.Readable.isDisturbed(stream)
-
stream<stream.Readable> | <ReadableStream> - Возвращает:
boolean
Возвращает, был ли поток прочитан или отменён.
stream.Duplex.from(src)
-
src<Stream> | <Blob> | <ArrayBuffer> | <string> | <Iterable> | <AsyncIterable> | <AsyncGeneratorFunction> | <AsyncFunction> | <Promise> | <Object>
Утилитарный метод для создания дуплексных потоков.
-
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игнорируется. - Возвращает: <stream.Duplex>
stream.addAbortSignal(signal, stream)
-
signal<AbortSignal> Сигнал, представляющий возможность отмены -
stream<Stream> поток, к которому необходимо подключить сигнал
Подключает AbortSignal к потоку чтения или записи. Это позволяет коду управлять уничтожением потока с помощью AbortController.
Вызов abort на AbortController соответствующем переданному AbortSignal будет вести себя так же, как вызов .destroy(new AbortError()) на потоке.
const fs = require('fs');
const controller = new AbortController();
const read = addAbortSignal(
controller.signal,
fs.createReadStream(('object.json'))
);
// Later, abort the operation closing the stream
controller.abort(); Или с использованием 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;
}
}
})(); API для разработчиков потоков
API модуля stream разработан для удобной реализации потоков с использованием прототипного наследования JavaScript.
Сначала разработчик потока объявляет новый JavaScript-класс, который расширяет один из четырёх базовых классов потоков (stream.Writable, stream.Readable, stream.Duplex, или stream.Transform), убедившись, что вызван соответствующий конструктор родительского класса:
const { Writable } = require('stream');
class MyWritable extends Writable {
constructor({ highWaterMark, ...options }) {
super({ highWaterMark });
// ...
}
} При расширении потоков учтите, какие параметры пользователь может и должен предоставить перед передачей их в базовый конструктор. Например, если реализация делает предположения относительно параметров 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('stream');
const myWritable = new Writable({
construct(callback) {
// Initialize state and load resources...
},
write(chunk, encoding, callback) {
// ...
},
destroy() {
// Free resources...
}
}); Реализация потока для записи
Класс stream.Writable расширяется для реализации потока Writable.
Пользовательские потоки Writable обязательно должны вызывать конструктор new stream.Writable([options]) и реализовывать метод writable._write() и/или writable._writev().
new stream.Writable([options])
-
options<Объект>-
highWaterMark<число> Уровень буфера, когдаstream.write()начинает возвращатьfalse. По умолчанию:16384(16 КБ) или16для потоковobjectMode. -
decodeStrings<логическое> Нужно ли кодировать значенияstring, переданные вstream.write(), вBuffer(с кодировкой, указанной в вызовеstream.write()) перед передачей их вstream._write(). Другие типы данных не преобразуются (например,Bufferне декодируются вstring). Установка в false предотвратит преобразованиеstringвtrue. По умолчанию:true. -
defaultEncoding<строка> По умолчанию кодировка, используемая, когда кодировка не указана как аргумент дляstream.write(). По умолчанию:'utf8'. -
objectMode<логическое> Является лиstream.write(anyObj)допустимой операцией. При установке этого параметра можно записать значения JavaScript, отличные от строк,BufferилиUint8Array, если это поддерживается реализацией потока. По умолчанию:false. -
emitClose<логическое> Нужно ли потоку отправлять'close'после уничтожения. По умолчанию:true. -
write<Функция> Реализация методаstream._write(). -
writev<Функция> Реализация методаstream._writev(). -
destroy<Функция> Реализация методаstream._destroy(). -
final<Функция> Реализация методаstream._final(). -
construct<Функция> Реализация методаstream._construct(). -
autoDestroy<логическое> Нужно ли потоку автоматически вызвать.destroy()после завершения. По умолчанию:true. -
signal<AbortSignal> Сигнал, представляющий возможность отмены.
-
const { Writable } = require('stream');
class MyWritable extends Writable {
constructor(options) {
// Calls the stream.Writable() constructor.
super(options);
// ...
}
} Или, при использовании конструкторов в стиле до ES6:
const { Writable } = require('stream');
const util = require('util');
function MyWritable(options) {
if (!(this instanceof MyWritable))
return new MyWritable(options);
Writable.call(this, options);
}
util.inherits(MyWritable, Writable); Или, используя упрощённый подход к конструктору:
const { Writable } = require('stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
}
}); Вызов abort на AbortController, соответствующем переданному AbortSignal, будет работать так же, как вызов .destroy(new AbortError()) на потоке для записи.
const { Writable } = require('stream');
const controller = new AbortController();
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
},
signal: controller.signal
});
// Later, abort the operation closing the stream
controller.abort();
writable._construct(callback)
-
callback<Функция> Вызовите эту функцию (по желанию, с аргументом ошибки), когда поток завершит инициализацию.
Метод _construct() НЕ ДОЛЖЕН вызываться напрямую. Его может реализовывать дочерний класс, и в этом случае он будет вызываться только внутренними методами класса Writable.
Эта необязательная функция будет вызвана в следующем цикле после возвращения конструктора потока, отложив любые вызовы _write(), _final() и _destroy() до тех пор, пока не будет вызван callback. Это полезно для инициализации состояния или асинхронной инициализации ресурсов перед использованием потока.
const { Writable } = require('stream');
const fs = require('fs');
class WriteStream extends Writable {
constructor(filename) {
super();
this.filename = filename;
}
_construct(callback) {
fs.open(this.filename, (err, fd) => {
if (err) {
callback(err);
} else {
this.fd = fd;
callback();
}
});
}
_write(chunk, encoding, callback) {
fs.write(this.fd, chunk, callback);
}
_destroy(err, callback) {
if (this.fd) {
fs.close(this.fd, (er) => callback(er || err));
} else {
callback(err);
}
}
}
writable._write(chunk, encoding, callback)
-
chunk<Буфер> | <строка> | <любой> Данные, которые необходимо записать, преобразованные изstring, переданных вstream.write(). Если у потока опцияdecodeStringsимеет значениеfalseили поток работает в объектном режиме, фрагмент не будет преобразован и останется таким, каким был передан вstream.write(). -
encoding<строка> Если фрагмент — строка, тоencoding— кодировка символов этой строки. Если фрагмент —Buffer, или если поток работает в объектном режиме,encodingможет быть проигнорирована. -
callback<Функция> Вызовите эту функцию (при необходимости, с аргументом ошибки) по завершении обработки переданного фрагмента.
Все реализации потоков 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<Массив объектов> Данные для записи. Это массив объектов <объект>, каждый из которых представляет собой отдельный фрагмент данных для записи. Свойства этих объектов:-
chunk<Буфер> | <строка> Экземпляр буфера или строка, содержащая данные для записи.chunkбудет строкой, еслиWritableбыл создан с опциейdecodeStrings, установленной в значениеfalse, и вwrite()была передана строка. -
encoding<строка> Кодировка символовchunk. Еслиchunk—Buffer, тоencodingбудет'buffer'.
-
-
callback<Функция> Функция обратного вызова (при необходимости с аргументом ошибки), вызываемая по завершении обработки переданных фрагментов.
Эту функцию ПРИМЕНЯТЬ НЕОБХОДИМО НЕПОСРЕДСТВЕННО из кода приложения. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса Writable.
Метод writable._writev() может быть реализован дополнительно или как альтернатива writable._write() в реализациях потоков, способных обрабатывать несколько фрагментов данных одновременно. Если он реализован и если есть буферизованные данные от предыдущих записей, _writev() вызывается вместо _write().
Метод writable._writev() имеет префикс подчеркивания, поскольку является внутренним для определяющего его класса и не должен вызываться напрямую пользовательскими программами.
writable._destroy(err, callback)
-
err<Ошибка> Возможная ошибка. -
callback<Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.
Метод _destroy() вызывается методом writable.destroy(). Его можно переопределить в дочерних классах, но нельзя вызывать напрямую.
writable._final(callback)
-
callback<Функция> Вызовите эту функцию (при необходимости, с аргументом ошибки) по завершении записи оставшихся данных.
Метод _final() нельзя вызывать напрямую. Он может быть реализован дочерними классами и, если реализован, вызывается только внутренними методами класса Writable.
Эта необязательная функция вызывается перед закрытием потока, откладывая событие 'finish' до вызова callback. Это полезно для закрытия ресурсов или записи буферизованных данных перед завершением потока.
Ошибки при записи
Ошибки, возникающие во время обработки методов writable._write(), writable._writev() и writable._final(), необходимо обрабатывать, вызвав функцию обратного вызова и передав ошибку в качестве первого аргумента. Выбрасывание исключения Error внутри этих методов или ручное генерирование события 'error' приводит к неопределенному поведению.
Если поток Readable подключается к потоку Writable при генерации ошибки в потоке Writable, поток Readable будет отключен от подключаемого потока.
const { Writable } = require('stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
}
}); Пример потока записи
Следующий пример демонстрирует достаточно упрощенную (и в какой-то мере бесполезную) реализацию пользовательского потока записи. Хотя этот конкретный экземпляр потока записи не обладает какой-либо особой практической ценностью, пример иллюстрирует каждый из необходимых элементов пользовательского экземпляра потока Writable:
const { Writable } = require('stream');
class MyWritable extends Writable {
_write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
}
} Декодирование буферов в потоке записи
Декодирование буферов — распространенная задача, например, при использовании преобразователей, входной параметр которых — строка. Это непростая задача при использовании кодировок с многобайтовыми символами, таких как UTF-8. Следующий пример демонстрирует, как декодировать многобайтовые строки с использованием StringDecoder и Writable.
const { Writable } = require('stream');
const { StringDecoder } = require('string_decoder');
class StringWritable extends Writable {
constructor(options) {
super(options);
this._decoder = new StringDecoder(options && 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: € Реализация потока чтения
Класс stream.Readable расширяется для реализации потока Readable.
Пользовательские потоки Readable обязаны вызывать конструктор new stream.Readable([options]) и реализовывать метод readable._read().
new stream.Readable([options])
-
options<Объект>-
highWaterMark<число> Максимальное количество байтов для хранения во внутренeм буфере перед прекращением чтения из базового ресурса. По умолчанию:16384(16 КБ) или16для потоковobjectMode. -
encoding<строка> Если указано, буферы будут декодированы в строки с использованием указанного кодирования. По умолчанию:null. -
objectMode<логическое_значение> Указывает, должен ли этот поток вести себя как поток объектов. Это означает, чтоstream.read(n)возвращает единственное значение вместоBufferразмераn. По умолчанию:false. -
emitClose<логическое_значение> Указывает, должен ли поток генерировать'close'после уничтожения. По умолчанию:true. -
read<Функция> Реализация методаstream._read(). -
destroy<Функция> Реализация методаstream._destroy(). -
construct<Функция> Реализация методаstream._construct(). -
autoDestroy<логическое_значение> Указывает, должен ли этот поток автоматически вызывать.destroy()на себе после завершения. По умолчанию:true. -
signal<AbortSignal> Сигнал, представляющий возможность отмены.
-
const { Readable } = require('stream');
class MyReadable extends Readable {
constructor(options) {
// Calls the stream.Readable(options) constructor.
super(options);
// ...
}
} Или, при использовании конструкторов в стиле до ES6:
const { Readable } = require('stream');
const util = require('util');
function MyReadable(options) {
if (!(this instanceof MyReadable))
return new MyReadable(options);
Readable.call(this, options);
}
util.inherits(MyReadable, Readable); Или, используя упрощенный подход к конструктору:
const { Readable } = require('stream');
const myReadable = new Readable({
read(size) {
// ...
}
}); Вызов abort на AbortController соответствующем переданному AbortSignal будет вести себя так же, как вызов .destroy(new AbortError()) на созданном readable.
const { Readable } = require('stream');
const controller = new AbortController();
const read = new Readable({
read(size) {
// ...
},
signal: controller.signal
});
// Later, abort the operation closing the stream
controller.abort();
readable._construct(callback)
-
callback<Функция> Вызовите эту функцию (по желанию с аргументом ошибки) при завершении инициализации потока.
Метод _construct() НЕ ДОЛЖЕН вызываться напрямую. Он может быть реализован дочерними классами и, если это так, будет вызываться только внутренними методами класса Readable.
Эта необязательная функция будет запланирована на следующий тик потоком-конструктором, откладывая любые вызовы _read() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной инициализации ресурсов до использования потока.
const { Readable } = require('stream');
const fs = require('fs');
class ReadStream extends Readable {
constructor(filename) {
super();
this.filename = filename;
this.fd = null;
}
_construct(callback) {
fs.open(this.filename, (err, fd) => {
if (err) {
callback(err);
} else {
this.fd = fd;
callback();
}
});
}
_read(n) {
const buf = Buffer.alloc(n);
fs.read(this.fd, buf, 0, n, null, (err, bytesRead) => {
if (err) {
this.destroy(err);
} else {
this.push(bytesRead > 0 ? buf.slice(0, bytesRead) : null);
}
});
}
_destroy(err, callback) {
if (this.fd) {
fs.close(this.fd, (er) => callback(er || err));
} else {
callback(err);
}
}
}
readable._read(size)
-
size<число> Количество байтов для асинхронного чтения
Эта функция НЕ ДОЛЖНА вызываться кодом приложения напрямую. Она должна быть реализована дочерними классами и вызываться только внутренними методами класса 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<Ошибка> Возможная ошибка. -
callback<Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.
Метод _destroy() вызывается методом readable.destroy(). Он может быть переопределен дочерними классами, но не должен вызываться напрямую.
readable.push(chunk[, encoding])
-
chunk<Буфер> | <Uint8Array> | <строка> | <null> | <любой> Чанк данных для помещения в очередь чтения. Для потоков, не работающих в режиме объектов,chunkдолжен быть строкой,BufferилиUint8Array. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript. -
encoding<строка> Кодировка строковых чанков. Должна быть допустимой кодировкойBuffer, такой как'utf8'или'ascii'. - Возвращает: <логическое_значение>
trueесли дополнительные чанки данных могут продолжить помещаться;falseв противном случае.
Когда chunk является Buffer, Uint8Array или 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();
}
} Метод readable.push() используется для помещения содержимого во внутренний буфер. Он может быть вызван методом readable._read().
Для потоков, не работающих в режиме объектов, если параметр chunk метода readable.push() равен undefined, он будет обработан как пустая строка или буфер. См. readable.push('') для получения дополнительной информации.
Ошибки при чтении
Ошибки, возникающие во время обработки метода readable._read(), должны передаваться через метод readable.destroy(err). Бросание Error изнутри readable._read() или ручное отправление события 'error' приводит к неопределенному поведению.
const { Readable } = require('stream');
const myReadable = new Readable({
read(size) {
const err = checkSomeErrorCondition();
if (err) {
this.destroy(err);
} else {
// Do some work.
}
}
}); Пример потока подсчёта
Ниже приведен базовый пример потока Readable, который испускает цифры от 1 до 1 000 000 в порядке возрастания, а затем завершается.
const { Readable } = require('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);
}
}
} Реализация дуплексного потока
Поток Duplex — это поток, который реализует как Readable, так и Writable, такой как подключение сокета TCP.
Поскольку JavaScript не поддерживает множественное наследование, класс stream.Duplex расширяется для реализации потока 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<Объект> Передается как в конструкторWritable, так и в конструкторReadable. Также содержит следующие поля:-
allowHalfOpen<логическое> Если установлено в значениеfalse, то поток автоматически завершит сторону записи при завершении стороны чтения. По умолчанию:true. -
readable<логическое> Устанавливает, должен ли потокDuplexбыть читаемым. По умолчанию:true. -
writable<логическое> Устанавливает, должен ли потокDuplexбыть записываемым. По умолчанию:true. -
readableObjectMode<логическое> УстанавливаетobjectModeдля стороны чтения потока. Не имеет эффекта, еслиobjectModeравноtrue. По умолчанию:false. -
writableObjectMode<логическое> УстанавливаетobjectModeдля стороны записи потока. Не имеет эффекта, еслиobjectModeравноtrue. По умолчанию:false. -
readableHighWaterMark<число> УстанавливаетhighWaterMarkдля стороны чтения потока. Не имеет эффекта, если указаноhighWaterMark. -
writableHighWaterMark<число> УстанавливаетhighWaterMarkдля стороны записи потока. Не имеет эффекта, если указаноhighWaterMark.
-
const { Duplex } = require('stream');
class MyDuplex extends Duplex {
constructor(options) {
super(options);
// ...
}
} Или при использовании конструкторов в стиле до ES6:
const { Duplex } = require('stream');
const util = require('util');
function MyDuplex(options) {
if (!(this instanceof MyDuplex))
return new MyDuplex(options);
Duplex.call(this, options);
}
util.inherits(MyDuplex, Duplex); Или при использовании упрощенного подхода к конструктору:
const { Duplex } = require('stream');
const myDuplex = new Duplex({
read(size) {
// ...
},
write(chunk, encoding, callback) {
// ...
}
}); При использовании конвейера:
const { Transform, pipeline } = require('stream');
const fs = require('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);
} catch (err) {
callback(err);
}
}
}),
fs.createWriteStream('valid-object.json'),
(err) => {
if (err) {
console.error('failed', err);
} else {
console.log('completed');
}
}
); Пример дуплексного потока
Следующий пример демонстрирует простой пример потока Duplex, который оборачивает гипотетический объект нижнего уровня, в который можно записывать данные и из которого можно читать данные, хотя используемый API несовместим с потоками Node.js. Следующий пример демонстрирует простой пример потока Duplex который буферизует входящие данные через интерфейс Writable, которые затем считываются через интерфейс Readable.
const { Duplex } = require('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));
});
}
} Самая важная особенность потока Duplex заключается в том, что стороны Readable и Writable работают независимо друг от друга, несмотря на совместное существование в одном экземпляре объекта.
Потоки дуплексного режима объекта
Для потоков Duplex режим объекта objectMode можно установить только для стороны Readable или Writable с помощью опций readableObjectMode и writableObjectMode соответственно.
Например, в следующем примере создается новый поток Transform (который является типом потока Duplex), у которого есть сторона Writable в режиме объекта, принимающая JavaScript-числа, которые преобразуются в шестнадцатеричные строки на стороне Readable.
const { Transform } = require('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 Реализация потока преобразования
Поток Transform — это поток Duplex, где выходные данные вычисляются каким-либо образом из входных данных. Примерами являются потоки zlib или crypto, которые сжимают, шифруют или дешифруют данные.
Нет требования, чтобы размер выходных данных был таким же, как размер входных данных, число фрагментов или время их поступления совпадали. Например, поток Hash будет иметь только один фрагмент выходных данных, который предоставляется при завершении входных данных. Поток zlib будет генерировать выходные данные, которые либо намного меньше, либо намного больше, чем входные данные.
Класс stream.Transform расширяется для реализации потока Transform.
Класс stream.Transform прототипически наследуется от stream.Duplex и реализует свои версии методов writable._write() и readable._read(). Реализации пользовательских потоков Transform обязательно должны реализовывать метод transform._transform(), а также могут реализовать метод transform._flush().
При использовании потоков Transform необходимо быть внимательным, так как данные, записанные в поток, могут привести к приостановке стороны Writable потока, если выходные данные на стороне Readable не потребляются.
new stream.Transform([options])
-
options<Объект> Передается как в конструкторWritable, так и в конструкторReadable. Также содержит следующие поля:-
transform<Функция> Реализация методаstream._transform(). -
flush<Функция> Реализация методаstream._flush().
-
const { Transform } = require('stream');
class MyTransform extends Transform {
constructor(options) {
super(options);
// ...
}
} Или при использовании конструкторов в стиле до ES6:
const { Transform } = require('stream');
const util = require('util');
function MyTransform(options) {
if (!(this instanceof MyTransform))
return new MyTransform(options);
Transform.call(this, options);
}
util.inherits(MyTransform, Transform); Или при использовании упрощенного подхода к конструктору:
const { Transform } = require('stream');
const myTransform = new Transform({
transform(chunk, encoding, callback) {
// ...
}
}); Событие: 'end'
Событие 'end' испускается классом stream.Readable. Событие 'end' испускается после вывода всех данных, что происходит после вызова обратного вызова в transform._flush(). В случае ошибки, 'end' не должно испускаться.
Событие: 'finish'
Событие 'finish' испускается классом stream.Writable. Событие 'finish' испускается после вызова stream.end() и обработки всех фрагментов методом stream._transform(). В случае ошибки, 'finish' не должно испускаться.
transform._flush(callback)
-
callback<Функция> Функция обратного вызова (необязательно с аргументом ошибки и данными), которая вызывается при сбросе оставшихся данных.
Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только внутренними методами класса Readable.
В некоторых случаях для операции преобразования может потребоваться вывести дополнительный фрагмент данных в конце потока. Например, поток сжатия zlib будет сохранять некоторое количество внутреннего состояния, используемого для оптимального сжатия выходных данных. Однако, когда поток завершается, эта дополнительная информация должна быть сброшена, чтобы сжатые данные были полными.
Пользовательские реализации Transform могут реализовать метод transform._flush(). Он будет вызываться, когда не осталось записываемых данных, но перед тем, как будет испущено событие 'end', сигнализирующее о завершении потока Readable.
В реализации метода transform._flush(), метод transform.push() может вызываться ноль или более раз, по мере необходимости. Функция callback должна вызываться при завершении операции сброса.
Метод transform._flush() имеет префикс с подчеркиванием, так как он является внутренним для определяющего его класса и никогда не должен вызываться непосредственно пользовательскими программами.
transform._transform(chunk, encoding, callback)
-
chunk<Buffer> | <строка> | <любой> ПередаваемыйBufferдля преобразования, конвертируемый изstring, переданного вstream.write(). Если опция потокаdecodeStringsравнаfalseили поток работает в режиме объектов, фрагмент не будет преобразован и будет таким, каким был передан вstream.write(). -
encoding<строка> Если фрагмент — строка, то это тип кодировки. Если фрагмент — буфер, то это специальное значение'buffer'. В этом случае игнорируйте его. -
callback<Функция> Функция обратного вызова (необязательно с аргументом ошибки и данными), которая будет вызвана после обработки переданного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);
}; Метод 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);
}
})(); Асинхронные итераторы регистрируют постоянный обработчик ошибок в потоке, чтобы предотвратить любые необработанные ошибки после уничтожения.
Создание читаемых потоков с асинхронными генераторами
Читаемый поток Node.js может быть создан из асинхронного генератора с помощью утилитарного метода Readable.from():
const { Readable } = require('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);
}); Перенаправление в потоки записи из асинхронных итераторов
При записи в поток записи из асинхронного итератора обеспечьте правильную обработку обратной задержки и ошибок. stream.pipeline() абстрагирует обработку обратной задержки и ошибок, связанных с обратной задержкой:
const fs = require('fs');
const { pipeline } = require('stream');
const { pipeline: pipelinePromise } = require('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();
}); Совместимость со старыми версиями 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); До 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); Помимо новых потоков 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('') не рекомендуется.
Добавление нулевой строки, Buffer или Uint8Array в поток, который не находится в режиме объектов, имеет интересный побочный эффект. Поскольку это вызов 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-v16.x/docs/api/stream.html