Поток[src]
Исходный код: lib/stream.js
Поток — это абстрактный интерфейс для работы со потоковыми данными в Node.js. Модуль node:stream предоставляет API для реализации интерфейса потока.
Node.js предоставляет множество объектов потока. Например, запрос к HTTP-серверу и process.stdout — оба являются экземплярами потоков.
Потоки могут быть читаемыми, записываемыми или и теми, и другими. Все потоки являются экземплярами EventEmitter.
Для доступа к модулю node:stream:
const stream = require('node:stream'); copy Модуль node:stream полезен для создания новых типов экземпляров потоков. Обычно нет необходимости использовать модуль node:stream для потребления потоков.
Структура данного документа
Данный документ содержит две основные секции и третью секцию для заметок. Первая секция объясняет, как использовать существующие потоки в приложении. Вторая секция объясняет, как создавать новые типы потоков.
Типы потоков
В Node.js существует четыре основных типа потоков:
-
Writable: потоки, в которые можно записывать данные (например,fs.createWriteStream()). -
Readable: потоки, из которых можно читать данные (например,fs.createReadStream()). -
Duplex: потоки, которые являются одновременноReadableиWritable(например,net.Socket). -
Transform:Duplexпотоки, которые могут изменять или преобразовывать данные при записи и чтении (например,zlib.createDeflate()).
Кроме того, этот модуль включает вспомогательные функции stream.pipeline(), stream.finished(), stream.Readable.from() и stream.addAbortSignal().
API потоков с обещаниями
API с обещаниями предоставляет альтернативный набор асинхронных вспомогательных функций для потоков, которые возвращают объекты обещаний вместо использования обратных вызовов. К API можно получить доступ через require('node:stream/promises') или require('node:stream').promises.
stream.pipeline(source[, ...transforms], destination[, options])
stream.pipeline(streams[, options])
-
streams<Поток[]> | <Итерируемый массив[]> | <Асинхронно итерируемый массив[]> | <Функция[]> -
source<Поток> | <Итерируемый> | <Асинхронно итерируемый> | <Функция>- Возвращает: <Обещание> | <Асинхронно итерируемый>
-
...transforms<Поток> | <Функция>-
source<Асинхронно итерируемый> - Возвращает: <Обещание> | <Асинхронно итерируемый>
-
-
destination<Поток> | <Функция>-
source<Асинхронно итерируемый> - Возвращает: <Обещание> | <Асинхронно итерируемый>
-
-
options<Объект> Параметры конвейера-
signal<Объект отмены> -
end<логическое> Завершить целевой поток при завершении исходного. Потоки преобразования всегда завершаются, даже если это значениеfalse. По умолчанию:true.
-
- Возвращает: <Обещание> Выполняется, когда конвейер завершен.
Модули CJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
const zlib = require('node:zlib');
async function run() {
await pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz'),
);
console.log('Pipeline succeeded.');
}
run().catch(console.error);
Модули MJS
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';
await pipeline(
createReadStream('archive.tar'),
createGzip(),
createWriteStream('archive.tar.gz'),
);
console.log('Pipeline succeeded.'); Чтобы использовать объект AbortSignal, передайте его внутри объекта options в качестве последнего аргумента. Когда сигнал отменён, destroy вызывается на базовом конвейере с AbortError.
Модули CJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
const zlib = require('node:zlib');
async function run() {
const ac = new AbortController();
const signal = ac.signal;
setImmediate(() => ac.abort());
await pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz'),
{ signal },
);
}
run().catch(console.error); // AbortError
Модули MJS
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';
const ac = new AbortController();
const { signal } = ac;
setImmediate(() => ac.abort());
try {
await pipeline(
createReadStream('archive.tar'),
createGzip(),
createWriteStream('archive.tar.gz'),
{ signal },
);
} catch (err) {
console.error(err); // AbortError
} API потоков также поддерживает асинхронные генераторы:
Модули CJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
async function run() {
await pipeline(
fs.createReadStream('lowercase.txt'),
async function* (source, { signal }) {
source.setEncoding('utf8'); // Work with strings rather than `Buffer`s.
for await (const chunk of source) {
yield await processChunk(chunk, { signal });
}
},
fs.createWriteStream('uppercase.txt'),
);
console.log('Pipeline succeeded.');
}
run().catch(console.error);
Модули MJS
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
await pipeline(
createReadStream('lowercase.txt'),
async function* (source, { signal }) {
source.setEncoding('utf8'); // Work with strings rather than `Buffer`s.
for await (const chunk of source) {
yield await processChunk(chunk, { signal });
}
},
createWriteStream('uppercase.txt'),
);
console.log('Pipeline succeeded.'); Не забудьте обработать аргумент signal переданный в асинхронный генератор. Особенно в случае, когда асинхронный генератор является источником для конвейера (т.е. первый аргумент) или конвейер никогда не завершится.
Модули CJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
async function run() {
await pipeline(
async function* ({ signal }) {
await someLongRunningfn({ signal });
yield 'asd';
},
fs.createWriteStream('uppercase.txt'),
);
console.log('Pipeline succeeded.');
}
run().catch(console.error);
Модули MJS
import { pipeline } from 'node:stream/promises';
import fs from 'node:fs';
await pipeline(
async function* ({ signal }) {
await someLongRunningfn({ signal });
yield 'asd';
},
fs.createWriteStream('uppercase.txt'),
);
console.log('Pipeline succeeded.'); API потоков предоставляет версию с обратными вызовами:
stream.finished(stream[, options])
-
stream<Поток> | <Поток чтения> | <Поток записи> Поток чтения и/или записи/веб-поток. -
options<Объект>-
error<логическое> | <неопределённо> -
readable<логическое> | <неопределённо> -
writable<логическое> | <неопределённо> -
signal: <Объект отмены> | <неопределённо>
-
- Возвращает: <Обещание> Выполняется, когда поток больше не читается или не записывается.
Модули CJS
const { finished } = require('node:stream/promises');
const fs = require('node:fs');
const rs = fs.createReadStream('archive.tar');
async function run() {
await finished(rs);
console.log('Stream is done reading.');
}
run().catch(console.error);
rs.resume(); // Drain the stream.
Модули MJS
import { finished } from 'node:stream/promises';
import { createReadStream } from 'node:fs';
const rs = createReadStream('archive.tar');
async function run() {
await finished(rs);
console.log('Stream is done reading.');
}
run().catch(console.error);
rs.resume(); // Drain the stream. API потоков также предоставляет версию с обратными вызовами.
Режим работы с объектами
Все потоки, созданные API Node.js, работают исключительно со строками, <Буфер>, <Массивом типов> и <Представлением данных>:
-
StringsиBuffers— наиболее распространённые типы, используемые с потоками. -
TypedArrayиDataViewпозволяют обрабатывать двоичные данные с типами, такими какInt32ArrayилиUint8Array. При записи массива типов или представления данных в поток Node.js обрабатывает исходные байты.
Однако реализация потоков может работать с другими типами JavaScript-значений (за исключением null, который служит специальной цели в потоках). Такие потоки считаются работающими в "режиме работы с объектами".
Экземпляры потоков переключаются в режим работы с объектами с помощью параметра objectMode при создании потока. Попытка переключить существующий поток в режим работы с объектами небезопасна.
Буферизация
Оба потока Writable и Readable будут хранить данные во внутреннем буфере.
Объём потенциально буферизованных данных зависит от параметра highWaterMark, переданного в конструктор потока. Для обычных потоков параметр highWaterMark определяет общее количество байтов. Для потоков, работающих в режиме работы с объектами, параметр highWaterMark определяет общее количество объектов. Для потоков, работающих со строками (но не декодирующих их), параметр highWaterMark определяет общее количество единиц кода UTF-16.
Данные буферизуются в Readable потоках, когда реализация вызывает stream.push(chunk). Если потребитель потока не вызывает stream.read(), данные будут храниться во внутренней очереди до момента их потребления.
Когда общий размер внутреннего буфера чтения достигает порога, заданного highWaterMark, поток временно прекращает чтение данных из базового ресурса до тех пор, пока буферизованные данные не будут потреблены (то есть, поток перестанет вызывать внутренний метод readable._read(), используемый для заполнения буфера чтения).
Данные буферизуются в Writable потоках при многократном вызове метода writable.write(chunk). Пока общий размер внутреннего буфера записи ниже порога, заданного highWaterMark, вызовы writable.write() будут возвращать true. После того, как размер внутреннего буфера достигнет или превысит highWaterMark, будет возвращено значение false.
Ключевой целью API stream, особенно метода stream.pipe(), является ограничение буферизации данных приемлемыми уровнями, чтобы источники и приемники с различной скоростью не перегружали доступную память.
Опция highWaterMark — это порог, а не ограничение: она определяет объем данных, который поток буферизует перед тем, как перестать запрашивать дополнительные данные. Она не накладывает строгих ограничений на использование памяти в целом. Конкретные реализации потоков могут выбрать применение более строгих ограничений, но это необязательно.
Так как потоки Duplex и Transform являются одновременно Readable и Writable, каждый из них поддерживает два отдельных внутренних буфера для чтения и записи, позволяя каждой стороне работать независимо от другой, поддерживая надлежащий и эффективный поток данных. Например, экземпляры net.Socket являются потоками Duplex, чья сторона Readable позволяет потреблять данные, полученные из сокета, а чья сторона Writable позволяет записывать данные в сокет. Поскольку данные могут записываться в сокет быстрее или медленнее, чем они получаются, каждая сторона должна работать (и буферизовать данные) независимо от другой.
Механизм внутренней буферизации является внутренней реализацией и может быть изменён в любое время. Однако для некоторых расширенных реализаций внутренние буферы можно получить, используя writable.writableBuffer или readable.readableBuffer. Использование этих недокументированных свойств не рекомендуется.
API для потребителей потоков
Практически все приложения Node.js, независимо от сложности, используют потоки каким-то образом. Ниже приведен пример использования потоков в приложении Node.js, реализующем HTTP-сервер:
const http = require('node:http');
const server = http.createServer((req, res) => {
// `req` is an http.IncomingMessage, which is a readable stream.
// `res` is an http.ServerResponse, which is a writable stream.
let body = '';
// Get the data as utf8 strings.
// If an encoding is not set, Buffer objects will be received.
req.setEncoding('utf8');
// Readable streams emit 'data' events once a listener is added.
req.on('data', (chunk) => {
body += chunk;
});
// The 'end' event indicates that the entire body has been received.
req.on('end', () => {
try {
const data = JSON.parse(body);
// Write back something interesting to the user:
res.write(typeof data);
res.end();
} catch (er) {
// uh oh! bad json!
res.statusCode = 400;
return res.end(`error: ${er.message}`);
}
});
});
server.listen(1337);
// $ curl localhost:1337 -d "{}"
// object
// $ curl localhost:1337 -d "\"foo\""
// string
// $ curl localhost:1337 -d "not json"
// error: Unexpected token 'o', "not json" is not valid JSON copy Writable потоки (например, res в примере) предоставляют методы, такие как write() и end(), которые используются для записи данных в поток.
Readable потоки используют API EventEmitter для уведомления кода приложения, когда данные доступны для чтения из потока. Эти данные можно прочитать из потока несколькими способами.
Оба Writable и Readable потока используют API EventEmitter различными способами для передачи текущего состояния потока.
Duplex и Transform потоки являются одновременно Writable и Readable.
Приложения, которые либо записывают данные в поток, либо потребляют данные из потока, не обязаны реализовывать интерфейсы потоков напрямую и, как правило, не имеют причин вызывать require('node:stream').
Разработчики, желающие реализовать новые типы потоков, должны обратиться к разделу API для разработчиков потоков.
Потоки записи
Потоки записи — это абстракция назначения, в которое записываются данные.
Примеры Writable потоков включают:
- HTTP-запросы (на стороне клиента)
- HTTP-ответы (на стороне сервера)
- Потоки записи в файлы (fs)
- zlib-потоки
- crypto-потоки
- TCP-сокеты
- стандартный ввод процесса-потомка
-
process.stdout,process.stderr
Некоторые из этих примеров фактически являются Duplex потоками, которые реализуют интерфейс Writable.
Все Writable потоки реализуют интерфейс, определённый классом stream.Writable.
Хотя конкретные экземпляры Writable потоков могут отличаться различными способами, все Writable потоки следуют одному фундаментальному шаблону использования, как показано в примере ниже:
const myStream = getWritableStreamSomehow();
myStream.write('some data');
myStream.write('some more data');
myStream.end('done writing data'); copy Класс: stream.Writable
Событие: 'close'
Событие 'close' излучается, когда поток и все его подчинённые ресурсы (например, дескриптор файла) были закрыты. Это событие указывает на то, что больше событий не будет излучаться и дальнейшие вычисления не будут происходить.
Writable поток всегда излучит событие 'close', если он был создан с параметром emitClose.
Событие: 'drain'
Если вызов stream.write(chunk) возвращает false, событие 'drain' будет излучено, когда будет уместно возобновить запись данных в поток.
// Write the data to the supplied writable stream one million times.
// Be attentive to back-pressure.
function writeOneMillionTimes(writer, data, encoding, callback) {
let i = 1000000;
write();
function write() {
let ok = true;
do {
i--;
if (i === 0) {
// Last time!
writer.write(data, encoding, callback);
} else {
// See if we should continue, or wait.
// Don't pass the callback, because we're not done yet.
ok = writer.write(data, encoding);
}
} while (i > 0 && ok);
if (i > 0) {
// Had to stop early!
// Write some more once it drains.
writer.once('drain', write);
}
}
} copy Событие: 'error'
Событие 'error' излучается, если при записи или передаче данных возникла ошибка. Обработчик вызывается с одним аргументом типа Error.
Поток закрывается, когда излучается событие 'error', если параметр autoDestroy не был установлен в значение false при создании потока.
После 'error', больше никаких событий, кроме 'close' не должно излучаться (включая события 'error').
Событие: 'finish'
Событие 'finish' излучается после вызова метода stream.end() и после того, как все данные были отправлены в подчинённую систему.
const writer = getWritableStreamSomehow();
for (let i = 0; i < 100; i++) {
writer.write(`hello, #${i}!\n`);
}
writer.on('finish', () => {
console.log('All writes are now complete.');
});
writer.end('This is the end\n'); copy Событие: 'pipe'
-
srcИсточник-поток <stream.Readable>, который направляет данные в этот поток записи
Событие 'pipe' излучается, когда метод stream.pipe() вызывается на потоке чтения, добавляя этот поток записи в список его назначений.
const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('pipe', (src) => {
console.log('Something is piping into the writer.');
assert.equal(src, reader);
});
reader.pipe(writer); copy Событие: 'unpipe'
-
srcИсточник-поток <stream.Readable>, который отменил перенаправление в этот поток записи
Событие 'unpipe' излучается, когда метод stream.unpipe() вызывается на потоке Readable, удаляя этот поток записи из списка его назначений.
Также излучается, если этот поток записи генерирует ошибку, когда в него перенаправляется поток Readable.
const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
console.log('Something has stopped piping into the writer.');
assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer); copy
writable.cork()
Метод writable.cork() принудительно буферирует все записанные данные в памяти. Буферизированные данные будут отправлены, когда будут вызваны методы stream.uncork() или stream.end().
Основное назначение writable.cork() — справиться со случаем, когда несколько небольших фрагментов записываются в поток в быстрой последовательности. Вместо немедленной передачи их в подчинённое назначение, writable.cork() буферирует все фрагменты, пока не будет вызвано writable.uncork(), что передаст их все в writable._writev(), если оно есть. Это предотвращает ситуацию блокировки очереди, когда данные буферизуются, ожидая обработки первого небольшого фрагмента. Однако использование writable.cork() без реализации writable._writev() может отрицательно сказаться на пропускной способности.
См. также: writable.uncork(), writable._writev().
writable.destroy([error])
Уничтожает поток. Необязательно излучает событие 'error', и излучает событие 'close' (если emitClose не установлено в значение false). После этого вызова поток записи завершён, и последующие вызовы write() или end() приведут к ошибке ERR_STREAM_DESTROYED. Это разрушительный и немедленный способ уничтожения потока. Предыдущие вызовы write() могут не быть обработанными, и могут вызвать ошибку ERR_STREAM_DESTROYED. Используйте end() вместо destroy, если данные должны быть очищены перед закрытием, или дождитесь события 'drain' перед уничтожением потока.
const { Writable } = require('node:stream');
const myStream = new Writable();
const fooErr = new Error('foo error');
myStream.destroy(fooErr);
myStream.on('error', (fooErr) => console.error(fooErr.message)); // foo error copy const { Writable } = require('node:stream');
const myStream = new Writable();
myStream.destroy();
myStream.on('error', function wontHappen() {}); copy const { Writable } = require('node:stream');
const myStream = new Writable();
myStream.destroy();
myStream.write('foo', (error) => console.error(error.code));
// ERR_STREAM_DESTROYED copy После вызова destroy() все последующие вызовы будут нет-операциями, и больше никаких ошибок, кроме _destroy(), не будут излучаться как 'error'.
Реализаторы не должны переопределять этот метод, а вместо этого реализовать writable._destroy().
writable.closed
Является ли поток true после того, как излучено событие 'close'.
writable.destroyed
Является ли поток true после вызова writable.destroy().
const { Writable } = require('node:stream');
const myStream = new Writable();
console.log(myStream.destroyed); // false
myStream.destroy();
console.log(myStream.destroyed); // true copy
writable.end([chunk[, encoding]][, callback])
-
chunk<строка> | <Буфер> | <Массив типов> | <DataView> | <любое> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжен быть <строкой>, <Буфером>, <Массивом типов> или <DataView>. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript, кромеnull. -
encoding<строка> Кодировка, еслиchunkявляется строкой -
callback<Функция> Обратный вызов, когда поток завершен. - Возвращает: <this>
Вызов метода writable.end() сигнализирует о том, что больше данных не будут записаны в Writable. Дополнительные аргументы chunk и encoding позволяют записать последний фрагмент данных непосредственно перед закрытием потока.
Вызов метода stream.write() после вызова stream.end() вызовет ошибку.
// Write 'hello, ' and then end with 'world!'.
const fs = require('node:fs');
const file = fs.createWriteStream('example.txt');
file.write('hello, ');
file.end('world!');
// Writing more now is not allowed! copy
writable.setDefaultEncoding(encoding)
Метод writable.setDefaultEncoding() устанавливает кодировку по умолчанию для потока Writable.
writable.uncork()
Метод writable.uncork() сбрасывает все данные, буферизованные с момента вызова stream.cork().
При использовании writable.cork() и writable.uncork() для управления буферизацией записей в поток, откладывайте вызовы writable.uncork() с помощью process.nextTick(). Это позволяет объединять все вызовы writable.write(), которые происходят в рамках фазы цикла событий Node.js.
stream.cork();
stream.write('some ');
stream.write('data ');
process.nextTick(() => stream.uncork()); copy Если метод writable.cork() вызывается несколько раз для потока, то необходимо вызвать такое же количество раз метод writable.uncork() для сброса буферизованных данных.
stream.cork();
stream.write('some ');
stream.cork();
stream.write('data ');
process.nextTick(() => {
stream.uncork();
// The data will not be flushed until uncork() is called a second time.
stream.uncork();
}); copy См. также: writable.cork().
writable.writable
Истина, если безопасно вызвать writable.write(), что означает, что поток не был уничтожен, не возникла ошибка и не было завершено.
writable.writableAborted
Возвращает значение true, если поток был уничтожен или произошла ошибка до эмиссии 'finish'.
writable.writableEnded
Истина после вызова writable.end(). Это свойство не указывает, были ли данные сброшены, для этого используйте writable.writableFinished.
writable.writableCorked
Число вызовов writable.uncork(), необходимых для полного разблокирования потока.
writable.errored
Возвращает ошибку, если поток был уничтожен с ошибкой.
writable.writableFinished
Устанавливается в true незадолго до эмиссии события 'finish'.
writable.writableHighWaterMark
Возвращает значение highWaterMark , переданное при создании этого Writable.
writable.writableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение предоставляет данные для интроспекции относительно состояния highWaterMark.
writable.writableNeedDrain
Истина, если буфер потока заполнен, и поток выведет 'drain'.
writable.writableObjectMode
Получатель свойства objectMode данного потока Writable.
writable.write(chunk[, encoding][, callback])
-
chunk<строка> | <Буфер> | <TypedArray> | <DataView> | <любое> Данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжно быть строкой <string>, буфером <Buffer>, <TypedArray> или <DataView>. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript, кромеnull. -
encoding<строка> | <null> Кодировка, еслиchunk— строка. По умолчанию:'utf8' -
callback<Функция> Обратный вызов при сбросе этого фрагмента данных. - Возвращает: <булево>
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(), можно соблюдать обратную загрузку и избегать проблем с памятью, используя событие 'drain':
function write(data, cb) {
if (!stream.write(data)) {
stream.once('drain', cb);
} else {
process.nextTick(cb);
}
}
// Wait for cb to be called before doing any other write.
write('hello', () => {
console.log('Write completed, do more writes now.');
}); copy Поток Writable в режиме объектов всегда игнорирует аргумент encoding.
Потоковые потоки для чтения
Потоки для чтения — абстракция для источника, из которого потребляются данные.
Примеры потоков для чтения включают:
- Ответы HTTP на стороне клиента
- Запросы HTTP на стороне сервера
- Потоки чтения fs
- Потоки zlib
- Потоки crypto
- Сокеты TCP
- Стандартный вывод и стандартная ошибка процесса дочернего процесса
process.stdin
Все потоки Readable реализуют интерфейс, определенный классом stream.Readable.
Два режима чтения
Потоки для чтения фактически работают в одном из двух режимов: активный и приостановленный. Эти режимы отделены от режима объектов. Поток Readable может быть в режиме объектов или нет, независимо от того, находится ли он в активном или приостановленном режиме.
-
В активном режиме данные считываются из подсистемы автоматически и предоставляются приложению как можно быстрее с помощью событий через интерфейс
EventEmitter. -
В режиме приостановки метод
stream.read()должен вызываться явно для чтения фрагментов данных из потока.
Все потоки Readable начинаются в режиме приостановки, но могут быть переведены в режим активного потока одним из следующих способов:
- Добавление обработчика события
'data'. - Вызов метода
stream.resume(). - Вызов метода
stream.pipe()для передачи данных вWritable.
Поток может вернуться в приостановленный режим следующим образом:
- Если нет целевых потоков, вызовом метода
stream.pause(). - Если есть целевые потоки, удаляя все целевые потоки. Несколько целевых потоков можно удалить, вызвав метод
stream.unpipe().
Важный момент: поток для чтения не будет генерировать данные, пока не будет предоставлен механизм для потребления или игнорирования этих данных. Если механизм потребления отключен или удалён, поток для чтения попытается остановить генерацию данных.
По соображениям обратной совместимости, удаление обработчиков событий 'data' не автоматически приостановит поток. Также, если есть целевые потоки, то вызов stream.pause() не гарантирует, что поток останется приостановленным, после того, как эти целевые потоки обработают данные и запросят новые.
Если поток Readable переведён в режим активного потока, и нет потребителей для обработки данных, эти данные будут потеряны. Это может произойти, например, когда метод readable.resume() вызывается без обработчика, привязанного к событию 'data', или когда обработчик события 'data' удаляется из потока.
Добавление обработчика события 'readable' автоматически останавливает активный поток, и данные необходимо потреблять с помощью readable.read(). Если обработчик события 'readable' удаляется, поток снова начнёт активный поток, если есть обработчик события 'data'.
Три состояния
Два режима работы потока для чтения представляют собой упрощённую абстракцию более сложного внутреннего управления состоянием, происходящего внутри реализации потока для чтения.
Конкретно, в любой момент времени каждый поток для чтения находится в одном из трёх возможных состояний:
readable.readableFlowing === nullreadable.readableFlowing === falsereadable.readableFlowing === true
Когда readable.readableFlowing находится в null, нет механизма для потребления данных потока. Поэтому поток не будет генерировать данные. В этом состоянии подключение обработчика для события 'data', вызов метода readable.pipe() или метода readable.resume() переведут readable.readableFlowing в true, вызвав Readable начать активную отправку событий по мере генерации данных.
Вызов readable.pause(), readable.unpipe() или получение обратной загрузки приведут к установке readable.readableFlowing как false, временно останавливая передачу событий, но не останавливая генерацию данных. В этом состоянии подключение обработчика для события 'data' не переведет readable.readableFlowing в true.
const { PassThrough, Writable } = require('node:stream');
const pass = new PassThrough();
const writable = new Writable();
pass.pipe(writable);
pass.unpipe(writable);
// readableFlowing is now false.
pass.on('data', (chunk) => { console.log(chunk.toString()); });
// readableFlowing is still false.
pass.write('ok'); // Will not emit 'data'.
pass.resume(); // Must be called to make stream emit 'data'.
// readableFlowing is now true. copy Пока readable.readableFlowing находится в false, данные могут накапливаться во внутреннем буфере потока.
Выбор одного стиля API
API потока для чтения эволюционировал в течение нескольких версий Node.js и предоставляет несколько способов потребления данных потока. В целом, разработчики должны выбрать один метод потребления данных и никогда не использовать несколько методов для потребления данных из одного потока. В частности, использование комбинации on('data'), on('readable'), pipe() или асинхронных итераторов может привести к неинтуитивному поведению.
Класс: 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.`);
}); copy Событие: 'end'
Событие 'end' генерируется, когда больше нет данных для потребления из потока.
Событие 'end' не будет сгенерировано, пока данные не будут полностью использованы. Это можно сделать, переключив поток в режим потоковой передачи или вызывая stream.read() многократно, пока все данные не будут использованы.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
});
readable.on('end', () => {
console.log('There will be no more data.');
}); copy Событие: 'error'
Событие 'error' может быть сгенерировано реализацией 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()) !== null) {
console.log(data);
}
}); copy Если достигнут конец потока, вызов stream.read() вернет null и сгенерирует событие 'end'. Это также верно, если данных для чтения никогда не было. Например, в следующем примере foo.txt - это пустой файл:
const fs = require('node:fs');
const rr = fs.createReadStream('foo.txt');
rr.on('readable', () => {
console.log(`readable: ${rr.read()}`);
});
rr.on('end', () => {
console.log('end');
}); copy Вывод выполнения этого скрипта:
$ node test.js readable: null end copy
В некоторых случаях присоединение обработчика события 'readable' приведет к чтению некоторого количества данных во внутренний буфер.
В целом, механизмы событий readable.pipe() и 'data' легче понять, чем событие 'readable'. Однако обработка 'readable' может привести к увеличению пропускной способности.
Если одновременно используются события 'readable' и 'data', то 'readable' имеет приоритет в управлении потоком, т.е. событие 'data' будет сгенерировано только при вызове stream.read(). Свойство readableFlowing станет false. Если есть 'data' обработчика, когда 'readable' удаляется, поток начнет передачу, т.е. события 'data' будут генерироваться без вызова .resume().
Событие: 'resume'
Событие 'resume' генерируется при вызове stream.resume(), если readableFlowing не true.
readable.destroy([error])
-
error<Ошибка> Ошибка, которая будет передана как полезная нагрузка в событии'error' - Возвращает: <текущий объект>
Уничтожить поток. При необходимости сгенерировать событие 'error' и событие 'close' (если emitClose не установлено в false). После этого вызова поток readable высвободит все внутренние ресурсы, и последующие вызовы push() будут игнорироваться.
После вызова destroy() любые дальнейшие вызовы будут являться пустой операцией, и не будет генерироваться никаких других ошибок, кроме тех, которые могут быть сгенерированы _destroy() в качестве 'error'.
Реализаторы не должны переопределять этот метод, а вместо этого реализовывать readable._destroy().
readable.closed
Является true после того, как сгенерировано 'close'.
readable.destroyed
Является true после вызова readable.destroy().
readable.isPaused()
- Возвращает: <логическое значение>
Метод readable.isPaused() возвращает текущее состояние работы Readable. Он используется в первую очередь механизмом, лежащим в основе метода readable.pipe(). В большинстве типичных случаев нет причины использовать этот метод напрямую.
const readable = new stream.Readable(); readable.isPaused(); // === false readable.pause(); readable.isPaused(); // === true readable.resume(); readable.isPaused(); // === false copy
readable.pause()
- Возвращает: <текущий объект>
Метод readable.pause() заставляет поток в режиме потоковой передачи прекратить генерировать события 'data', выходя из режима потоковой передачи. Любые доступные данные останутся во внутреннем буфере.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
readable.pause();
console.log('There will be no additional data for 1 second.');
setTimeout(() => {
console.log('Now data will start flowing again.');
readable.resume();
}, 1000);
}); copy Метод readable.pause() не оказывает никакого влияния, если имеется обработчик события 'readable'.
readable.pipe(destination[, options])
-
destination<stream.Writable> Назначение для записи данных -
options<Объект> Параметры для соединения-
end<логическое значение> Завершить писатель, когда читатель завершит. По умолчанию:true.
-
- Возвращает: <stream.Writable> Назначение, позволяющее создавать цепочку соединений, если это поток
DuplexилиTransform
Метод readable.pipe() подключает поток Writable к readable, заставляя его автоматически переключиться в режим потоковой передачи и передать все свои данные присоединенному потоку Writable. Поток данных будет автоматически управляться таким образом, чтобы целевой поток Writable не перегружался более быстрым потоком Readable.
Следующий пример направляет все данные из readable в файл с именем file.txt:
const fs = require('node:fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt'.
readable.pipe(writable); copy Можно подключить несколько Writable потоков к одному Readable потоку.
Метод readable.pipe() возвращает ссылку на поток назначения, что позволяет создавать цепочки потоков с перенаправлением:
const fs = require('node:fs');
const zlib = require('node:zlib');
const r = fs.createReadStream('file.txt');
const z = zlib.createGzip();
const w = fs.createWriteStream('file.txt.gz');
r.pipe(z).pipe(w); copy По умолчанию, метод stream.end() вызывается в потоке назначения Writable при получении события 'end' в источнике Readable потоке, чтобы поток назначения перестал быть записываемым. Чтобы отключить это поведение по умолчанию, можно передать опцию end как false, что позволит оставить поток назначения открытым:
reader.pipe(writer, { end: false });
reader.on('end', () => {
writer.end('Goodbye\n');
}); copy Важно отметить, что если поток 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.');
}); copy Каждый вызов 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('');
}); copy Поток Readable в режиме объектов всегда возвращает единственный элемент при вызове readable.read(size), независимо от значения аргумента size.
Если метод readable.read() возвращает фрагмент данных, также будет излучено событие 'data'.
Вызов stream.read([size]) после излучения события 'end' вернёт null. Ошибка во время выполнения не произойдёт.
readable.readable
Истинно, если безопасно вызвать readable.read(), что означает, что поток не был уничтожен и не излучил события 'error' или 'end'.
readable.readableAborted
Возвращает, был ли поток уничтожен или произошла ошибка до излучения события 'end'.
readable.readableDidRead
Возвращает, было ли излучено событие 'data'.
readable.readableEncoding
Геттер для свойства encoding заданного потока Readable. Свойство encoding можно установить с помощью метода readable.setEncoding().
readable.readableEnded
Становится true при излучении события 'end'.
readable.errored
Возвращает ошибку, если поток был уничтожен с ошибкой.
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.');
}); copy Метод readable.resume() не имеет эффекта, если существует обработчик события 'readable'.
readable.setEncoding(encoding)
Метод readable.setEncoding() устанавливает кодировку символов для данных, считываемых из потока Readable.
По умолчанию кодировка не задана, и данные потока будут возвращаться как объекты Buffer. Установка кодировки приводит к тому, что данные потока возвращаются в виде строк указанной кодировки, а не как объекты Buffer. Например, вызов readable.setEncoding('utf8') приведет к тому, что выходные данные будут интерпретированы как данные UTF-8 и переданы как строки. Вызов readable.setEncoding('hex') приведет к кодированию данных в шестнадцатеричном формате.
Поток Readable будет правильно обрабатывать многобайтовые символы, передаваемые через поток, которые в противном случае были бы неправильно декодированы, если бы их просто извлекли из потока как Buffer объекты.
const readable = getReadableStreamSomehow();
readable.setEncoding('utf8');
readable.on('data', (chunk) => {
assert.equal(typeof chunk, 'string');
console.log('Got %d characters of string data:', chunk.length);
}); copy
readable.unpipe([destination])
-
destination<stream.Writable> Необязательный конкретный поток для отсоединения от канала - Возвращает: <this>
Метод readable.unpipe() отсоединяет поток Writable, ранее подключенный с помощью метода stream.pipe().
Если destination не указан, то все каналы отсоединяются.
Если destination указан, но для него нет подключенного канала, то метод ничего не делает.
const fs = require('node:fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt',
// but only for the first second.
readable.pipe(writable);
setTimeout(() => {
console.log('Stop writing to file.txt.');
readable.unpipe(writable);
console.log('Manually close the file stream.');
writable.end();
}, 1000); copy
readable.unshift(chunk[, encoding])
-
chunk<Буфер> | <Массив типизированных данных> | <DataView> | <строка> | <null> | <любое> Чанк данных для добавления в начало очереди чтения. Для потоков, не работающих в режиме объектов,chunkдолжен быть <строкой>, <Буфером>, <Массивом типизированных данных>, <DataView> илиnull. Для потоков в режиме объектовchunkможет быть любым значением JavaScript. -
encoding<строка> Кодировка чанков строк. Должна быть допустимой кодировкойBuffer, такой как'utf8'или'ascii'.
Передача chunk в качестве null сигнализирует об окончании потока (EOF) и ведет себя так же, как readable.push(null), после чего больше данных нельзя записать. Сигнал EOF помещается в конец буфера, и любые буферизованные данные все равно будут сброшены.
Метод readable.unshift() вставляет кусок данных обратно во внутренний буфер. Это полезно в определенных ситуациях, когда поток потребляется кодом, которому нужно «отменить потребление» некоторого количества данных, которые он оптимистично извлек из источника, чтобы данные могли быть переданы какой-либо другой стороне.
Метод stream.unshift(chunk) не может быть вызван после того, как событие 'end' было выпущено, в противном случае будет выброшено исключение.
Разработчики, использующие stream.unshift() часто должны рассмотреть возможность переключения на использование потока Transform вместо него. См. раздел API для разработчиков потоков для получения дополнительной информации.
// Pull off a header delimited by \n\n.
// Use unshift() if we get too much.
// Call the callback with (error, header, stream).
const { StringDecoder } = require('node:string_decoder');
function parseHeader(stream, callback) {
stream.on('error', callback);
stream.on('readable', onReadable);
const decoder = new StringDecoder('utf8');
let header = '';
function onReadable() {
let chunk;
while (null !== (chunk = stream.read())) {
const str = decoder.write(chunk);
if (str.includes('\n\n')) {
// Found the header boundary.
const split = str.split(/\n\n/);
header += split.shift();
const remaining = split.join('\n\n');
const buf = Buffer.from(remaining, 'utf8');
stream.removeListener('error', callback);
// Remove the 'readable' listener before unshifting.
stream.removeListener('readable', onReadable);
if (buf.length)
stream.unshift(buf);
// Now the body of the message can be read from the stream.
callback(null, header, stream);
return;
}
// Still reading the header.
header += str;
}
}
} copy В отличие от stream.push(chunk), stream.unshift(chunk) не завершит процесс чтения, сбросив внутреннее состояние чтения потока. Это может привести к непредвиденным результатам, если readable.unshift() вызывается во время чтения (например, из реализации stream._read() в пользовательском потоке). Однако, вызов readable.unshift() с последующим немедленным вызовом stream.push('') сбросит состояние чтения должным образом, но лучше просто избегать вызова readable.unshift() во время чтения.
readable.wrap(stream)
До Node.js 0.10 потоки не реализовывали весь API модуля node:stream, как он определен в настоящее время. (См. Совместимость для получения дополнительной информации.)
При использовании более старой библиотеки Node.js, которая испускает события 'data' и имеет метод stream.pause(), который является только рекомендательным, метод readable.wrap() может быть использован для создания потока Readable, который использует старый поток в качестве источника данных.
Использование readable.wrap() редко необходимо, но метод предоставлен для удобства взаимодействия со старыми приложениями и библиотеками Node.js.
const { OldReader } = require('./old-api-module.js');
const { Readable } = require('node:stream');
const oreader = new OldReader();
const myReader = new Readable().wrap(oreader);
myReader.on('readable', () => {
myReader.read(); // etc.
}); copy
readable[Symbol.asyncIterator]()
- Возвращает: <AsyncIterator> для полного потребления потока.
const fs = require('node:fs');
async function print(readable) {
readable.setEncoding('utf8');
let data = '';
for await (const chunk of readable) {
data += chunk;
}
console.log(data);
}
print(fs.createReadStream('file')).catch(console.error); copy Если цикл завершается с break, return, или throw, поток будет уничтожен. Другими словами, итерирование по потоку полностью потребляет поток. Поток будет читаться кусками размером, равным параметру highWaterMark. В приведенном выше примере данные будут в одном куске, если файл имеет размер меньше 64 КБ, потому что параметр highWaterMark не предоставляется методу fs.createReadStream().
readable[Symbol.asyncDispose]()
Вызывает readable.destroy() с AbortError и возвращает обещание, которое выполняется, когда поток завершен.
readable.compose(stream[, options])
-
stream<Поток> | <Итерируемый объект> | <AsyncIterable> | <Функция> -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Duplex> составной поток
stream.
import { Readable } from 'node:stream';
async function* splitToWords(source) {
for await (const chunk of source) {
const words = String(chunk).split(' ');
for (const word of words) {
yield word;
}
}
}
const wordsStream = Readable.from(['this is', 'compose as operator']).compose(splitToWords);
const words = await wordsStream.toArray();
console.log(words); // prints ['this', 'is', 'compose', 'as', 'operator'] copy См. stream.compose для получения дополнительной информации.
readable.iterator([options])
-
options<Объект>-
destroyOnReturn<логическое значение> Если установлено вfalse, вызовreturnдля асинхронного итератора или выход из циклаfor await...ofитерирования с помощьюbreak,return, илиthrowне уничтожит поток. По умолчанию:true.
-
- Возвращает: <AsyncIterator> для потребления потока.
Созданный этим методом итератор дает пользователям возможность отменить уничтожение потока, если цикл for await...of завершается с return, break, или throw, или если итератор должен уничтожить поток, если поток выпустил ошибку во время итерации.
const { Readable } = require('node:stream');
async function printIterator(readable) {
for await (const chunk of readable.iterator({ destroyOnReturn: false })) {
console.log(chunk); // 1
break;
}
console.log(readable.destroyed); // false
for await (const chunk of readable.iterator({ destroyOnReturn: false })) {
console.log(chunk); // Will print 2 and then 3
}
console.log(readable.destroyed); // True, stream was totally consumed
}
async function printSymbolAsyncIterator(readable) {
for await (const chunk of readable) {
console.log(chunk); // 1
break;
}
console.log(readable.destroyed); // true
}
async function showBoth() {
await printIterator(Readable.from([1, 2, 3]));
await printSymbolAsyncIterator(Readable.from([1, 2, 3]));
}
showBoth(); copy
readable.map(fn[, options])
-
fn<Функция> | <AsyncФункция> функция для обработки каждого фрагмента в потоке.-
data<любой> фрагмент данных из потока. -
options<Объект>-
signal<СигналПрерывания> прерывается, если поток уничтожен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное вызовfnдля вызова потока за раз. По умолчанию:1. -
highWaterMark<число> сколько элементов буферизовать, ожидая обработки пользователем сопоставленных элементов. По умолчанию:concurrency * 2 - 1. -
signal<СигналПрерывания> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <ПотокНаЧтение> поток, сопоставленный с функцией
fn.
Этот метод позволяет сопоставить поток. Функция fn будет вызываться для каждого фрагмента в потоке. Если функция fn возвращает промис, этот промис будет await перед передачей в результирующий поток.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).map((x) => x * 2)) {
console.log(chunk); // 2, 4, 6, 8
}
// With an asynchronous mapper, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).map((domain) => resolver.resolve4(domain), { concurrency: 2 });
for await (const result of dnsResults) {
console.log(result); // Logs the DNS result of resolver.resolve4.
} copy
readable.filter(fn[, options])
-
fn<Функция> | <AsyncФункция> функция для фильтрации фрагментов из потока.-
data<любой> фрагмент данных из потока. -
options<Объект>-
signal<СигналПрерывания> прерывается, если поток уничтожен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное вызовfnдля вызова потока за раз. По умолчанию:1. -
highWaterMark<число> сколько элементов буферизовать, ожидая обработки пользователем отфильтрованных элементов. По умолчанию:concurrency * 2 - 1. -
signal<СигналПрерывания> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <ПотокНаЧтение> отфильтрованный поток с предикатом
fn.
Этот метод позволяет отфильтровать поток. Для каждого фрагмента в потоке будет вызвана функция fn, и если она возвращает истинное значение, фрагмент будет передан в результирующий поток. Если функция fn возвращает промис, этот промис будет await.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).filter(async (domain) => {
const { address } = await resolver.resolve4(domain, { ttl: true });
return address.ttl > 60;
}, { concurrency: 2 });
for await (const result of dnsResults) {
// Logs domains with more than 60 seconds on the resolved dns record.
console.log(result);
} copy
readable.forEach(fn[, options])
-
fn<Функция> | <AsyncФункция> функция для вызова для каждого фрагмента потока.-
data<любой> фрагмент данных из потока. -
options<Объект>-
signal<СигналПрерывания> прерывается, если поток уничтожен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное вызовfnдля вызова потока за раз. По умолчанию:1. -
signal<СигналПрерывания> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Промис> промис для завершения потока.
Этот метод позволяет перебирать поток. Для каждого фрагмента в потоке будет вызвана функция fn. Если функция fn возвращает промис, этот промис будет await.
Этот метод отличается от циклов for await...of тем, что он может обрабатывать фрагменты параллельно. Кроме того, итерация forEach может быть остановлена только передачей опции signal и прерыванием связанного AbortController в то время как for await...of можно остановить с помощью break или return. В любом случае поток будет уничтожен.
Этот метод отличается от прослушивания события 'data' тем, что он использует событие readable в подлежащей машине и может ограничить количество одновременных вызовов fn.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).map(async (domain) => {
const { address } = await resolver.resolve4(domain, { ttl: true });
return address;
}, { concurrency: 2 });
await dnsResults.forEach((result) => {
// Logs result, similar to `for await (const result of dnsResults)`
console.log(result);
});
console.log('done'); // Stream has finished copy
readable.toArray([options])
-
options<Объект>-
signal<СигналПрерывания> позволяет отменить операцию toArray, если сигнал прерван.
-
- Возвращает: <Промис> промис, содержащий массив с содержимым потока.
Этот метод позволяет легко получить содержимое потока.
Поскольку этот метод читает весь поток в память, он аннулирует преимущества потоков. Он предназначен для совместимости и удобства, а не как основной способ потребления потоков.
import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';
await Readable.from([1, 2, 3, 4]).toArray(); // [1, 2, 3, 4]
// Make dns queries concurrently using .map and collect
// the results into an array using toArray
const dnsResults = await Readable.from([
'nodejs.org',
'openjsf.org',
'www.linuxfoundation.org',
]).map(async (domain) => {
const { address } = await resolver.resolve4(domain, { ttl: true });
return address;
}, { concurrency: 2 }).toArray(); copy
readable.some(fn[, options])
-
fn<Функция> | <AsyncФункция> функция для вызова для каждого фрагмента потока.-
data<любой> фрагмент данных из потока. -
options<Объект>-
signal<СигналПрерывания> прерывается, если поток уничтожен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное вызовfnдля вызова потока за раз. По умолчанию:1. -
signal<СигналПрерывания> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Промис> промис, равный
true, еслиfnвернуло истинное значение хотя бы для одного из фрагментов.
Этот метод похож на Array.prototype.some и вызывает fn для каждого фрагмента в потоке до тех пор, пока ожидаемое возвращаемое значение не станет true (или любым истинным значением). Как только вызов fn для фрагмента вернет истинное значение, поток уничтожается, и обещание выполняется со значением true. Если ни один из вызовов fn для фрагментов не вернёт истинное значение, обещание выполняется со значением false.
import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';
// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).some((x) => x > 2); // true
await Readable.from([1, 2, 3, 4]).some((x) => x < 0); // false
// With an asynchronous predicate, making at most 2 file checks at a time.
const anyBigFile = await Readable.from([
'file1',
'file2',
'file3',
]).some(async (fileName) => {
const stats = await stat(fileName);
return stats.size > 1024 * 1024;
}, { concurrency: 2 });
console.log(anyBigFile); // `true` if any file in the list is bigger than 1MB
console.log('done'); // Stream has finished copy
readable.find(fn[, options])
-
fn<Функция> | <AsyncFunction> функция для вызова для каждого фрагмента потока.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток уничтожен, позволяя прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное обращение кfnдля вызова по потоку одновременно. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Обещание> обещание, оценивающее первый фрагмент, для которого
fnвычислилось с истинным значением илиundefined, если элемент не был найден.
Этот метод похож на Array.prototype.find и вызывает fn для каждого фрагмента в потоке, чтобы найти фрагмент с истинным значением для fn. Как только ожидаемое возвращаемое значение вызова fn станет истинным, поток уничтожается, и обещание выполняется со значением, для которого fn вернуло истинное значение. Если все вызовы fn для фрагментов возвращают ложное значение, обещание выполняется со значением undefined.
import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';
// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).find((x) => x > 2); // 3
await Readable.from([1, 2, 3, 4]).find((x) => x > 0); // 1
await Readable.from([1, 2, 3, 4]).find((x) => x > 10); // undefined
// With an asynchronous predicate, making at most 2 file checks at a time.
const foundBigFile = await Readable.from([
'file1',
'file2',
'file3',
]).find(async (fileName) => {
const stats = await stat(fileName);
return stats.size > 1024 * 1024;
}, { concurrency: 2 });
console.log(foundBigFile); // File name of large file, if any file in the list is bigger than 1MB
console.log('done'); // Stream has finished copy
readable.every(fn[, options])
-
fn<Функция> | <AsyncFunction> функция для вызова для каждого фрагмента потока.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток уничтожен, позволяя прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное обращение кfnдля вызова по потоку одновременно. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Обещание> обещание, оценивающее
true, еслиfnвернуло истинное значение для всех фрагментов.
Этот метод похож на Array.prototype.every и вызывает fn для каждого фрагмента в потоке, чтобы проверить, являются ли все ожидаемые возвращаемые значения истинными значениями для fn. Как только ожидаемое возвращаемое значение вызова fn для фрагмента станет ложным, поток уничтожается, и обещание выполняется со значением false. Если все вызовы fn для фрагментов возвращают истинное значение, обещание выполняется со значением true.
import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';
// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).every((x) => x > 2); // false
await Readable.from([1, 2, 3, 4]).every((x) => x > 0); // true
// With an asynchronous predicate, making at most 2 file checks at a time.
const allBigFiles = await Readable.from([
'file1',
'file2',
'file3',
]).every(async (fileName) => {
const stats = await stat(fileName);
return stats.size > 1024 * 1024;
}, { concurrency: 2 });
// `true` if all files in the list are bigger than 1MiB
console.log(allBigFiles);
console.log('done'); // Stream has finished copy
readable.flatMap(fn[, options])
-
fn<Функция> | <AsyncGeneratorFunction> | <AsyncFunction> функция для отображения каждого фрагмента в потоке.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток уничтожен, позволяя прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное обращение кfnдля вызова по потоку одновременно. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Чтение> поток, сглаженный с помощью функции
fn.
Этот метод возвращает новый поток, применяя заданный обратный вызов к каждому фрагменту потока, а затем сплющивая результат.
Возможна передача потока или другого итерируемого или асинхронно итерируемого объекта из fn, и результирующие потоки будут объединены (сплющены) в возвращаемый поток.
import { Readable } from 'node:stream';
import { createReadStream } from 'node:fs';
// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).flatMap((x) => [x, x])) {
console.log(chunk); // 1, 1, 2, 2, 3, 3, 4, 4
}
// With an asynchronous mapper, combine the contents of 4 files
const concatResult = Readable.from([
'./1.mjs',
'./2.mjs',
'./3.mjs',
'./4.mjs',
]).flatMap((fileName) => createReadStream(fileName));
for await (const result of concatResult) {
// This will contain the contents (all chunks) of all 4 files
console.log(result);
} copy
readable.drop(limit[, options])
-
limit<число> количество фрагментов для удаления из потока. -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Чтение> поток с
limitопущенными фрагментами.
Этот метод возвращает новый поток, с limit опущенными первыми фрагментами.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).drop(2).toArray(); // [3, 4] copy
readable.take(limit[, options])
-
limit<число> количество фрагментов для извлечения из потока. -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Чтение> поток с
limitвзятыми фрагментами.
Этот метод возвращает новый поток с первыми limit фрагментами.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).take(2).toArray(); // [1, 2] copy
readable.reduce(fn[, initial[, options]])
-
fn<Функция> | <AsyncFunction> функция-редуктор для вызова над каждым фрагментом в потоке.-
previous<любой> значение, полученное из последнего вызоваfnили значениеinitial, если указано, или первый фрагмент потока в противном случае. -
data<любой> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывает, если поток уничтожен, позволяя прервать вызовfnраньше.
-
-
-
initial<любой> начальное значение для использования в редукции. -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Promise> обещание для конечного значения редукции.
Этот метод вызывает fn для каждого фрагмента потока в порядке, передавая результат вычисления по предыдущему элементу. Он возвращает обещание для конечного значения редукции.
Если начальное значение initial не указано, используется первый фрагмент потока. Если поток пуст, обещание отклоняется с TypeError с ERR_INVALID_ARGS свойством кода.
import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
const directoryPath = './src';
const filesInDir = await readdir(directoryPath);
const folderSize = await Readable.from(filesInDir)
.reduce(async (totalSize, file) => {
const { size } = await stat(join(directoryPath, file));
return totalSize + size;
}, 0);
console.log(folderSize); copy Функция-редуктор итерирует элементы потока по одному, что означает, что параметр concurrency или параллелизм отсутствуют. Для выполнения reduce конкуртентно, вы можете извлечь асинхронную функцию в метод readable.map.
import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
const directoryPath = './src';
const filesInDir = await readdir(directoryPath);
const folderSize = await Readable.from(filesInDir)
.map((file) => stat(join(directoryPath, file)), { concurrency: 2 })
.reduce((totalSize, { size }) => totalSize + size, 0);
console.log(folderSize); copy Потоки duplex и transform
Класс: stream.Duplex
Потоки duplex — это потоки, которые реализуют как интерфейсы Readable, так и Writable.
Примеры потоков duplex:
duplex.allowHalfOpen
Если false , поток автоматически завершит сторону записи, когда сторона чтения завершится. Начальное значение задается опцией конструктора allowHalfOpen, по умолчанию true.
Это можно изменить вручную, чтобы изменить поведение полуоткрытого потока Duplex , но это нужно сделать до выдачи события 'end'.
Класс: stream.Transform
Потоки transform — это потоки Duplex, где вывод каким-то образом связан с вводом. Как и все потоки Duplex, потоки Transform реализуют как интерфейсы Readable, так и Writable.
Примеры потоков transform:
transform.destroy([error])
Уничтожить поток и, необязательно, выдать событие 'error'. После этого вызова поток transform освободит все внутренние ресурсы. Реализаторы не должны переопределять этот метод, а вместо этого реализовать readable._destroy(). По умолчанию реализация _destroy() для Transform также выдает 'close' , если emitClose не установлено в false.
После вызова destroy() все последующие вызовы будут no-op и не будут генерировать дополнительных ошибок, кроме тех, что могут быть сгенерированы _destroy() как 'error'.
stream.finished(stream[, options], callback)
-
stream<Поток> | <ReadableStream> | <WritableStream> Поток чтения и/или записи/веб-поток. -
options<Объект>-
error<логическое> Если установлено вfalse, вызовemit('error', err)не считается завершением. По умолчанию:true. -
readable<логическое> При установке вfalse, обратный вызов будет вызван при завершении потока, даже если поток по-прежнему может быть читаемым. По умолчанию:true. -
writable<логическое> При установке вfalse, обратный вызов будет вызван при завершении потока, даже если поток по-прежнему может быть записываемым. По умолчанию:true. -
signal<AbortSignal> позволяет прервать ожидание завершения потока. Подлежащий поток не прерывается, если сигнал прерывается. Обратный вызов вызывается сAbortError. Все зарегистрированные слушатели, добавленные этой функцией, также будут удалены. -
cleanup<логическое> удалить все зарегистрированные слушатели потока. По умолчанию:false.
-
-
callback<Функция> Функция обратного вызова, принимающая необязательный аргумент ошибки. - Возвращает: <Функция> Функция очистки, которая удаляет все зарегистрированные слушатели.
Функция для получения уведомлений, когда поток больше не читается, не записывается, или произошла ошибка или событие преждевременного закрытия.
const { finished } = require('node:stream');
const fs = require('node:fs');
const rs = fs.createReadStream('archive.tar');
finished(rs, (err) => {
if (err) {
console.error('Stream failed.', err);
} else {
console.log('Stream is done reading.');
}
});
rs.resume(); // Drain the stream. copy Особо полезно в сценариях обработки ошибок, где поток уничтожается преждевременно (например, прерванный HTTP-запрос) и не будет издавать события 'end' или 'finish'.
API finished предоставляет версию обещания.
stream.finished() оставляет висящие слушатели событий (в частности, 'error', 'end', 'finish' и 'close') после вызова callback. Причина в том, что непредвиденные события 'error' (из-за неправильной реализации потоков) не вызывают непредвиденных сбоев. Если это поведение нежелательно, функция очистки, возвращаемая функцией, должна вызываться в обратном вызове:
const cleanup = finished(rs, (err) => {
cleanup();
// ...
}); copy
stream.pipeline(source[, ...transforms], destination, callback)
stream.pipeline(streams, callback)
-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> -
source<Stream> | <Iterable> | <AsyncIterable> | <Function> | <ReadableStream>- Возвращает: <Iterable> | <AsyncIterable>
-
...transforms<Stream> | <Function> | <TransformStream>-
source<AsyncIterable> - Возвращает: <AsyncIterable>
-
-
destination<Stream> | <Function> | <WritableStream>-
source<AsyncIterable> - Возвращает: <AsyncIterable> | <Promise>
-
-
callback<Function> Вызывается, когда обработка данных завершена.-
err<Error> -
valЗначение, возвращенное методомPromiseвызываемымdestination.
-
- Возвращает: <Stream>
Метод модуля для перенаправления данных между потоками и генераторами, перенаправляя ошибки и правильно очищая ресурсы, а также предоставляя обратный вызов по завершении обработки.
const { pipeline } = require('node:stream');
const fs = require('node:fs');
const zlib = require('node:zlib');
// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.
// A pipeline to gzip a potentially huge tar file efficiently:
pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz'),
(err) => {
if (err) {
console.error('Pipeline failed.', err);
} else {
console.log('Pipeline succeeded.');
}
},
); copy API pipeline предоставляет версию с promise.
stream.pipeline() вызовет stream.destroy(err) для всех потоков, кроме:
-
Readableпотоков, которые издали'end'или'close'. -
Writableпотоков, которые издали'finish'или'close'.
stream.pipeline() оставляет висящие обработчики событий в потоках после вызова callback. В случае повторного использования потоков после ошибки это может привести к утечкам обработчиков событий и необработанным ошибкам. Если последний поток является читаемым, висящие обработчики событий будут удалены, чтобы последний поток можно было обработать позже.
stream.pipeline() закрывает все потоки при возникновении ошибки. Использование IncomingRequest с pipeline может привести к неожиданному поведению, когда сокет будет уничтожен без отправки ожидаемого ответа. Смотрите пример ниже:
const fs = require('node:fs');
const http = require('node:http');
const { pipeline } = require('node:stream');
const server = http.createServer((req, res) => {
const fileStream = fs.createReadStream('./fileNotExist.txt');
pipeline(fileStream, res, (err) => {
if (err) {
console.log(err); // No such file
// this message can't be sent once `pipeline` already destroyed the socket
return res.end('error!!!');
}
});
}); copy
stream.compose(...streams)
stream.compose является экспериментальным.-
streams<Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> | <Duplex[]> | <Function> - Возвращает: <stream.Duplex>
Комбинирует два или более потоков в поток Duplex, который записывает в первый поток и считывает из последнего. Каждый предоставленный поток перенаправляется в следующий с использованием stream.pipeline. Если любой из потоков вызывает ошибку, все потоки уничтожаются, включая внешний поток Duplex.
Так как stream.compose возвращает новый поток, который в свою очередь может (и должен) быть перенаправлен в другие потоки, это позволяет создавать композиции. В отличие от этого, при передаче потоков в stream.pipeline, обычно первый поток является потоком чтения, а последний — потоком записи, образуя замкнутый цикл.
Если передан Function, он должен быть фабричным методом, принимающим source Iterable.
import { compose, Transform } from 'node:stream';
const removeSpaces = new Transform({
transform(chunk, encoding, callback) {
callback(null, String(chunk).replace(' ', ''));
},
});
async function* toUpper(source) {
for await (const chunk of source) {
yield String(chunk).toUpperCase();
}
}
let res = '';
for await (const buf of compose(removeSpaces, toUpper).end('hello world')) {
res += buf;
}
console.log(res); // prints 'HELLOWORLD' copy stream.compose может использоваться для преобразования асинхронных итераторов, генераторов и функций в потоки.
-
AsyncIterableпреобразует в читаемыйDuplex. Не может возвращатьnull. -
AsyncGeneratorFunctionпреобразует в читаемый/записываемый трансформирующийDuplex. Должен принимать исходныйAsyncIterableв качестве первого параметра. Не может возвращатьnull. -
AsyncFunctionпреобразует в записываемыйDuplex. Должен возвращать либоnull, либоundefined.
import { compose } from 'node:stream';
import { finished } from 'node:stream/promises';
// Convert AsyncIterable into readable Duplex.
const s1 = compose(async function*() {
yield 'Hello';
yield 'World';
}());
// Convert AsyncGenerator into transform Duplex.
const s2 = compose(async function*(source) {
for await (const chunk of source) {
yield String(chunk).toUpperCase();
}
});
let res = '';
// Convert AsyncFunction into writable Duplex.
const s3 = compose(async function(source) {
for await (const chunk of source) {
res += chunk;
}
});
await finished(compose(s1, s2, s3));
console.log(res); // prints 'HELLOWORLD' copy См. readable.compose(stream) для stream.compose как оператора.
stream.Readable.from(iterable[, options])
-
iterable<Iterable> Объект, реализующий протокол итерацииSymbol.asyncIteratorилиSymbol.iterator. Издает событие 'error', если передан null. -
options<Object> Опции, предоставляемыеnew stream.Readable([options]). По умолчаниюReadable.from()установитoptions.objectModeвtrue, если это явно не отменено установкойoptions.objectModeвfalse. - Возвращает: <stream.Readable>
Утилита для создания читаемых потоков из итераторов.
const { Readable } = require('node:stream');
async function * generate() {
yield 'hello';
yield 'streams';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
}); copy Вызов Readable.from(string) или Readable.from(buffer) не приведет к итерации строк или буферов для соответствия семантике других потоков по причинам производительности.
Если в качестве аргумента передан объект Iterable содержащий промисы, это может привести к необработанному отклонению.
const { Readable } = require('node:stream');
Readable.from([
new Promise((resolve) => setTimeout(resolve('1'), 1500)),
new Promise((_, reject) => setTimeout(reject(new Error('2')), 1000)), // Unhandled rejection
]); copy
stream.Readable.fromWeb(readableStream[, options])
-
readableStream<ReadableStream> -
options<Object>-
encoding<string> -
highWaterMark<number> -
objectMode<boolean> -
signal<AbortSignal>
-
- Возвращает: <stream.Readable>
stream.Readable.isDisturbed(stream)
-
stream<stream.Readable> | <ReadableStream> - Возвращает:
boolean
Возвращает, был ли поток прочитан или отменён.
stream.isErrored(stream)
-
stream<Readable> | <Writable> | <Duplex> | <WritableStream> | <ReadableStream> - Возвращает: <boolean>
Возвращает, столкнулся ли поток с ошибкой.
stream.isReadable(stream)
-
stream<Readable> | <Duplex> | <ReadableStream> - Возвращает: <boolean>
Возвращает, является ли поток читаемым.
stream.Readable.toWeb(streamReadable[, options])
-
streamReadable<stream.Readable> -
options<Object>-
strategy<Object>-
highWaterMark<number> Максимальный размер внутренней очереди (созданногоReadableStream) перед применением обратной связи при чтении из данногоstream.Readable. Если значение не указано, оно будет взято из данногоstream.Readable. -
size<Function> Функция, вычисляющая размер данного фрагмента данных. Если значение не указано, размер будет1для всех фрагментов.
-
-
- Возвращает: <ReadableStream>
stream.Writable.fromWeb(writableStream[, options])
-
writableStream<WritableStream> -
options<Object>-
decodeStrings<boolean> -
highWaterMark<number> -
objectMode<boolean> -
signal<AbortSignal>
-
- Возвращает: <stream.Writable>
stream.Writable.toWeb(streamWritable)
-
streamWritable<stream.Writable> - Возвращает: <WritableStream>
stream.Duplex.from(src)
-
src<Stream> | <Blob> | <ArrayBuffer> | <string> | <Iterable> | <AsyncIterable> | <AsyncGeneratorFunction> | <AsyncFunction> | <Promise> | <Object> | <ReadableStream> | <WritableStream>
Утилитарный метод для создания дуплексных потоков.
-
Streamпреобразует поток записи в потоки записиDuplexи поток чтения вDuplex. -
Blobпреобразует в читаемый потокDuplex. -
stringпреобразует в читаемый потокDuplex. -
ArrayBufferпреобразует в читаемый потокDuplex. -
AsyncIterableпреобразует в читаемыйDuplex. Нельзя использоватьnull. -
AsyncGeneratorFunctionпреобразует в преобразующийDuplexпоток чтения/записи. Должен принимать исходныйAsyncIterableв качестве первого параметра. Нельзя использоватьnull. -
AsyncFunctionпреобразует в записывающийDuplex. Должен возвращать либоnull, либоundefined. -
Object ({ writable, readable })преобразуетreadableиwritableвStreamи затем объединяет их вDuplex, гдеDuplexбудет записывать вwritableи читать изreadable. -
Promiseпреобразует в читаемыйDuplex. Значениеnullигнорируется. -
ReadableStreamпреобразует в читаемыйDuplex. -
WritableStreamпреобразует в записывающийDuplex. - Возвращает: <stream.Duplex>
Если в качестве аргумента передаётся объект Iterable , содержащий обещания, это может привести к необработанному отклонению.
const { Duplex } = require('node:stream');
Duplex.from([
new Promise((resolve) => setTimeout(resolve('1'), 1500)),
new Promise((_, reject) => setTimeout(reject(new Error('2')), 1000)), // Unhandled rejection
]); copy
stream.Duplex.fromWeb(pair[, options])
-
pair<Object>-
readable<ReadableStream> -
writable<WritableStream>
-
-
options<Object> - Возвращает: <stream.Duplex>
Модули MJS
import { Duplex } from 'node:stream';
import {
ReadableStream,
WritableStream,
} from 'node:stream/web';
const readable = new ReadableStream({
start(controller) {
controller.enqueue('world');
},
});
const writable = new WritableStream({
write(chunk) {
console.log('writable', chunk);
},
});
const pair = {
readable,
writable,
};
const duplex = Duplex.fromWeb(pair, { encoding: 'utf8', objectMode: true });
duplex.write('hello');
for await (const chunk of duplex) {
console.log('readable', chunk);
}
Модули CJS
const { Duplex } = require('node:stream');
const {
ReadableStream,
WritableStream,
} = require('node:stream/web');
const readable = new ReadableStream({
start(controller) {
controller.enqueue('world');
},
});
const writable = new WritableStream({
write(chunk) {
console.log('writable', chunk);
},
});
const pair = {
readable,
writable,
};
const duplex = Duplex.fromWeb(pair, { encoding: 'utf8', objectMode: true });
duplex.write('hello');
duplex.once('readable', () => console.log('readable', duplex.read()));
stream.Duplex.toWeb(streamDuplex)
-
streamDuplex<stream.Duplex> - Возвращает: <Объект>
-
readable<ReadableStream> -
writable<WritableStream>
-
Модули MJS
import { Duplex } from 'node:stream';
const duplex = Duplex({
objectMode: true,
read() {
this.push('world');
this.push(null);
},
write(chunk, encoding, callback) {
console.log('writable', chunk);
callback();
},
});
const { readable, writable } = Duplex.toWeb(duplex);
writable.getWriter().write('hello');
const { value } = await readable.getReader().read();
console.log('readable', value);
Модули CJS
const { Duplex } = require('node:stream');
const duplex = Duplex({
objectMode: true,
read() {
this.push('world');
this.push(null);
},
write(chunk, encoding, callback) {
console.log('writable', chunk);
callback();
},
});
const { readable, writable } = Duplex.toWeb(duplex);
writable.getWriter().write('hello');
readable.getReader().read().then((result) => {
console.log('readable', result.value);
});
stream.addAbortSignal(signal, stream)
-
signal<AbortSignal> Сигнал, представляющий возможную отмену -
stream<Поток> | <ReadableStream> | <WritableStream> Поток, к которому нужно прикрепить сигнал.
Прикрепляет AbortSignal к потоку чтения или записи. Это позволяет коду управлять уничтожением потока, используя AbortController.
Вызов abort на AbortController, соответствующем переданному AbortSignal будет эквивалентен вызову .destroy(new AbortError()) на потоке, и controller.error(new AbortError()) для webstreams.
const fs = require('node:fs');
const controller = new AbortController();
const read = addAbortSignal(
controller.signal,
fs.createReadStream(('object.json')),
);
// Later, abort the operation closing the stream
controller.abort(); copy Или с использованием AbortSignal с потоком чтения как асинхронной итерируемой последовательности:
const controller = new AbortController();
setTimeout(() => controller.abort(), 10_000); // set a timeout
const stream = addAbortSignal(
controller.signal,
fs.createReadStream(('object.json')),
);
(async () => {
try {
for await (const chunk of stream) {
await process(chunk);
}
} catch (e) {
if (e.name === 'AbortError') {
// The operation was cancelled
} else {
throw e;
}
}
})(); copy Или с использованием AbortSignal с ReadableStream:
const controller = new AbortController();
const rs = new ReadableStream({
start(controller) {
controller.enqueue('hello');
controller.enqueue('world');
controller.close();
},
});
addAbortSignal(controller.signal, rs);
finished(rs, (err) => {
if (err) {
if (err.name === 'AbortError') {
// The operation was cancelled
}
}
});
const reader = rs.getReader();
reader.read().then(({ value, done }) => {
console.log(value); // hello
console.log(done); // false
controller.abort();
}); copy
stream.getDefaultHighWaterMark(objectMode)
-
objectMode<логическое> - Возвращает: <целое число>
Возвращает значение по умолчанию highWaterMark, используемое потоками. По умолчанию 16384 (16 КБ), или 16 для objectMode.
stream.setDefaultHighWaterMark(objectMode, value)
-
objectMode<логическое> -
value<целое число> Значение highWaterMark
Устанавливает значение по умолчанию highWaterMark, используемое потоками.
API для разработчиков потоков
API модуля node:stream разработан для упрощения реализации потоков с использованием прототипного наследования JavaScript.
Сначала разработчик потока должен объявить новый класс JavaScript, расширяющий один из четырёх базовых классов потоков (stream.Writable, stream.Readable, stream.Duplex, или stream.Transform ), убедившись, что он вызывает соответствующий конструктор родительского класса:
const { Writable } = require('node:stream');
class MyWritable extends Writable {
constructor({ highWaterMark, ...options }) {
super({ highWaterMark });
// ...
}
} copy При расширении потоков следует учитывать, какие параметры пользователь может и должен предоставить перед передачей их в базовый конструктор. Например, если реализация делает предположения относительно параметров autoDestroy и emitClose, не позволяйте пользователю их переопределять. Будьте явными в отношении передаваемых параметров, вместо неявной передачи всех параметров.
Новый класс потока должен реализовать один или несколько методов, в зависимости от типа создаваемого потока, как подробно описано в таблице ниже:
| Сценарий использования | Класс | Методы для реализации |
|---|---|---|
| Только чтение | Readable |
_read() |
| Только запись | Writable |
_write(), _writev(), _final()
|
| Чтение и запись | Duplex |
_read(), _write(), _writev(), _final()
|
| Обработка данных, записанных, затем чтение результата | Transform |
_transform(), _flush(), _final()
|
Код реализации потока никогда не должен вызывать «публичные» методы потока, предназначенные для использования потребителями (как описано в разделе API для потребителей потоков). Это может привести к нежелательным побочным эффектам в прикладном коде, использующем поток.
Избегайте переопределения публичных методов, таких как write(), end(), cork(), uncork(), read() и destroy(), или испускания внутренних событий, таких как 'error', 'data', 'end', 'finish' и 'close' через .emit(). Это может нарушить текущие и будущие инварианты потока, что приведёт к проблемам с поведением и/или совместимостью с другими потоками, утилитами потоков и ожиданиями пользователей.
Упрощённое создание
Во многих простых случаях можно создать поток, не прибегая к наследованию. Это можно сделать, непосредственно создав экземпляры объектов stream.Writable, stream.Readable, stream.Duplex, или stream.Transform и передав соответствующие методы в качестве параметров конструктора.
const { Writable } = require('node:stream');
const myWritable = new Writable({
construct(callback) {
// Initialize state and load resources...
},
write(chunk, encoding, callback) {
// ...
},
destroy() {
// Free resources...
},
}); copy Реализация потока записи
Класс stream.Writable расширяется для реализации потока Writable.
Пользовательские потоки Writable обязаны вызывать конструктор new stream.Writable([options]) и реализовывать метод writable._write() и/или writable._writev().
new stream.Writable([options])
-
options<Объект>-
highWaterMark<число> Уровень буфера, когдаstream.write()начинает возвращатьfalse. По умолчанию:65536(64 КБ), или16для потоковobjectMode. -
decodeStrings<логическое значение> Необходимо ли кодироватьstringзначения, передаваемые вstream.write()вBuffer(с кодировкой, указанной в вызовеstream.write()) перед передачей их вstream._write(). Другие типы данных не преобразуются (например,Bufferне декодируются вstring). Установка значения в false предотвратит преобразованиеstring. По умолчанию:true. -
defaultEncoding<строка> Кодировка по умолчанию, используемая, когда кодировка не указана как аргумент вstream.write(). По умолчанию:'utf8'. -
objectMode<логическое значение> Является ли операцияstream.write(anyObj)корректной. В случае установки параметра, появляется возможность записи в поток значений JavaScript, отличных от строк, <Buffer>, <TypedArray> или <DataView>, если это поддерживается реализацией потока. По умолчанию: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('node:stream');
class MyWritable extends Writable {
constructor(options) {
// Calls the stream.Writable() constructor.
super(options);
// ...
}
} copy Или, при использовании конструкторов в стиле до ES6:
const { Writable } = require('node:stream');
const util = require('node:util');
function MyWritable(options) {
if (!(this instanceof MyWritable))
return new MyWritable(options);
Writable.call(this, options);
}
util.inherits(MyWritable, Writable); copy Или, используя упрощённый подход к конструктору:
const { Writable } = require('node:stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
},
}); copy Вызов abort на объекте AbortController, соответствующем переданному AbortSignal, будет работать так же, как и вызов .destroy(new AbortError()) на потоке записи.
const { Writable } = require('node:stream');
const controller = new AbortController();
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
},
signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort(); copy
writable._construct(callback)
-
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки) после завершения инициализации потока.
Метод _construct() НЕЛЬЗЯ вызывать напрямую. Он может быть реализован производными классами, и в этом случае будет вызываться только внутренними методами класса Writable.
Эта необязательная функция будет вызвана через tick после возвращения конструктора потока, откладывая любые вызовы _write(), _final() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной инициализации ресурсов перед использованием потока.
const { Writable } = require('node:stream');
const fs = require('node:fs');
class WriteStream extends Writable {
constructor(filename) {
super();
this.filename = filename;
this.fd = null;
}
_construct(callback) {
fs.open(this.filename, (err, fd) => {
if (err) {
callback(err);
} else {
this.fd = fd;
callback();
}
});
}
_write(chunk, encoding, callback) {
fs.write(this.fd, chunk, callback);
}
_destroy(err, callback) {
if (this.fd) {
fs.close(this.fd, (er) => callback(er || err));
} else {
callback(err);
}
}
} copy
writable._write(chunk, encoding, callback)
-
chunk<Буфер> | <строка> | <любой> ДанныеBufferдля записи, преобразованные изstring, переданного в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<Буфер> | <строка> Экземпляр буфера или строка, содержащая данные для записи. Значение будет строкой, если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(). Он может быть переопределён дочерними классами, но НЕ должен вызываться напрямую. Кроме того, callback не следует смешивать с async/await после его выполнения при разрешении промиса.
writable._final(callback)
-
callback<Функция> Вызов этой функции (при необходимости с аргументом ошибки) по завершении записи всех оставшихся данных.
Метод _final() НЕЛЬЗЯ вызывать напрямую. Он может быть реализован дочерними классами и, если реализован, вызывается только внутренними методами класса Writable.
Эта необязательная функция вызывается перед закрытием потока, откладывая событие 'finish' до вызова callback. Это полезно для закрытия ресурсов или записи буферизованных данных перед завершением потока.
Ошибки при записи
Ошибки, возникающие во время обработки методов writable._write(), writable._writev() и writable._final(), должны обрабатываться вызовом обратного вызова с ошибкой в качестве первого аргумента. Исключение Error внутри этих методов или ручное генерирование события 'error' приводит к неопределённому поведению.
Если поток Readable подключается к потоку Writable и Writable генерирует ошибку, то поток Readable будет отключён.
const { Writable } = require('node:stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
},
}); copy Пример потока на запись
Следующий пример иллюстрирует довольно упрощённую (и несколько бесполезную) реализацию пользовательского потока на запись Writable. Хотя этот конкретный экземпляр потока на запись Writable не обладает практической ценностью, пример иллюстрирует каждый из необходимых элементов пользовательского экземпляра потока Writable:
const { Writable } = require('node:stream');
class MyWritable extends Writable {
_write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
}
} copy Декодирование буферов в потоке на запись
Декодирование буферов — распространённая задача, например, при использовании трансформаторов, вход которых является строкой. Это нетривиальный процесс при использовании кодировок символов с несколькими байтами, таких как UTF-8. Следующий пример демонстрирует, как декодировать многобайтовые строки с использованием StringDecoder и Writable.
const { Writable } = require('node:stream');
const { StringDecoder } = require('node:string_decoder');
class StringWritable extends Writable {
constructor(options) {
super(options);
this._decoder = new StringDecoder(options && options.defaultEncoding);
this.data = '';
}
_write(chunk, encoding, callback) {
if (encoding === 'buffer') {
chunk = this._decoder.write(chunk);
}
this.data += chunk;
callback();
}
_final(callback) {
this.data += this._decoder.end();
callback();
}
}
const euro = [[0xE2, 0x82], [0xAC]].map(Buffer.from);
const w = new StringWritable();
w.write('currency: ');
w.write(euro[0]);
w.end(euro[1]);
console.log(w.data); // currency: € copy Реализация потока на чтение
Класс stream.Readable расширяется для реализации потока Readable.
Пользовательские потоки Readable ДОЛЖНЫ вызывать конструктор new stream.Readable([options]) и реализовывать метод readable._read().
new stream.Readable([options])
-
options<Объект>-
highWaterMark<число> Максимальное количество байтов для хранения во внутрене буфере перед прекращением чтения из базового ресурса. По умолчанию:65536(64 КБ) или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('node:stream');
class MyReadable extends Readable {
constructor(options) {
// Calls the stream.Readable(options) constructor.
super(options);
// ...
}
} copy Или, при использовании конструкторов в стиле до ES6:
const { Readable } = require('node:stream');
const util = require('node:util');
function MyReadable(options) {
if (!(this instanceof MyReadable))
return new MyReadable(options);
Readable.call(this, options);
}
util.inherits(MyReadable, Readable); copy Или, используя упрощённый подход к конструктору:
const { Readable } = require('node:stream');
const myReadable = new Readable({
read(size) {
// ...
},
}); copy Вызов abort на AbortController соответствующем переданному AbortSignal будет вести себя так же, как вызов .destroy(new AbortError()) на создаваемом чтении.
const { Readable } = require('node:stream');
const controller = new AbortController();
const read = new Readable({
read(size) {
// ...
},
signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort(); copy
readable._construct(callback)
-
callback<Функция> Вызовите эту функцию (по желанию с аргументом ошибки) при завершении инициализации потока.
Метод _construct() НЕ ДОЛЖЕН вызываться напрямую. Он может быть реализован дочерними классами и, если это так, вызывается только внутренними методами класса Readable.
Эта необязательная функция будет запланирована на следующий тик конструктором потока, откладывая любые вызовы _read() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной инициализации ресурсов до использования потока.
const { Readable } = require('node:stream');
const fs = require('node:fs');
class ReadStream extends Readable {
constructor(filename) {
super();
this.filename = filename;
this.fd = null;
}
_construct(callback) {
fs.open(this.filename, (err, fd) => {
if (err) {
callback(err);
} else {
this.fd = fd;
callback();
}
});
}
_read(n) {
const buf = Buffer.alloc(n);
fs.read(this.fd, buf, 0, n, null, (err, bytesRead) => {
if (err) {
this.destroy(err);
} else {
this.push(bytesRead > 0 ? buf.slice(0, bytesRead) : null);
}
});
}
_destroy(err, callback) {
if (this.fd) {
fs.close(this.fd, (er) => callback(er || err));
} else {
callback(err);
}
}
} copy
readable._read(size)
-
size<число> Количество байтов для асинхронного чтения
Эта функция НЕ ДОЛЖНА вызываться кодом приложения напрямую. Она должна реализовываться дочерними классами и вызываться только внутренними методами класса 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<Буфер> | <Массив типизированных данных> | <DataView> | <строка> | <null> | <любое значение> Чанк данных для добавления в очередь чтения. Для потоков, не работающих в режиме объектов,chunkдолжен быть <строкой>, <буфером>, <массивом типизированных данных> или <DataView>. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript. -
encoding<строка> Кодировка строковых чанков. Должна быть допустимой кодировкой, такой какBufferили'utf8'или'ascii'. - Возвращает: <логическое значение>
trueесли можно продолжить добавлять чанки данных;falseв противном случае.
Если chunk — это <буфер>, <массив типизированных данных>, <DataView> или <строка>, данные будут добавлены во внутреннюю очередь для потребления пользователями потока. Передача chunk в качестве null сигнализирует об окончании потока (EOF), после чего больше данных не может быть написано.
Когда Readable работает в приостановленном режиме, добавленные данные с помощью readable.push() могут быть прочитаны, вызвав метод readable.read() при выводе события 'readable'.
Если Readable работает в режиме потока, данные, добавленные с помощью readable.push() будут переданы путём вывода события 'data'.
Метод readable.push() разработан для максимальной гибкости. Например, при обёртке более низкого уровня источника, предоставляющего механизм паузы/возобновления и обратный вызов данных, низкоуровневый источник можно обернуть с помощью настраиваемого экземпляра Readable:
// `_source` is an object with readStop() and readStart() methods,
// and an `ondata` member that gets called when it has data, and
// an `onend` member that gets called when the data is over.
class SourceWrapper extends Readable {
constructor(options) {
super(options);
this._source = getLowLevelSourceObject();
// Every time there's data, push it into the internal buffer.
this._source.ondata = (chunk) => {
// If push() returns false, then stop reading from source.
if (!this.push(chunk))
this._source.readStop();
};
// When the source ends, push the EOF-signaling `null` chunk.
this._source.onend = () => {
this.push(null);
};
}
// _read() will be called when the stream wants to pull more data in.
// The advisory size argument is ignored in this case.
_read(size) {
this._source.readStart();
}
} copy Метод readable.push() используется для помещения содержимого во внутренний буфер. Он может быть вызван методом readable._read().
Для потоков, не работающих в режиме объектов, если параметр chunk метода readable.push() равен undefined, он будет обработан как пустая строка или буфер. Дополнительную информацию см. в readable.push('').
Ошибки при чтении
Ошибки, возникающие во время обработки метода readable._read(), должны быть обработаны через метод readable.destroy(err). Бросание исключения Error внутри метода readable._read() или ручное излучение события 'error' приводит к неопределенному поведению.
const { Readable } = require('node:stream');
const myReadable = new Readable({
read(size) {
const err = checkSomeErrorCondition();
if (err) {
this.destroy(err);
} else {
// Do some work.
}
},
}); copy Пример счетного потока
Ниже приведен базовый пример потока Readable, который излучает числа от 1 до 1 000 000 в возрастающем порядке, а затем завершается.
const { Readable } = require('node:stream');
class Counter extends Readable {
constructor(opt) {
super(opt);
this._max = 1000000;
this._index = 1;
}
_read() {
const i = this._index++;
if (i > this._max)
this.push(null);
else {
const str = String(i);
const buf = Buffer.from(str, 'ascii');
this.push(buf);
}
}
} copy Реализация дуплексного потока
Поток Duplex — это поток, реализующий как Readable, так и Writable, например, соединение TCP-соккета.
Поскольку JavaScript не поддерживает множественное наследование, класс 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указано. По умолчанию:false. -
readableHighWaterMark<число> УстанавливаетhighWaterMarkдля стороны чтения потока. Не имеет эффекта, если указаноhighWaterMark. -
writableHighWaterMark<число> УстанавливаетhighWaterMarkдля стороны записи потока. Не имеет эффекта, если указаноhighWaterMark.
-
const { Duplex } = require('node:stream');
class MyDuplex extends Duplex {
constructor(options) {
super(options);
// ...
}
} copy Или при использовании конструкторов в стиле до ES6:
const { Duplex } = require('node:stream');
const util = require('node:util');
function MyDuplex(options) {
if (!(this instanceof MyDuplex))
return new MyDuplex(options);
Duplex.call(this, options);
}
util.inherits(MyDuplex, Duplex); copy Или используя упрощённый подход к конструктору:
const { Duplex } = require('node:stream');
const myDuplex = new Duplex({
read(size) {
// ...
},
write(chunk, encoding, callback) {
// ...
},
}); copy При использовании конвейера:
const { Transform, pipeline } = require('node:stream');
const fs = require('node:fs');
pipeline(
fs.createReadStream('object.json')
.setEncoding('utf8'),
new Transform({
decodeStrings: false, // Accept string input rather than Buffers
construct(callback) {
this.data = '';
callback();
},
transform(chunk, encoding, callback) {
this.data += chunk;
callback();
},
flush(callback) {
try {
// Make sure is valid json.
JSON.parse(this.data);
this.push(this.data);
callback();
} catch (err) {
callback(err);
}
},
}),
fs.createWriteStream('valid-object.json'),
(err) => {
if (err) {
console.error('failed', err);
} else {
console.log('completed');
}
},
); copy Пример дуплексного потока
Ниже приведён простой пример потока Duplex , который оборачивает гипотетический объект нижнего уровня, к которому можно записывать данные и из которого можно читать данные, хотя API несовместим с потоками Node.js. Ниже приведён простой пример потока Duplex , который буферизует входящие данные с помощью интерфейса Writable , которые считываются обратно с помощью интерфейса Readable.
const { Duplex } = require('node:stream');
const kSource = Symbol('source');
class MyDuplex extends Duplex {
constructor(source, options) {
super(options);
this[kSource] = source;
}
_write(chunk, encoding, callback) {
// The underlying source only deals with strings.
if (Buffer.isBuffer(chunk))
chunk = chunk.toString();
this[kSource].writeSomeData(chunk);
callback();
}
_read(size) {
this[kSource].fetchSomeData(size, (data, encoding) => {
this.push(Buffer.from(data, encoding));
});
}
} copy Самый важный аспект дуплексного потока заключается в том, что стороны чтения и записи работают независимо друг от друга, несмотря на совместное существование в одном экземпляре объекта.
Дуплексные потоки в режиме объектов
Для потоков Duplex параметр objectMode может быть установлен исключительно для стороны чтения или записи с помощью параметров readableObjectMode и writableObjectMode соответственно.
Например, в следующем примере создаётся новый поток Transform (который является типом потока Duplex), у которого сторона чтения в режиме объектов принимает числа JavaScript, которые преобразуются в шестнадцатеричные строки на стороне записи.
const { Transform } = require('node:stream');
// All Transform streams are also Duplex Streams.
const myTransform = new Transform({
writableObjectMode: true,
transform(chunk, encoding, callback) {
// Coerce the chunk to a number if necessary.
chunk |= 0;
// Transform the chunk into something else.
const data = chunk.toString(16);
// Push the data onto the readable queue.
callback(null, '0'.repeat(data.length % 2) + data);
},
});
myTransform.setEncoding('ascii');
myTransform.on('data', (chunk) => console.log(chunk));
myTransform.write(1);
// Prints: 01
myTransform.write(10);
// Prints: 0a
myTransform.write(100);
// Prints: 64 copy Реализация потока преобразования
Поток Transform — это поток Duplex, где вывод вычисляется каким-либо способом из входных данных. Примеры включают потоки zlib или crypto, которые сжимают, шифруют или дешифруют данные.
Нет требований, чтобы размер вывода был таким же, как размер входных данных, количество фрагментов было таким же или данные поступали одновременно. Например, поток Hash всегда будет иметь только один фрагмент вывода, который предоставляется при завершении входных данных. Поток zlib будет генерировать вывод, который может быть значительно меньше или значительно больше, чем входные данные.
Класс stream.Transform расширяется для реализации потока Transform.
Класс stream.Transform прототипически наследуется от stream.Duplex и реализует свои собственные версии методов writable._write() и readable._read(). Пользовательские реализации Transform обязаны реализовать метод transform._transform() и могут также реализовать метод transform._flush().
При использовании потоков Transform следует соблюдать осторожность, так как запись данных в поток может привести к приостановке стороны чтения потока, если вывод на стороне записи не обрабатывается.
new stream.Transform([options])
-
options<Объект> Передаётся в конструкторы какWritable, так иReadable. Также имеет следующие поля:-
transform<Функция> Реализация методаstream._transform(). -
flush<Функция> Реализация методаstream._flush().
-
const { Transform } = require('node:stream');
class MyTransform extends Transform {
constructor(options) {
super(options);
// ...
}
} copy Или при использовании конструкторов в стиле до ES6:
const { Transform } = require('node:stream');
const util = require('node:util');
function MyTransform(options) {
if (!(this instanceof MyTransform))
return new MyTransform(options);
Transform.call(this, options);
}
util.inherits(MyTransform, Transform); copy Или используя упрощённый подход к конструктору:
const { Transform } = require('node:stream');
const myTransform = new Transform({
transform(chunk, encoding, callback) {
// ...
},
}); copy Событие: 'end'
Событие 'end' исходит от класса stream.Readable. Событие 'end' излучается после вывода всех данных, что происходит после вызова обратного вызова в методе transform._flush(). В случае ошибки событие 'end' не должно излучаться.
Событие: 'finish'
Событие 'finish' исходит от класса stream.Writable. Событие 'finish' излучается после вызова stream.end() и обработки всех фрагментов методом stream._transform(). В случае ошибки событие 'finish' не должно излучаться.
transform._flush(callback)
-
callback<Функция> Функция обратного вызова (по желанию с аргументом ошибки и данными), которая вызывается при сбросе оставшихся данных.
Эта функция НЕ должна вызываться кодом приложения напрямую. Она должна реализовываться дочерними классами и вызываться только внутренними методами класса Readable.
В некоторых случаях операция преобразования может потребовать вывода дополнительных данных в конце потока. Например, поток сжатия zlib будет хранить объем внутреннего состояния, используемого для оптимального сжатия вывода. Однако, когда поток завершается, эти дополнительные данные необходимо сбросить, чтобы сжатые данные были полными.
Пользовательские реализации Transform могут реализовывать метод transform._flush(). Он будет вызван, когда больше нет данных для чтения, но перед тем, как будет выведено событие 'end', сигнализирующее о конце потока Readable.
В реализации transform._flush(), метод transform.push() может быть вызван ноль или более раз, по мере необходимости. Функция callback должна быть вызвана, когда операция сброса завершена.
Метод transform._flush() имеет префикс подчеркивания, потому что он является внутренним для класса, который его определяет, и никогда не должен вызываться непосредственно программами пользователя.
transform._transform(chunk, encoding, callback)
-
chunk<Буфер> | <строка> | <любой> Данные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);
}; copy Метод transform._transform() имеет префикс подчеркивания, потому что он является внутренним для класса, который его определяет, и никогда не должен вызываться непосредственно программами пользователя.
transform._transform() никогда не вызывается параллельно; потоки реализуют механизм очереди, и для получения следующего фрагмента необходимо вызвать callback, синхронно или асинхронно.
Класс: stream.PassThrough
Класс stream.PassThrough — это тривиальная реализация потока Transform, который просто передает входные байты на вывод. Его цель в первую очередь состоит в предоставлении примеров и проведении тестов, но в некоторых случаях stream.PassThrough полезен как строительный блок для новых типов потоков.
Дополнительные заметки
Совместимость потоков с асинхронными генераторами и асинхронными итераторами
С поддержкой асинхронных генераторов и итераторов в JavaScript, асинхронные генераторы на данный момент являются конструкцией потока первого класса на уровне языка.
Ниже представлены некоторые распространенные случаи взаимодействия при использовании потоков Node.js с асинхронными генераторами и асинхронными итераторами.
Использование асинхронных итераторов для потребления потоков чтения
(async function() {
for await (const chunk of readable) {
console.log(chunk);
}
})(); copy Асинхронные итераторы регистрируют постоянный обработчик ошибок в потоке для предотвращения любых необработанных ошибок после уничтожения.
Создание потоков чтения с помощью асинхронных генераторов
Поток чтения Node.js можно создать из асинхронного генератора, используя служебный метод Readable.from():
const { Readable } = require('node:stream');
const ac = new AbortController();
const signal = ac.signal;
async function * generate() {
yield 'a';
await someLongRunningFn({ signal });
yield 'b';
yield 'c';
}
const readable = Readable.from(generate());
readable.on('close', () => {
ac.abort();
});
readable.on('data', (chunk) => {
console.log(chunk);
}); copy Перенаправление в потоки записи из асинхронных итераторов
При записи в поток записи из асинхронного итератора убедитесь в правильной обработке обратной задержки и ошибок. stream.pipeline() абстрагирует обработку обратной задержки и ошибок, связанных с обратной задержкой:
const fs = require('node:fs');
const { pipeline } = require('node:stream');
const { pipeline: pipelinePromise } = require('node:stream/promises');
const writable = fs.createWriteStream('./file');
const ac = new AbortController();
const signal = ac.signal;
const iterator = createIterator({ signal });
// Callback Pattern
pipeline(iterator, writable, (err, value) => {
if (err) {
console.error(err);
} else {
console.log(value, 'value returned');
}
}).on('close', () => {
ac.abort();
});
// Promise Pattern
pipelinePromise(iterator, writable)
.then((value) => {
console.log(value, 'value returned');
})
.catch((err) => {
console.error(err);
ac.abort();
}); copy Совместимость со старыми версиями Node.js
До версии Node.js 0.10 интерфейс потока Readable был проще, но также менее мощным и менее полезным.
- Вместо ожидания вызовов метода
stream.read(), события'data'начинали излучаться немедленно. Приложениям, которым потребовалась бы определенная работа для принятия решения о том, как обработать данные, потребовалось бы сохранять считываемые данные в буферах, чтобы данные не потерялись. - Метод
stream.pause()был рекомендательным, а не гарантированным. Это означало, что все еще необходимо было быть готовым к получению событий'data'даже когда поток находился в приостановленном состоянии.
В Node.js 0.10 был добавлен класс Readable. Для обратной совместимости с более старыми программами Node.js потоки Readable переключаются в режим «потока» при добавлении обработчика события 'data' или вызове метода stream.resume(). Это означает, что даже при отсутствии использования нового метода stream.read() и события 'readable', больше не нужно беспокоиться о потере фрагментов 'data'.
Хотя большинство приложений будут продолжать работать нормально, это создает частный случай в следующих условиях:
- Ни один обработчик события
'data'не добавлен. - Метод
stream.resume()никогда не вызывался. - Поток не перенаправлен ни на какое место назначения потока записи.
Например, рассмотрим следующий код:
// WARNING! BROKEN!
net.createServer((socket) => {
// We add an 'end' listener, but never consume the data.
socket.on('end', () => {
// It will never get here.
socket.end('The message was received but was not processed.\n');
});
}).listen(1337); copy До Node.js 0.10 входящие данные сообщения просто отбрасывались. Однако в Node.js 0.10 и более поздних версиях сокет остается приостановленным навсегда.
Решение в этой ситуации — вызвать метод stream.resume(), чтобы начать поток данных:
// Workaround.
net.createServer((socket) => {
socket.on('end', () => {
socket.end('The message was received but was not processed.\n');
});
// Start the flow of data, discarding it.
socket.resume();
}).listen(1337); copy Помимо новых потоков Readable, которые переключаются в режим потока, потоки в стиле до 0.10 можно обернуть в класс Readable с помощью метода readable.wrap().
readable.read(0)
В некоторых случаях необходимо запустить обновление механизмов подлежащего потока чтения, не потребляя при этом никаких данных. В таких случаях можно вызвать readable.read(0), что всегда вернёт null.
Если внутренний буфер чтения находится ниже highWaterMark, и поток в данный момент не читает, то вызов stream.read(0) запустит вызов stream._read() на низком уровне.
Хотя большинство приложений практически никогда не будут нуждаться в этом, существуют ситуации в Node.js, где это делается, особенно во внутренних частях класса потока Readable.
readable.push('')
Использование readable.push('') не рекомендуется.
Передача нулевого байтового значения <строки>, <буфера>, <TypedArray> или <DataView> в поток, который не находится в режиме объектов, имеет интересный побочный эффект. Поскольку это вызов readable.push(), вызов завершит процесс чтения. Однако, поскольку аргументом является пустая строка, данные не добавляются в буфер чтения, поэтому пользователь ничего не может прочитать.
highWaterMark расхождение после вызова readable.setEncoding()
Использование readable.setEncoding() изменит поведение работы highWaterMark в режиме, отличном от режима объектов.
Обычно размер текущего буфера измеряется относительно highWaterMark в байтах. Однако после вызова setEncoding() функция сравнения начнёт измерять размер буфера в символах.
Это не проблема в общих случаях с latin1 или ascii. Но рекомендуется быть внимательным к этому поведению, работая со строками, которые могут содержать многобайтовые символы.
© Joyent, Inc. and other Node contributors
Licensed under the MIT License.
Node.js is a trademark of Joyent, Inc. and is used with its permission.
We are not endorsed by or affiliated with Joyent.
https://nodejs.org/api/stream.html