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