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