Поток[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<строка> | <Buffer> | <TypedArray> | <DataView> | <любой> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжен быть <строкой>, <Buffer>, <TypedArray> или <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() устанавливает кодировку по умолчанию encoding для потока Writable.
writable.uncork()
Метод writable.uncork() сбрасывает все данные, буферизованные с момента вызова stream.cork().
При использовании writable.cork() и writable.uncork() для управления буферизацией записей в поток, отложите вызовы writable.uncork() используя process.nextTick(). Это позволяет объединять все вызовы writable.write() в рамках фазы цикла событий Node.js.
stream.cork();
stream.write('some ');
stream.write('data ');
process.nextTick(() => stream.uncork()); copy Если метод writable.cork() вызывается несколько раз для потока, необходимо вызвать такое же количество раз метод writable.uncork() для сброса буферизованных данных.
stream.cork();
stream.write('some ');
stream.cork();
stream.write('data ');
process.nextTick(() => {
stream.uncork();
// The data will not be flushed until uncork() is called a second time.
stream.uncork();
}); copy См. также: writable.cork().
writable.writable
Является true если безопасно вызвать writable.write(), что означает, что поток не был уничтожен, не получил ошибку и не завершен.
writable.writableAborted
Возвращает значение true, если поток был уничтожен или получил ошибку до испускания события 'finish'.
writable.writableEnded
Принимает значение true после вызова writable.end(). Это свойство не указывает, были ли данные сброшены. Для этой цели используйте свойство writable.writableFinished.
writable.writableCorked
Число вызовов writable.uncork(), необходимых для полного разбуривания потока.
writable.errored
Возвращает ошибку, если поток был уничтожен с ошибкой.
writable.writableFinished
Устанавливается в true непосредственно перед испусканием события 'finish'.
writable.writableHighWaterMark
Возвращает значение highWaterMark переданное при создании этого потока Writable.
writable.writableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение предоставляет данные для интроспекции состояния буфера highWaterMark.
writable.writableNeedDrain
Принимает значение true, если буфер потока заполнен, и поток испустит событие '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 когда источник Readable поток генерирует событие 'end', чтобы поток назначения перестал быть доступным для записи. Чтобы отключить это поведение по умолчанию, можно передать опцию 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<Поток> | <Итерируемый объект> | <Асинхронно итерируемый объект> | <Функция> -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Дуплексный> поток, составленный из потока
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<Функция> | <AsyncFunction> функция для обработки каждого фрагмента в потоке.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток разрушен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное выполнение вызоваfnдля потока. По умолчанию:1. -
highWaterMark<число> количество элементов для буферизации, пока ожидается обработка сопоставленных элементов пользователем. По умолчанию:concurrency * 2 - 1. -
signal<AbortSignal> позволяет разрушить поток, если сигнал прерван.
-
- Возвращает: <Чтение> поток, отображённый с помощью функции
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<Функция> | <AsyncFunction> функция для фильтрации фрагментов из потока.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток разрушен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное выполнение вызоваfnдля потока. По умолчанию:1. -
highWaterMark<число> количество элементов для буферизации, пока ожидается обработка отфильтрованных элементов пользователем. По умолчанию:concurrency * 2 - 1. -
signal<AbortSignal> позволяет разрушить поток, если сигнал прерван.
-
- Возвращает: <Чтение> поток, отфильтрованный с помощью предиката
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<Функция> | <AsyncFunction> функция для вызова для каждого фрагмента потока.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток разрушен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное выполнение вызоваfnдля потока. По умолчанию:1. -
signal<AbortSignal> позволяет разрушить поток, если сигнал прерван.
-
- Возвращает: <Промис> промис для завершения потока.
Этот метод позволяет итерировать поток. Для каждого фрагмента в потоке вызывается функция 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<AbortSignal> позволяет отменить операцию 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<Функция> | <AsyncFunction> функция для вызова для каждого фрагмента потока.-
data<любой тип> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток разрушен, что позволяет прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное выполнение вызоваfnдля потока. По умолчанию:1. -
signal<AbortSignal> позволяет разрушить поток, если сигнал прерван.
-
- Возвращает: <Промис> промис, возвращающий значение
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<Функция> | <AsyncФункция> функция, которая вызывается для каждого фрагмента потока.-
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<Функция> | <AsyncФункция> функция, которая вызывается для каждого фрагмента потока.-
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<Функция> | <AsyncГенераторФункция> | <AsyncФункция> функция для отображения каждого фрагмента в потоке.-
data<любое> фрагмент данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток уничтожается, позволяя прервать вызовfnраньше.
-
-
-
options<Объект>-
concurrency<число> максимальное одновременное вызовfnдля вызова в потоке сразу. По умолчанию:1. -
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Readable> поток, отображенный с помощью функции
fn.
Этот метод возвращает новый поток, применяя заданный обратный вызов к каждому фрагменту потока, а затем уплощая результат.
Возможен возврат потока или другого итерируемого или асинхронного итерируемого объекта из fn , и результирующие потоки будут объединены (уплощены) в возвращаемый поток.
import { Readable } from 'node:stream';
import { createReadStream } from 'node:fs';
// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).flatMap((x) => [x, x])) {
console.log(chunk); // 1, 1, 2, 2, 3, 3, 4, 4
}
// With an asynchronous mapper, combine the contents of 4 files
const concatResult = Readable.from([
'./1.mjs',
'./2.mjs',
'./3.mjs',
'./4.mjs',
]).flatMap((fileName) => createReadStream(fileName));
for await (const result of concatResult) {
// This will contain the contents (all chunks) of all 4 files
console.log(result);
} copy
readable.drop(limit[, options])
-
limit<число> количество фрагментов, которые нужно выбросить из readable. -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Readable> поток с выброшенными
limitфрагментами.
Этот метод возвращает новый поток с первыми limit фрагментами, выброшенными.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).drop(2).toArray(); // [3, 4] copy
readable.take(limit[, options])
-
limit<число> количество фрагментов для взятия из readable. -
options<Объект>-
signal<AbortSignal> позволяет уничтожить поток, если сигнал прерван.
-
- Возвращает: <Readable> поток с взятыми
limitфрагментами.
Этот метод возвращает новый поток с первыми limit фрагментами.
import { Readable } from 'node:stream';
await Readable.from([1, 2, 3, 4]).take(2).toArray(); // [1, 2] copy
readable.asIndexedPairs([options])
-
options<Объект>-
signal<AbortSignal> позволяет разрушить поток, если сигнал прерван.
-
- Возвращает: <Чтение> поток индексированных пар.
Этот метод возвращает новый поток с частями исходного потока, соединёнными с счётчиком в формате [index, chunk]. Первое значение индекса — 0, и оно увеличивается на 1 для каждой сгенерированной части.
import { Readable } from 'node:stream';
const pairs = await Readable.from(['a', 'b', 'c']).asIndexedPairs().toArray();
console.log(pairs); // [[0, 'a'], [1, 'b'], [2, 'c']] copy
readable.reduce(fn[, initial[, options]])
-
fn<Функция> | <АсинхроннаяФункция> функция-редуктор для вызова каждой части потока.-
previous<любой> значение, полученное от последнего вызоваfnили значениеinitial, если задано, или первая часть потока в противном случае. -
data<любой> часть данных из потока. -
options<Объект>-
signal<AbortSignal> прерывается, если поток разрушен, что позволяет прервать вызовfnраньше времени.
-
-
-
initial<любой> начальное значение для использования в редукции. -
options<Объект>-
signal<AbortSignal> позволяет разрушить поток, если сигнал прерван.
-
- Возвращает: <Обещание> обещание для конечного значения редукции.
Этот метод вызывает fn для каждой части потока в порядке, передавая результат вычисления по предыдущему элементу. Он возвращает обещание для конечного значения редукции.
Если значение initial не указано, первая часть потока используется в качестве начального значения. Если поток пуст, обещание отклоняется с TypeError со свойством кода ERR_INVALID_ARGS.
import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
const directoryPath = './src';
const filesInDir = await readdir(directoryPath);
const folderSize = await Readable.from(filesInDir)
.reduce(async (totalSize, file) => {
const { size } = await stat(join(directoryPath, file));
return totalSize + size;
}, 0);
console.log(folderSize); copy Функция-редуктор итерирует элементы потока по одному, что означает, что параметр concurrency или параллелизм отсутствуют. Чтобы выполнить reduce параллельно, можно выделить асинхронную функцию в метод readable.map.
import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';
const directoryPath = './src';
const filesInDir = await readdir(directoryPath);
const folderSize = await Readable.from(filesInDir)
.map((file) => stat(join(directoryPath, file)), { concurrency: 2 })
.reduce((totalSize, { size }) => totalSize + size, 0);
console.log(folderSize); copy Потоки типа «дуплекс» и «трансформация»
Класс: stream.Duplex
Потоки типа «дуплекс» — это потоки, которые реализуют как интерфейс Readable, так и интерфейс Writable.
Примеры потоков типа «дуплекс» включают:
duplex.allowHalfOpen
Если false , то поток автоматически завершит сторону записи, когда закончится сторона чтения. Изначально устанавливается опцией конструктора allowHalfOpen, которая по умолчанию равна true.
Это можно вручную изменить, чтобы изменить поведение полуоткрытого потока типа «дуплекс» существующего экземпляра потока, но это необходимо сделать до выдачи события 'end'.
Класс: stream.Transform
Потоки типа «трансформация» — это потоки типа «дуплекс» Duplex, где выходные данные каким-то образом связаны с входными. Как и все потоки типа «дуплекс» Duplex, потоки типа «трансформация» Transform реализуют интерфейсы Readable и Writable.
Примеры потоков типа «трансформация» включают:
transform.destroy([error])
Разрушить поток и, по желанию, выдать событие 'error'. После этого вызова поток типа «трансформация» освободит все внутренние ресурсы. Разработчики не должны переопределять этот метод, а вместо этого должны реализовать readable._destroy(). По умолчанию реализация _destroy() для Transform также генерирует 'close', если emitClose не установлено в false.
После вызова destroy(), все последующие вызовы будут бесполезны, и никакие другие ошибки, кроме ошибок _destroy(), не могут быть выведены как 'error'.
stream.finished(stream[, options], callback)
-
stream<Поток> | <ПотокЧтения> | <ПотокЗаписи> Поток чтения и/или записи/веб-поток. -
options<Объект>-
error<логическое> Если установлено вfalse, то вызовemit('error', err)не обрабатывается как завершённый. По умолчанию:true. -
readable<логическое> При установке вfalse, обратный вызов будет вызван при завершении потока, даже если поток всё ещё может быть читаемым. По умолчанию:true. -
writable<логическое> При установке вfalse, обратный вызов будет вызван при завершении потока, даже если поток всё ещё может быть записываемым. По умолчанию:true. -
signal<AbortSignal> позволяет прервать ожидание завершения потока. Базовый поток не прерывается, если сигнал прерван. Обратный вызов будет вызван сAbortError. Все зарегистрированные обработчики, добавленные этой функцией, также будут удалены. -
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', 'finish', 'end' и 'close') после вызова callback. Причина в том, чтобы неожиданные 'error' события (из-за неправильной реализации потоков) не приводили к непредвиденным сбоям. Если это нежелательное поведение, то возвращаемая функция очистки должна быть вызвана в обратном вызове:
const cleanup = finished(rs, (err) => {
cleanup();
// ...
}); copy
stream.pipeline(source[, ...transforms], destination, callback)
stream.pipeline(streams, callback)
-
streams<Поток[]> | <Итерируемый[]> | <АсинхронноИтерируемый[]> | <Функция[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> -
source<Поток> | <Итерируемый> | <АсинхронноИтерируемый> | <Функция> | <ReadableStream>- Возвращает: <Итерируемый> | <АсинхронноИтерируемый>
-
...transforms<Поток> | <Функция> | <TransformStream>-
source<АсинхронноИтерируемый> - Возвращает: <АсинхронноИтерируемый>
-
-
destination<Поток> | <Функция> | <WritableStream>-
source<АсинхронноИтерируемый> - Возвращает: <АсинхронноИтерируемый> | <Promise>
-
-
callback<Функция> Вызывается, когда конвейер полностью завершен.-
err<Ошибка> -
valЗначение, возвращаемоеPromiseвызовомdestination.
-
- Возвращает: <Поток>
Метод модуля для передачи данных между потоками и генераторами, перенаправляя ошибки и должным образом очищая и предоставляя обратный вызов при завершении конвейера.
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) для всех потоков, кроме:
- потоков, которые уже выпустили
'end'или'close'. - потоков, которые уже выпустили
'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<Поток[]> | <Итерируемый[]> | <АсинхронноИтерируемый[]> | <Функция[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> | <Duplex[]> | <Функция> - Возвращает: <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<Итерируемый> Объект, реализующийSymbol.asyncIteratorилиSymbol.iteratorитерируемый протокол. Генерирует событие 'error', если передано значение null. -
options<Объект> Опции, предоставленные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<Объект>-
encoding<строка> -
highWaterMark<число> -
objectMode<логическое значение> -
signal<AbortSignal>
-
- Возвращает: <stream.Readable>
stream.Readable.isDisturbed(stream)
-
stream<stream.Readable> | <ReadableStream> - Возвращает:
boolean
Возвращает значение, указывающее, была ли потоковая передача прочитана или отменена.
stream.isErrored(stream)
-
stream<Readable> | <Writable> | <Duplex> | <WritableStream> | <ReadableStream> - Возвращает: <логическое значение>
Возвращает значение, указывающее, произошла ли ошибка в потоке.
stream.isReadable(stream)
-
stream<Readable> | <Duplex> | <ReadableStream> - Возвращает: <логическое значение>
Возвращает значение, указывающее, является ли поток читаемым.
stream.Readable.toWeb(streamReadable[, options])
-
streamReadable<stream.Readable> -
options<Объект>-
strategy<Объект>-
highWaterMark<число> Максимальный размер внутренней очереди (созданногоReadableStream) перед применением обратной связи о переполнении при чтении из данногоstream.Readable. Если значение не указано, оно будет взято из данногоstream.Readable. -
size<Функция> Функция, определяющая размер данного фрагмента данных. Если значение не указано, размер будет1для всех фрагментов.
-
-
- Возвращает: <ReadableStream>
stream.Writable.fromWeb(writableStream[, options])
-
writableStream<WritableStream> -
options<Объект>-
decodeStrings<логическое значение> -
highWaterMark<число> -
objectMode<логическое значение> -
signal<AbortSignal>
-
- Возвращает: <stream.Writable>
stream.Writable.toWeb(streamWritable)
-
streamWritable<stream.Writable> - Возвращает: <WritableStream>
stream.Duplex.from(src)
-
src<Поток> | <Blob> | <ArrayBuffer> | <строка> | <Итерируемый объект> | <Асинхронно итерируемый объект> | <Функция асинхронного генератора> | <Асинхронная функция> | <Promise> | <Объект> | <ReadableStream> | <WritableStream>
Утилитарный метод для создания потоков duplex.
-
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<Объект>-
readable<ReadableStream> -
writable<WritableStream>
-
-
options<Объект> - Возвращает: <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)
Возвращает значение по умолчанию highWaterMark, используемое потоками. По умолчанию 16384 (16 Кбайт) или 16 для objectMode.
stream.setDefaultHighWaterMark(objectMode, value)
Устанавливает значение по умолчанию 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. По умолчанию:16384(16 КБ) или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. - и т.д.
-
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.
Эта необязательная функция будет вызвана в следующем цикле после возвращения конструктора потока, откладывая все вызовы _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> | <строка> | <любой> Данные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<Буфер> | <строка> Экземпляр буфера или строка, содержащая данные для записи.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:
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.
Пользовательские потоки на чтение обязательно должны вызывать конструктор new stream.Readable([options]) и реализовывать метод readable._read().
new stream.Readable([options])
-
options<Object>-
highWaterMark<число> Максимальное количество байтов, которые будут храниться во внутреннем буфере, прежде чем прекратится чтение из базового ресурса. По умолчанию:16384(16 КБ) или16дляobjectModeпотоков. -
encoding<строка> Если указано, буферы будут декодированы в строки с использованием указанной кодировки. По умолчанию:null. -
objectMode<логическое> Указывает, должен ли этот поток вести себя как поток объектов. Это означает, чтоstream.read(n)возвращает единственное значение вместоBufferразмераn. По умолчанию:false. -
emitClose<логическое> Указывает, должен ли поток генерировать'close'после уничтожения. По умолчанию:true. -
read<Функция> Реализация методаstream._read(). -
destroy<Функция> Реализация методаstream._destroy(). -
construct<Функция> Реализация методаstream._construct(). -
autoDestroy<логическое> Указывает, должен ли этот поток автоматически вызывать.destroy()на себе после завершения. По умолчанию:true. -
signal<AbortSignal> Сигнал, представляющий возможную отмену.
-
const { Readable } = require('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()) на создаваемом readable.
const { Readable } = require('node:stream');
const controller = new AbortController();
const read = new Readable({
read(size) {
// ...
},
signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort(); 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равноtrue. По умолчанию: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 Наиболее важной частью дуплексного потока является то, что стороны чтения Readable и записи Writable работают независимо друг от друга, несмотря на совместное существование в одном экземпляре объекта.
Дуплексные потоки в объектном режиме
Для потоков Duplex режим объектов objectMode может быть установлен исключительно для стороны чтения Readable или стороны записи Writable с помощью опций readableObjectMode и writableObjectMode соответственно.
Например, в следующем примере создается новый поток Transform (который является типом потока Duplex), у которого сторона чтения Writable в объектном режиме принимает JavaScript-числа, которые преобразуются в шестнадцатеричные строки на стороне Readable.
const { Transform } = require('node:stream');
// All Transform streams are also Duplex Streams.
const myTransform = new Transform({
writableObjectMode: true,
transform(chunk, encoding, callback) {
// Coerce the chunk to a number if necessary.
chunk |= 0;
// Transform the chunk into something else.
const data = chunk.toString(16);
// Push the data onto the readable queue.
callback(null, '0'.repeat(data.length % 2) + data);
},
});
myTransform.setEncoding('ascii');
myTransform.on('data', (chunk) => console.log(chunk));
myTransform.write(1);
// Prints: 01
myTransform.write(10);
// Prints: 0a
myTransform.write(100);
// Prints: 64 copy Реализация потока преобразования
Поток Transform — это дуплексный поток Duplex, где выход вычисляется каким-либо образом из входных данных. Примеры включают потоки zlib или crypto, которые сжимают, шифруют или дешифруют данные.
Нет требования, чтобы размер выходных данных был таким же, как у входных, чтобы количество блоков было одинаковым или чтобы они поступали одновременно. Например, поток Hash будет иметь только один блок вывода, который предоставляется при завершении входных данных. Поток zlib будет генерировать выходные данные, которые могут быть намного меньше или намного больше, чем входные.
Класс stream.Transform расширяется для реализации потока Transform.
Класс stream.Transform прототипически наследуется от stream.Duplex и реализует собственные версии методов writable._write() и readable._read(). Пользовательские реализации Transform обязаны реализовать метод transform._transform() и могут также реализовать метод transform._flush().
При использовании потоков преобразования необходимо учитывать, что данные, записанные в поток, могут привести к приостановке стороны Writable потока, если выходные данные на стороне Readable не потребляются.
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/dist/latest-v20.x/docs/api/stream.html