Поток[src]
Поток — это абстрактный интерфейс для работы со с потоковыми данными в Node.js. Модуль stream предоставляет базовый API, который упрощает создание объектов, реализующих интерфейс потока.
Node.js предоставляет множество объектов потоков. Например, запрос к HTTP-серверу и process.stdout являются экземплярами потоков.
Потоки могут быть читаемыми, записываемыми или и тем, и другим. Все потоки являются экземплярами EventEmitter.
Модуль stream доступен с помощью:
const stream = require('stream');
Хотя важно понимать, как работают потоки, модуль stream сам по себе наиболее полезен для разработчиков, создающих новые типы экземпляров потоков. Разработчикам, которые в основном работают с потреблением объектов потоков, редко нужно использовать модуль stream напрямую.
Структура этого документа
Этот документ разделен на два основных раздела и третий раздел для дополнительных примечаний. В первом разделе объясняются элементы API потоков, необходимые для использования потоков в приложении. Во втором разделе объясняются элементы API, необходимые для реализации новых типов потоков.
Типы потоков
В Node.js существует четыре основных типа потоков:
-
Writable— потоки, в которые можно записывать данные (например,fs.createWriteStream()). -
Readable— потоки, из которых можно читать данные (например,fs.createReadStream()). -
Duplex— потоки, которые являются какReadable, так иWritable(например,net.Socket). -
Transform— потоки, которые могутDuplexданные по мере записи и чтения (например,zlib.createDeflate()).
Кроме того, этот модуль включает в себя вспомогательные функции pipeline, finished и Readable.from.
Режим работы с объектами
Все потоки, созданные API Node.js, работают только со строками и Buffer (или Uint8Array) объектами. Однако реализации потоков могут работать с другими типами значений JavaScript (за исключением null, которое имеет специальное назначение в потоках). Такие потоки считаются работающими в «режиме работы с объектами».
Экземпляры потоков переключаются в режим работы с объектами с помощью параметра objectMode при создании потока. Попытка переключения существующего потока в режим работы с объектами небезопасна.
Буферизация
Оба потока Writable и Readable будут хранить данные в внутреннем буфере, который можно получить с помощью writable.writableBuffer или readable.readableBuffer, соответственно.
Количество данных, потенциально буферизуемых, зависит от параметра 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(), является ограничение буферизации данных приемлемым уровнем, чтобы источники и получатели с различной скоростью не перегружали доступную память.
Поскольку потоки Duplex и Transform являются как Readable, так и Writable, каждый поддерживает два отдельных внутренних буфера для чтения и записи, что позволяет каждой стороне работать независимо от другой при одновременном поддержании эффективного потока данных. Например, экземпляры net.Socket являются потоками Duplex, чья сторона Readable позволяет потреблять данные, полученные из сокета, а чья сторона Writable позволяет записывать данные в сокет. Поскольку данные могут записываться в сокет с большей или меньшей скоростью, чем данные поступают, важно, чтобы каждая сторона работала (и буферизовалась) независимо от другой.
API для потребителей потоков
Практически все приложения Node.js, независимо от сложности, используют потоки каким-либо образом. Ниже приведен пример использования потоков в приложении Node.js, реализующем HTTP-сервер:
const http = require('http');
const server = http.createServer((req, res) => {
// req is an http.IncomingMessage, which is a Readable Stream
// res is an http.ServerResponse, which is a Writable Stream
let body = '';
// Get the data as utf8 strings.
// If an encoding is not set, Buffer objects will be received.
req.setEncoding('utf8');
// Readable streams emit 'data' events once a listener is added
req.on('data', (chunk) => {
body += chunk;
});
// the 'end' event indicates that the entire body has been received
req.on('end', () => {
try {
const data = JSON.parse(body);
// write back something interesting to the user:
res.write(typeof data);
res.end();
} catch (er) {
// uh oh! bad json!
res.statusCode = 400;
return res.end(`error: ${er.message}`);
}
});
});
server.listen(1337);
// $ curl localhost:1337 -d "{}"
// object
// $ curl localhost:1337 -d "\"foo\""
// string
// $ curl localhost:1337 -d "not json"
// error: Unexpected token o in JSON at position 1
Потоки Writable (например, res в примере) предоставляют методы, такие как write() и end(), которые используются для записи данных в поток.
Потоки Readable используют API EventEmitter для уведомления кода приложения, когда данные доступны для чтения из потока. Эти данные можно прочитать из потока различными способами.
И потоки Writable, и потоки Readable используют API EventEmitter различными способами для передачи текущего состояния потока.
Потоки Duplex и Transform являются как потоками Writable, так и потоками Readable.
Приложениям, которые либо записывают данные в поток, либо потребляют данные из потока, не требуется реализовывать интерфейсы потоков напрямую, и обычно у них нет причин вызывать require('stream').
Разработчики, желающие реализовать новые типы потоков, должны обратиться к разделу API для разработчиков потоков.
Потоки записи
Потоки записи — это абстракция для получателя, в который записываются данные.
Примеры потоков Writable включают:
- HTTP-запросы (клиент)
- HTTP-ответы (сервер)
- Потоки записи fs
- Потоки zlib
- Потоки 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');
Класс: stream.Writable
Событие: 'close'
Событие 'close' генерируется, когда поток и любые его базовые ресурсы (например, дескриптор файла) были закрыты. Событие указывает, что больше событий не будет, и дальнейшие вычисления не будут выполняться.
Поток Writable всегда будет генерировать событие 'close', если он был создан с параметром emitClose.
Событие: 'drain'
Если вызов stream.write(chunk) возвращает false, событие 'drain' будет сгенерировано, когда будет уместно возобновить запись данных в поток.
// Write the data to the supplied writable stream one million times.
// Be attentive to back-pressure.
function writeOneMillionTimes(writer, data, encoding, callback) {
let i = 1000000;
write();
function write() {
let ok = true;
do {
i--;
if (i === 0) {
// last time!
writer.write(data, encoding, callback);
} else {
// see if we should continue, or wait
// don't pass the callback, because we're not done yet.
ok = writer.write(data, encoding);
}
} while (i > 0 && ok);
if (i > 0) {
// had to stop early!
// write some more once it drains
writer.once('drain', write);
}
}
}
Событие: 'error'
Событие 'error' генерируется, если при записи или передаче данных произошла ошибка. При вызове обработчика события передаётся единственный аргумент Error.
Поток не закрывается при генерации события 'error'.
Событие: 'finish'
Событие 'finish' генерируется после вызова метода stream.end() и после того, как все данные были переданы в подлежащую систему.
const writer = getWritableStreamSomehow();
for (let i = 0; i < 100; i++) {
writer.write(`hello, #${i}!\n`);
}
writer.end('This is the end\n');
writer.on('finish', () => {
console.log('All writes are now complete.');
});
Событие: 'pipe'
-
srcПоток-источник <stream.Readable>, который передаёт данные в этот поток записи
Событие 'pipe' генерируется, когда для потока чтения вызывается метод stream.pipe(), добавляющий этот поток записи в список получателей.
const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('pipe', (src) => {
console.log('Something is piping into the writer.');
assert.equal(src, reader);
});
reader.pipe(writer);
Событие: 'unpipe'
-
srcПоток-источник <stream.Readable>, который отсоединил этот поток записи от себя с помощью unpipe
Событие 'unpipe' генерируется, когда для потока чтения вызывается метод stream.unpipe(), удаляя этот поток записи из списка получателей потока чтения.
Это событие также генерируется, если этот поток записи Writable генерирует ошибку, когда поток чтения Readable пытается передать данные в него.
const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
console.log('Something has stopped piping into the writer.');
assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer);
writable.cork()
Метод writable.cork() заставляет все записанные данные быть буферизованными в памяти. Буферизованные данные будут переданы в систему при вызове метода stream.uncork() или stream.end().
Основной целью writable.cork() является предотвращение ситуации, когда запись многих небольших фрагментов данных в поток не приводит к задержке во внутренней буферизации, что негативно скажется на производительности. В таких ситуациях реализации, использующие метод writable._writev(), могут выполнять буферизованную запись более эффективным способом.
См. также: writable.uncork().
writable.destroy([error])
Уничтожает поток и генерирует переданное 'error' и событие 'close'. После этого вызова поток записи завершен, и последующие вызовы write() или end() приведут к ошибке ERR_STREAM_DESTROYED. Реализаторы не должны переопределять этот метод, а вместо этого должны реализовать writable._destroy().
writable.end([chunk][, encoding][, callback])
-
chunk<строка> | <Буфер> | <Uint8Array> | <любой> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжен быть строкой,BufferилиUint8Array. Для потоков в режиме объектовchunkможет быть любым значением JavaScript, кромеnull. -
encoding<строка> Кодировка, еслиchunkявляется строкой -
callback<Функция> Необязательный обратный вызов, когда поток завершен - Возвращает: <this>
Вызов метода writable.end() сигнализирует о том, что больше данных в поток Writable не будет передаваться. Необязательные аргументы chunk и encoding позволяют записать последний фрагмент данных перед закрытием потока. Если указана, необязательная функция callback добавляется как обработчик события 'finish'.
Вызов метода stream.write() после вызова stream.end() вызовет ошибку.
// write 'hello, ' and then end with 'world!'
const fs = require('fs');
const file = fs.createWriteStream('example.txt');
file.write('hello, ');
file.end('world!');
// writing more now is not allowed!
writable.setDefaultEncoding(encoding)
Метод writable.setDefaultEncoding() устанавливает кодировку по умолчанию для потока Writable.
writable.uncork()
Метод writable.uncork() передаёт все данные, буферизованные с момента вызова stream.cork().
При использовании writable.cork() и writable.uncork() для управления буферизацией записей в поток, рекомендуется отложить вызовы writable.uncork() с использованием process.nextTick(). Это позволяет объединить все вызовы writable.write() в рамках одной фазы цикла событий Node.js.
stream.cork();
stream.write('some ');
stream.write('data ');
process.nextTick(() => stream.uncork());
Если метод writable.cork() вызывается несколько раз для потока, то для сброса буферизованных данных необходимо вызвать метод writable.uncork() такое же количество раз.
stream.cork();
stream.write('some ');
stream.cork();
stream.write('data ');
process.nextTick(() => {
stream.uncork();
// The data will not be flushed until uncork() is called a second time.
stream.uncork();
});
См. также: writable.cork().
writable.writable
Имеет значение true, если безопасно вызвать [writable.write()][].
writable.writableHighWaterMark
Возвращает значение highWaterMark, переданное при создании этого Writable.
writable.writableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение предоставляет данные для интроспекции состояния очереди highWaterMark.
writable.write(chunk[, encoding][, callback])
-
chunk<string> | <Buffer> | <Uint8Array> | <any> Данные для записи. Для потоков, не работающих в режиме объектов,chunkдолжно быть строкой,BufferилиUint8Array. Для потоков в режиме объектовchunkможет быть любым значением JavaScript, отличным отnull. -
encoding<string> Кодировка, еслиchunkявляется строкой -
callback<Функция> Обработчик события, когда этот фрагмент данных будет обработан - Возвращает: <логическое значение>
falseесли поток хочет, чтобы вызывающий код ожидал события'drain'перед продолжением записи дополнительных данных; в противном случаеtrue.
Метод writable.write() записывает данные в поток и вызывает переданный callback после того, как данные будут полностью обработаны. Если произошла ошибка, callback может или не может быть вызван с ошибкой в качестве первого аргумента. Для надёжного обнаружения ошибок записи добавьте обработчик события 'error'.
Значение возврата - true если внутренний буфер меньше highWaterMark , настроенного при создании потока после добавления chunk. Если возвращается false, дальнейшие попытки записи данных в поток должны быть остановлены до тех пор, пока не будет испущен 'drain' событие.
Пока поток не завершает обработку, вызовы write() будут буферизовать chunk, и возвращать false. Как только все текущие буферизованные фрагменты данных будут обработаны (приняты для передачи операционной системой), будет испущено событие 'drain' . Рекомендуется, чтобы после того, как write() вернул false, больше никаких фрагментов данных не записывалось до тех пор, пока не будет испущено событие 'drain' . Вызов write() для потока, не завершающего обработку, разрешён, но Node.js будет буферизовать все записанные фрагменты до достижения максимального использования памяти, после чего он безусловно прервёт работу. Даже до прерывания высокое использование памяти приведёт к плохой работе сборщика мусора и высокой RSS (которая обычно не возвращается системе, даже после того, как память больше не требуется). Поскольку сокеты TCP могут никогда не завершить обработку, если удалённый peer не читает данные, запись в сокет, который не завершает обработку, может привести к удалённо эксплуатируемой уязвимости.
Запись данных, пока поток не завершает обработку, особенно проблематична для Transform, поскольку потоки Transform приостановлены по умолчанию до тех пор, пока они не будут переданы или не будет добавлен обработчик событий 'data' или 'readable'.
Если данные, которые нужно записать, могут быть сгенерированы или получены по требованию, рекомендуется инкапсулировать логику в Readable и использовать stream.pipe(). Однако если вызов write() предпочтительнее, можно уважать backpressure и избежать проблем с памятью, используя событие 'drain':
function write(data, cb) {
if (!stream.write(data)) {
stream.once('drain', cb);
} else {
process.nextTick(cb);
}
}
// Wait for cb to be called before doing any other write.
write('hello', () => {
console.log('Write completed, do more writes now.');
});
Поток Writable в режиме объектов всегда игнорирует аргумент encoding.
Потоки чтения
Потоки чтения — это абстракция источника, из которого потребляются данные.
Примеры потоков Readable:
- Ответы HTTP на стороне клиента
- Запросы HTTP на сервере
- Потоки чтения fs
- Потоки zlib
- Потоки crypto
- Сокеты TCP
- stdout и stderr дочернего процесса
process.stdin
Все потоки Readable реализуют интерфейс, определённый классом stream.Readable.
Два режима чтения
Потоки Readable эффективно работают в одном из двух режимов: потоковый и приостановленный. Эти режимы отделены от режима объектов. Поток Readable может быть в режиме объектов или нет, независимо от того, в потоковом режиме он находится или приостановлен.
-
В потоковом режиме данные считываются из базовой системы автоматически и предоставляются приложению как можно быстрее с помощью событий через интерфейс
EventEmitter. -
В режиме паузы метод
stream.read()должен быть вызван явно для чтения фрагментов данных из потока.
Все потоки Readable начинаются в режиме паузы, но могут быть переключены в потоковый режим следующим образом:
- Добавление обработчика события
'data'. - Вызов метода
stream.resume(). - Вызов метода
stream.pipe()для передачи данных вWritable.
Поток может переключиться обратно в приостановленный режим следующим образом:
- Если нет мест назначения для передачи, вызвав метод
stream.pause(). - Если есть места назначения для передачи, удалив все места назначения. Несколько мест назначения могут быть удалены, вызвав метод
stream.unpipe().
Важно помнить, что поток не сгенерирует данные до тех пор, пока не будет предоставлен механизм либо потребления, либо игнорирования этих данных. Если механизм потребления отключён или удалён, поток попытается прекратить генерацию данных.
Для обратной совместимости удаление обработчиков событий 'data' не автоматически приостановит поток. Также, если есть места назначения для передачи, вызов stream.pause() не гарантирует, что поток останется приостановленным, после того как эти места назначения завершат обработку и запросят новые данные.
Если поток Readable переключается в потоковый режим, и нет потребителей для обработки данных, эти данные будут потеряны. Это может произойти, например, при вызове метода readable.resume() без обработчика события 'data', или когда обработчик события 'data' удалён из потока.
Добавление обработчика события 'readable' автоматически приостанавливает поток, и данные потребляются через readable.read(). Если обработчик события 'readable' удалён, то поток начнёт потоковый режим снова, если есть обработчик события 'data'.
Три состояния
Два "режима" работы потока Readable — это упрощённая абстракция для более сложного внутреннего управления состоянием в реализации потока Readable.
Конкретно, в любой момент времени каждый поток Readable находится в одном из трёх возможных состояний:
readable.readableFlowing === nullreadable.readableFlowing === falsereadable.readableFlowing === true
Когда поток readable.readableFlowing находится в null, механизм потребления данных потока отсутствует. Поэтому поток не генерирует данные. В этом состоянии добавление обработчика события 'data', вызов метода readable.pipe() или метода readable.resume() переключат readable.readableFlowing в true, заставив поток Readable начать активно генерировать события по мере генерации данных.
Вызов readable.pause(), readable.unpipe() или получение backpressure установят состояние readable.readableFlowing как false, временно приостановив поток событий, но не приостанавливая генерацию данных. В этом состоянии добавление обработчика события 'data' не переключится readable.readableFlowing на true.
const { PassThrough, Writable } = require('stream');
const pass = new PassThrough();
const writable = new Writable();
pass.pipe(writable);
pass.unpipe(writable);
// readableFlowing is now false
pass.on('data', (chunk) => { console.log(chunk.toString()); });
pass.write('ok'); // will not emit 'data'
pass.resume(); // must be called to make stream emit 'data'
Когда readable.readableFlowing находится в false, данные могут накапливаться во внутреннем буфере потока.
Выберите один стиль API
API потоков Readable эволюционировал на протяжении нескольких версий Node.js и предоставляет несколько способов потребления данных потоков. В общем случае разработчики должны выбрать один способ потребления данных и никогда не использовать несколько методов для потребления данных из одного потока. В частности, использование комбинации on('data'), on('readable'), pipe() или асинхронных итераторов может привести к неинтуитивному поведению.
Использование метода readable.pipe() рекомендуется для большинства пользователей, поскольку он реализован для обеспечения наилучшего способа потребления данных потока. Разработчики, которым требуется более тонкий контроль над передачей и генерацией данных, могут использовать EventEmitter и readable.on('readable')/readable.read() или API readable.pause()/readable.resume().
Класс: stream.Readable
Событие: 'close'
Событие 'close' генерируется, когда поток и любые его базовые ресурсы (например, дескриптор файла) закрыты. Событие указывает, что больше событий не будет генерироваться и дальнейшие вычисления не будут производиться.
Поток Readable всегда будет генерировать событие 'close' , если он создан с опцией emitClose.
Событие: 'data'
-
chunk<Буфер> | <строка> | <любой> Чанк данных. Для потоков, работающих не в объектном режиме, чанк будет либо строкой, либоBuffer. Для потоков в объектном режиме чанк может быть любым значением JavaScript, кромеnull.
Событие 'data' генерируется всякий раз, когда поток отказывается от владения куском данных потребителю. Это может происходить всякий раз, когда поток переключается в режим потока вызовом readable.pipe(), readable.resume(), или при добавлении обработчика обратного вызова к событию 'data'. Событие 'data' также будет генерироваться всякий раз, когда вызывается метод readable.read() и доступен чанк данных для возврата.
Прикрепление обработчика события 'data' к потоку, который не был явно приостановлен, переключает поток в режим потока. Данные будут передаваться сразу по их доступности.
Обработчик обратного вызова будет получать чанк данных в виде строки, если для потока задано значение кодировки по умолчанию с помощью метода readable.setEncoding(); в противном случае данные будут переданы как Buffer.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
});
Событие: 'end'
Событие 'end' генерируется, когда больше нет данных для потребления из потока.
Событие 'end' не будет генерироваться, если данные не полностью потреблены. Это можно сделать, переключив поток в режим потока или вызвав stream.read() многократно, пока все данные не будут потреблены.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
});
readable.on('end', () => {
console.log('There will be no more data.');
});
Событие: 'error'
Событие 'error' может генерироваться реализацией Readable в любое время. Обычно это может произойти, если базовый поток не может сгенерировать данные из-за внутренней ошибки или когда реализация потока пытается передать недопустимый чанк данных.
Обработчик обратного вызова получит объект Error.
Событие: 'readable'
Событие 'readable' генерируется, когда данные доступны для чтения из потока. В некоторых случаях прикрепление обработчика события 'readable' приведет к чтению некоторого количества данных в внутренний буфер.
const readable = getReadableStreamSomehow();
readable.on('readable', function() {
// there is some data to read now
let data;
while (data = this.read()) {
console.log(data);
}
});
Событие 'readable' также будет генерироваться, когда конец данных потока достигнут, но до генерации события 'end'.
По сути, событие 'readable' указывает, что поток содержит новую информацию: доступны новые данные или достигнут конец потока. В первом случае stream.read() вернет доступные данные. Во втором случае stream.read() вернет null. Например, в следующем примере foo.txt — это пустой файл:
const fs = require('fs');
const rr = fs.createReadStream('foo.txt');
rr.on('readable', () => {
console.log(`readable: ${rr.read()}`);
});
rr.on('end', () => {
console.log('end');
});
Вывод выполнения этого скрипта:
$ node test.js readable: null end
В целом, механизмы событий readable.pipe() и 'data' легче понять, чем событие 'readable'. Однако обработка события 'readable' может привести к увеличению пропускной способности.
Если одновременно используются события 'readable' и 'data', событие 'readable' имеет приоритет в управлении потоком, т.е. событие 'data' будет генерироваться только при вызове stream.read(). Свойство readableFlowing станет false. Если есть 'data' обработчика при удалении 'readable', поток начнёт генерировать события 'data' без вызова .resume().
readable.destroy([error])
-
error<Ошибка> Ошибка, которая будет передана в качестве полезной нагрузки в событии'error' - Возвращает: <this>
Уничтожить поток и сгенерировать события 'error' и 'close'. После этого вызова, потоковое чтение освободит все внутренние ресурсы, и последующие вызовы к push() будут игнорироваться. Реализаторы не должны переопределять этот метод, а вместо этого реализовать readable._destroy().
readable.isPaused()
- Возвращает: <логическое>
Метод readable.isPaused() возвращает текущее рабочее состояние Readable. Он используется в основном механизмом, лежащим в основе метода readable.pipe(). В большинстве типичных случаев нет необходимости использовать этот метод напрямую.
const readable = new stream.Readable(); readable.isPaused(); // === false readable.pause(); readable.isPaused(); // === true readable.resume(); readable.isPaused(); // === false
readable.pause()
- Возвращает: <this>
Метод readable.pause() заставит поток в режиме потока прекратить генерировать события 'data', переключившись из режима потока. Любые доступные данные останутся во внутреннем буфере.
const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
console.log(`Received ${chunk.length} bytes of data.`);
readable.pause();
console.log('There will be no additional data for 1 second.');
setTimeout(() => {
console.log('Now data will start flowing again.');
readable.resume();
}, 1000);
});
Метод readable.pause() не имеет эффекта, если есть обработчик события 'readable'.
readable.pipe(destination[, options])
-
destination<stream.Писабельный> Назначение для записи данных -
options<Объект> Опции конвейера-
end<логическое> Завершить запись при завершении чтения. По умолчанию:true.
-
- Возвращает: <stream.Писабельный> Назначение, позволяющее организовать цепочку конвейеров, если это поток
DuplexилиTransform
Метод readable.pipe() присоединяет поток Writable к readable, вызывая автоматический переход в режим потока и отправку всех данных в присоединённый поток Writable. Поток данных будет автоматически управляться таким образом, чтобы целевой поток Writable не перегружался более быстрым потоком Readable.
Следующий пример передает все данные из readable в файл с именем file.txt:
const fs = require('fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt'
readable.pipe(writable);
Возможно подключение нескольких потоков Writable к одному потоку Readable.
Метод readable.pipe() возвращает ссылку на целевой поток, позволяя создавать цепочки конвейеров потоков:
const fs = require('fs');
const r = fs.createReadStream('file.txt');
const z = zlib.createGzip();
const w = fs.createWriteStream('file.txt.gz');
r.pipe(z).pipe(w);
По умолчанию, stream.end() вызывается на целевом потоке Writable, когда исходный поток Readable генерирует событие 'end', чтобы целевой поток больше не был доступен для записи. Чтобы отключить это поведение по умолчанию, опция end может быть передана как false, что приведет к сохранению открытым целевого потока:
reader.pipe(writer, { end: false });
reader.on('end', () => {
writer.end('Goodbye\n');
});
Важный момент: если поток Readable генерирует ошибку во время обработки, целевой поток Writable не закрывается автоматически. В случае ошибки необходимо руконечно закрыть каждый поток, чтобы предотвратить утечки памяти.
Потоки process.stderr и process.stdout Writable никогда не закрываются до выхода процесса Node.js, независимо от указанных опций.
readable.read([size])
-
size<number> Дополнительный аргумент для указания количества считываемых данных. - Возвращает: <string> | <Buffer> | <null> | <any>
Метод readable.read() извлекает данные из внутреннего буфера и возвращает их. Если данных для чтения нет, возвращается null. По умолчанию данные возвращаются как объект Buffer, если не задано кодирование с помощью метода readable.setEncoding() или поток не работает в объектном режиме.
Дополнительный аргумент size задаёт конкретное количество байтов для чтения. Если size байтов недоступно для чтения, возвращается null, *за исключением* случая, когда поток завершён, в этом случае возвращаются все данные, оставшиеся во внутреннем буфере.
Если аргумент size не указан, возвращаются все данные, содержащиеся во внутреннем буфере.
Аргумент size должен быть меньше или равен 1 ГБ.
Метод readable.read() должен вызываться только для потоков Readable в режиме приостановки. В режиме потоковой передачи readable.read() вызывается автоматически, пока внутренний буфер не будет полностью исчерпан.
const readable = getReadableStreamSomehow();
readable.on('readable', () => {
let chunk;
while (null !== (chunk = readable.read())) {
console.log(`Received ${chunk.length} bytes of data.`);
}
});
Обратите внимание, что цикл while необходим при обработке данных с помощью readable.read(). Только после того, как readable.read() вернёт null, 'readable' будет отправлен.
Поток Readable в объектном режиме всегда возвращает один элемент при вызове readable.read(size), независимо от значения аргумента size.
Если метод readable.read() возвращает фрагмент данных, также будет отправлено событие 'data'.
Вызов stream.read([size]) после отправки события 'end' вернёт null. Ошибка выполнения не будет поднята.
readable.readable
Истинно, если безопасно вызывать [readable.read()][].
readable.readableFlowing
Это свойство отражает текущее состояние потока Readable как описано в разделе Состояния потока.
readable.readableHighWaterMark
Возвращает значение highWaterMark , переданное при создании этого Readable.
readable.readableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к чтению. Значение предоставляет данные интроспекции относительно состояния highWaterMark.
readable.resume()
- Возвращает: <this>
Метод readable.resume() заставляет явно приостановленный поток Readable возобновить отправку событий 'data', переключая поток в режим потоковой передачи.
Метод readable.resume() можно использовать для полного потребления данных из потока, не обрабатывая сами данные:
getReadableStreamSomehow()
.resume()
.on('end', () => {
console.log('Reached the end, but did not read anything.');
});
Метод readable.resume() не имеет эффекта, если существует слушатель события 'readable'.
readable.setEncoding(encoding)
Метод readable.setEncoding() устанавливает кодировку символов для данных, считанных из потока Readable.
По умолчанию кодировка не задана, и данные потока будут возвращаться как объекты Buffer. Установка кодировки приводит к возвращению данных потока в виде строк заданной кодировки, а не как объектов Buffer. Например, вызов readable.setEncoding('utf8') приведет к интерпретации выходных данных как данных UTF-8 и передаче их как строк. Вызов readable.setEncoding('hex') приведет к кодированию данных в шестнадцатеричном строковом формате.
Поток Readable корректно обрабатывает многобайтовые символы, передаваемые через поток, которые в противном случае были бы неправильно декодированы, если бы просто извлекались из потока как объекты Buffer.
const readable = getReadableStreamSomehow();
readable.setEncoding('utf8');
readable.on('data', (chunk) => {
assert.equal(typeof chunk, 'string');
console.log('Got %d characters of string data:', chunk.length);
});
readable.unpipe([destination])
-
destination<stream.Writable> Дополнительный конкретный поток для отсоединения - Возвращает: <this>
Метод readable.unpipe() отсоединяет поток Writable , ранее подключённый с помощью метода stream.pipe().
Если destination не указан, то все подключения отсоединяются.
Если destination указан, но подключение для него не установлено, метод ничего не делает.
const fs = require('fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt',
// but only for the first second
readable.pipe(writable);
setTimeout(() => {
console.log('Stop writing to file.txt.');
readable.unpipe(writable);
console.log('Manually close the file stream.');
writable.end();
}, 1000);
readable.unshift(chunk)
-
chunk<Buffer> | <Uint8Array> | <string> | <any> Фрагмент данных, который необходимо добавить в начало очереди чтения. Для потоков, не работающих в объектном режиме,chunkдолжно быть строкой,BufferилиUint8Array. Для потоков в объектном режимеchunkможет быть любым значением JavaScript, кромеnull.
Метод readable.unshift() помещает фрагмент данных обратно во внутренний буфер. Это полезно в некоторых ситуациях, когда поток потребляется кодом, которому необходимо «отменить» потребление некоторого количества данных, которые он оптимистично извлёк из источника, чтобы данные могли быть переданы другой стороне.
Метод stream.unshift(chunk) нельзя вызывать после отправки события 'end', в противном случае произойдёт ошибка выполнения.
Разработчикам, использующим stream.unshift() , следует рассмотреть возможность переключения на использование потока Transform. См. раздел API для разработчиков потоков для получения дополнительной информации.
// Pull off a header delimited by \n\n
// use unshift() if we get too much
// Call the callback with (error, header, stream)
const { StringDecoder } = require('string_decoder');
function parseHeader(stream, callback) {
stream.on('error', callback);
stream.on('readable', onReadable);
const decoder = new StringDecoder('utf8');
let header = '';
function onReadable() {
let chunk;
while (null !== (chunk = stream.read())) {
const str = decoder.write(chunk);
if (str.match(/\n\n/)) {
// found the header boundary
const split = str.split(/\n\n/);
header += split.shift();
const remaining = split.join('\n\n');
const buf = Buffer.from(remaining, 'utf8');
stream.removeListener('error', callback);
// remove the 'readable' listener before unshifting
stream.removeListener('readable', onReadable);
if (buf.length)
stream.unshift(buf);
// now the body of the message can be read from the stream.
callback(null, header, stream);
} else {
// still reading the header.
header += str;
}
}
}
}
В отличие от stream.push(chunk), stream.unshift(chunk) не завершит процесс чтения, сбросив внутреннее состояние чтения потока. Это может привести к неожиданным результатам, если readable.unshift() вызывается во время чтения (например, из реализации stream._read() в пользовательском потоке). Вызов readable.unshift() сразу же после stream.push('') сбросит состояние чтения должным образом, однако лучше всего просто избегать вызова readable.unshift() во время чтения.
readable.wrap(stream)
До Node.js 0.10 потоки не реализовывали весь API модуля stream, как он определён сейчас. (См. Совместимость для получения дополнительной информации.)
При использовании более старой библиотеки Node.js, которая отправляет события 'data' и имеет метод stream.pause(), который является только рекомендательным, метод readable.wrap() можно использовать для создания потока Readable, который использует старый поток в качестве источника данных.
Редко потребуется использовать readable.wrap(), но метод предоставлен для удобства взаимодействия со старыми приложениями и библиотеками Node.js.
const { OldReader } = require('./old-api-module.js');
const { Readable } = require('stream');
const oreader = new OldReader();
const myReader = new Readable().wrap(oreader);
myReader.on('readable', () => {
myReader.read(); // etc.
});
readable[Symbol.asyncIterator]()
- Возвращает: <AsyncIterator> для полного потребления потока.
const fs = require('fs');
async function print(readable) {
readable.setEncoding('utf8');
let data = '';
for await (const k of readable) {
data += k;
}
console.log(data);
}
print(fs.createReadStream('file')).catch(console.log);
Если цикл завершается с break или throw, поток будет уничтожен. Другими словами, итерирование по потоку полностью потребляет поток. Поток будет считываться частями размером, равным значению параметра highWaterMark. В приведенном выше примере данные будут в одном блоке, если файл содержит меньше 64 КБ данных, так как параметр highWaterMark не предоставлен для fs.createReadStream().
Потоки Duplex и Transform
Класс: stream.Duplex
Потоки Duplex — это потоки, которые реализуют оба интерфейса Readable и Writable.
Примеры потоков Duplex включают:
Класс: stream.Transform
Потоки Transform — это Duplex потоки, где выходные данные каким-либо образом связаны с входными. Как и все Duplex потоки, потоки Transform реализуют оба интерфейса Readable и Writable.
Примеры потоков Transform включают:
transform.destroy([error])
-
error<Ошибка>
Уничтожить поток и выпустить 'error'. После этого вызова поток transform высвободит все внутренние ресурсы. Реализаторы не должны переопределять этот метод, а вместо этого реализовать readable._destroy(). По умолчанию реализация _destroy() для Transform также выпустит 'close'.
stream.finished(stream, callback)
-
stream<Поток> Поток чтения и/или записи. -
callback<Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.
Функция, уведомляющая о том, что поток больше не читается, не записывается или столкнулся с ошибкой или преждевременным закрытием.
const { finished } = require('stream');
const rs = fs.createReadStream('archive.tar');
finished(rs, (err) => {
if (err) {
console.error('Stream failed.', err);
} else {
console.log('Stream is done reading.');
}
});
rs.resume(); // drain the stream
Особенно полезно в сценариях обработки ошибок, когда поток уничтожается преждевременно (например, при прерывании HTTP-запроса) и не выпустит 'end' или 'finish'.
API finished также можно использовать с промисами;
const finished = util.promisify(stream.finished);
const rs = fs.createReadStream('archive.tar');
async function run() {
await finished(rs);
console.log('Stream is done reading.');
}
run().catch(console.error);
rs.resume(); // drain the stream
stream.pipeline(...streams[, callback])
-
...streams<Поток> Два или более потоков для соединения. -
callback<Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.
Метод модуля для соединения потоков, передающий ошибки и должным образом очищающий потоки, а также предоставляющий обратный вызов при завершении конвейера.
const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('zlib');
// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.
// A pipeline to gzip a potentially huge tar file efficiently:
pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz'),
(err) => {
if (err) {
console.error('Pipeline failed.', err);
} else {
console.log('Pipeline succeeded.');
}
}
);
API pipeline также можно использовать с промисами:
const pipeline = util.promisify(stream.pipeline);
async function run() {
await pipeline(
fs.createReadStream('archive.tar'),
zlib.createGzip(),
fs.createWriteStream('archive.tar.gz')
);
console.log('Pipeline succeeded.');
}
run().catch(console.error);
Readable.from(iterable, [options])
-
iterable<Итерируемый объект> Объект, реализующий протокол итерацииSymbol.asyncIteratorилиSymbol.iterator. -
options<Объект> Параметры, предоставляемые дляnew stream.Readable([options]). По умолчаниюReadable.from()установитoptions.objectModeвtrue, если это не отменено явно, установивoptions.objectModeвfalse.
Утилита для создания потоков Readable из итераторов.
const { Readable } = require('stream');
async function * generate() {
yield 'hello';
yield 'streams';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
});
API для разработчиков потоков
API модуля stream разработан для того, чтобы сделать лёгкой реализацию потоков с использованием прототипного наследования JavaScript.
Сначала разработчик потоков объявит новый класс JavaScript, который расширяет один из четырёх основных классов потоков (stream.Writable, stream.Readable, stream.Duplex, или stream.Transform), убедившись, что они вызывают соответствующий конструктор родительского класса:
const { Writable } = require('stream');
class MyWritable extends Writable {
constructor(options) {
super(options);
// ...
}
}
Затем новый класс потоков должен реализовать один или несколько конкретных методов, в зависимости от типа создаваемого потока, как подробно указано в таблице ниже:
| Сценарий использования | Класс | Реализуемый(ые) метод(ы) |
|---|---|---|
| Только чтение | Readable |
_read |
| Только запись | Writable |
_write, _writev, _final
|
| Чтение и запись | Duplex |
_read, _write, _writev, _final
|
| Обработка записанных данных, затем чтение результата | Transform |
_transform, _flush, _final
|
Код реализации потока никогда не должен вызывать «публичные» методы потока, предназначенные для использования потребителями (как описано в разделе «API для потребителей потоков»). Это может привести к нежелательным побочным эффектам в прикладном коде, использующем поток.
Упрощённое создание
Во многих простых случаях можно создать поток, не используя наследование. Это можно сделать путём непосредственного создания экземпляров stream.Writable, stream.Readable, stream.Duplex или stream.Transform объектов и передачи соответствующих методов в качестве параметров конструктора.
const { Writable } = require('stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
}
});
Реализация потока записи
Класс stream.Writable расширяется для реализации потока Writable.
Пользовательские потоки Writable должны вызывать конструктор new stream.Writable([options]) и реализовывать метод writable._write(). Метод writable._writev() также может быть реализован.
Конструктор: new stream.Writable([options])
-
options<Объект>-
highWaterMark<число> Уровень буфера, когдаstream.write()начинает возвращатьfalse. По умолчанию:16384(16кб) или16дляobjectModeпотоков. -
decodeStrings<логическое> Необходимо ли кодироватьstringпередаваемые вstream.write()вBuffer(с кодировкой, указанной в вызовеstream.write()) перед передачей вstream._write(). Данные других типов не преобразуются (т.е.Bufferне декодируются вstring). Установка в значение false предотвратит преобразованиеstring. По умолчанию:true. -
defaultEncoding<строка> По умолчанию используется кодировка, когда кодировка не указана в качестве аргумента кstream.write(). По умолчанию:'utf8'. -
objectMode<логическое> Является ли операцияstream.write(anyObj)валидной. При установке в значение true, становится возможным писать JavaScript-значения, отличные от строк,BufferилиUint8Array(если это поддерживается реализацией потока). По умолчанию:false. -
emitClose<логическое> Должен ли поток генерировать'close'после уничтожения. По умолчанию:true. -
write<Функция> Реализация методаstream._write(). -
writev<Функция> Реализация методаstream._writev(). -
destroy<Функция> Реализация методаstream._destroy(). -
final<Функция> Реализация методаstream._final(). -
autoDestroy<логическое> Должен ли данный поток автоматически вызывать.destroy()на себе после завершения. По умолчанию:false.
-
const { Writable } = require('stream');
class MyWritable extends Writable {
constructor(options) {
// Calls the stream.Writable() constructor
super(options);
// ...
}
}
Или, при использовании конструкторов в стиле до ES6:
const { Writable } = require('stream');
const util = require('util');
function MyWritable(options) {
if (!(this instanceof MyWritable))
return new MyWritable(options);
Writable.call(this, options);
}
util.inherits(MyWritable, Writable);
Или, используя упрощённый подход к конструктору:
const { Writable } = require('stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
// ...
},
writev(chunks, callback) {
// ...
}
});
writable._write(chunk, encoding, callback)
-
chunk<Буфер> | <строка> | <любой> Данные, подлежащие записи, преобразованные изstringпереданных вstream.write(). Если опция потокаdecodeStringsимеет значениеfalse, или поток работает в режиме объектов, чанк не будет преобразован и будет иметь то же значение, что и переданное вstream.write(). -
encoding<строка> Если чанк является строкой, тоencodingпредставляет собой кодировку символов этой строки. Если чанк являетсяBuffer, или если поток работает в режиме объектов,encodingможет быть проигнорировано. -
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки) по завершении обработки предоставленного чанка.
Все реализации потоков Writable должны предоставлять метод writable._write() для отправки данных на подлежащий ресурс.
Transform потоки предоставляют собственную реализацию метода writable._write().
Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только внутренними методами класса Writable.
Необходимо вызвать метод callback для указания того, завершилась ли запись успешно или произошла ошибка. Первый аргумент, передаваемый в 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: ..., encoding: ... }. -
callback<Функция> Функция обратного вызова (при необходимости с аргументом ошибки), которая вызывается по завершении обработки предоставленных чанков.
Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только внутренними методами класса Writable.
Метод writable._writev() может быть реализован дополнительно к writable._write() в реализациях потоков, способных обрабатывать несколько чанков данных одновременно. Если реализован, метод будет вызван со всеми чанками данных, в настоящее время буферизованными в очереди записи.
Метод writable._writev() имеет префикс подчеркивания, так как является внутренним для определяющего его класса и не должен вызываться напрямую пользовательскими программами.
writable._destroy(err, callback)
-
err<Ошибка> Возможная ошибка. -
callback<Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.
Метод _destroy() вызывается методом writable.destroy(). Он может быть переопределён дочерними классами, но не должен вызываться напрямую.
writable._final(callback)
-
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки) по завершении записи всех оставшихся данных.
Метод _final() не должен вызываться напрямую. Он может быть реализован дочерними классами и, если реализован, будет вызываться только внутренними методами класса Writable.
Эта необязательная функция вызывается перед закрытием потока, откладывая событие 'finish' до вызова callback. Это полезно для закрытия ресурсов или записи буферизованных данных перед завершением потока.
Возникновение ошибок при записи
Рекомендуется, чтобы ошибки, возникающие во время обработки методов writable._write() и writable._writev(), сообщались путём вызова обратного вызова и передачи ошибки в качестве первого аргумента. Это вызовет событие 'error' для Writable. Бросание исключения Error изнутри writable._write() может привести к непредсказуемому и несогласованному поведению, в зависимости от того, как используется поток. Использование обратного вызова гарантирует согласованную и предсказуемую обработку ошибок.
Если поток Readable подключается к потоку Writable, когда Writable генерирует ошибку, поток Readable будет отключён.
const { Writable } = require('stream');
const myWritable = new Writable({
write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
}
});
Пример потока записи
Следующий пример демонстрирует довольно упрощенную (и несколько бесполезную) реализацию потока Writable. Хотя этот конкретный Writable экземпляр потока не обладает реальной практической ценностью, пример иллюстрирует каждый из необходимых элементов экземпляра пользовательского потока Writable:
const { Writable } = require('stream');
class MyWritable extends Writable {
_write(chunk, encoding, callback) {
if (chunk.toString().indexOf('a') >= 0) {
callback(new Error('chunk is invalid'));
} else {
callback();
}
}
}
Декодирование буферов в потоке записи
Декодирование буферов — распространенная задача, например, при использовании преобразователей, входными данными которых является строка. Это нетривиальный процесс при использовании кодировки символов с несколькими байтами, такой как UTF-8. Следующий пример демонстрирует, как декодировать многобайтовые строки, используя StringDecoder и Writable.
const { Writable } = require('stream');
const { StringDecoder } = require('string_decoder');
class StringWritable extends Writable {
constructor(options) {
super(options);
this._decoder = new StringDecoder(options && options.defaultEncoding);
this.data = '';
}
_write(chunk, encoding, callback) {
if (encoding === 'buffer') {
chunk = this._decoder.write(chunk);
}
this.data += chunk;
callback();
}
_final(callback) {
this.data += this._decoder.end();
callback();
}
}
const euro = [[0xE2, 0x82], [0xAC]].map(Buffer.from);
const w = new StringWritable();
w.write('currency: ');
w.write(euro[0]);
w.end(euro[1]);
console.log(w.data); // currency: €
Реализация потока чтения
Класс stream.Readable расширяется для реализации потока Readable.
Пользовательские потоки Readable обязаны вызывать конструктор new stream.Readable([options]) и реализовывать метод readable._read().
new stream.Readable([options])
-
options<Объект>-
highWaterMark<число> Максимальное количество байтов для хранения во внутреннем буфере перед прекращением чтения из базового ресурса. По умолчанию:16384(16 КБ) или16для потоковobjectMode. -
encoding<строка> Если указано, буферы будут декодированы в строки с использованием указанной кодировки. По умолчанию:null. -
objectMode<логическое> Нужно ли этому потоку вести себя как потоку объектов. Это означает, чтоstream.read(n)возвращает одно значение вместоBufferразмеромn. По умолчанию:false. -
read<Функция> Реализация методаstream._read(). -
destroy<Функция> Реализация методаstream._destroy(). -
autoDestroy<логическое> Нужно ли этому потоку автоматически вызывать.destroy()после завершения. По умолчанию:false.
-
const { Readable } = require('stream');
class MyReadable extends Readable {
constructor(options) {
// Calls the stream.Readable(options) constructor
super(options);
// ...
}
}
Или, при использовании конструкторов в стиле до ES6:
const { Readable } = require('stream');
const util = require('util');
function MyReadable(options) {
if (!(this instanceof MyReadable))
return new MyReadable(options);
Readable.call(this, options);
}
util.inherits(MyReadable, Readable);
Или, используя упрощенный подход к конструктору:
const { Readable } = require('stream');
const myReadable = new Readable({
read(size) {
// ...
}
});
readable._read(size)
-
size<число> Количество байтов для асинхронного чтения
Этот метод НЕ должен вызываться напрямую кодом приложения. Он должен реализовываться дочерними классами и вызываться только методами внутреннего класса Readable.
Все реализации потоков Readable должны предоставить реализацию метода readable._read() для получения данных из базового ресурса.
Когда вызывается readable._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 в качестве null сигнализирует об окончании потока (EOF), после чего больше данных нельзя записать.
Когда поток Readable приостановлен, данные, добавленные с помощью readable.push(), могут быть прочитаны, вызвав метод readable.read(), когда будет выведен событие 'readable'.
Когда поток Readable находится в режиме потока, данные, добавленные с помощью readable.push(), будут переданы путём выдачи события 'data'.
Метод readable.push() разработан для максимальной гибкости. Например, при обёртке низкоуровневого источника, предоставляющего механизм приостановки/возобновления и обратный вызов данных, низкоуровневый источник может быть обернут экземпляром пользовательского потока Readable:
// source is an object with readStop() and readStart() methods,
// and an `ondata` member that gets called when it has data, and
// an `onend` member that gets called when the data is over.
class SourceWrapper extends Readable {
constructor(options) {
super(options);
this._source = getLowlevelSourceObject();
// Every time there's data, push it into the internal buffer.
this._source.ondata = (chunk) => {
// if push() returns false, then stop reading from source
if (!this.push(chunk))
this._source.readStop();
};
// When the source ends, push the EOF-signaling `null` chunk
this._source.onend = () => {
this.push(null);
};
}
// _read will be called when the stream wants to pull more data in
// the advisory size argument is ignored in this case.
_read(size) {
this._source.readStart();
}
}
Метод readable.push() предназначен для вызова только реализаторами Readable и только изнутри метода readable._read().
Для потоков, не работающих в режиме объектов, если параметр chunk метода readable.push() имеет значение undefined, он будет обработан как пустая строка или буфер. См. readable.push('') для получения дополнительной информации.
Ошибки при чтении
Рекомендуется, чтобы ошибки, возникающие во время обработки метода readable._read(), передавались с помощью события 'error', а не выбрасывались. Выбрасывание ошибки Error изнутри readable._read() может привести к непредсказуемому и несовместимому поведению в зависимости от того, работает ли поток в режиме потока или приостановки. Использование события 'error' гарантирует согласованную и предсказуемую обработку ошибок.
const { Readable } = require('stream');
const myReadable = new Readable({
read(size) {
if (checkSomeErrorCondition()) {
process.nextTick(() => this.emit('error', err));
return;
}
// do some work
}
});
Пример счётного потока
Следующий пример — это базовый пример потока Readable , который выводит числа от 1 до 1 000 000 в порядке возрастания, а затем завершается.
const { Readable } = require('stream');
class Counter extends Readable {
constructor(opt) {
super(opt);
this._max = 1000000;
this._index = 1;
}
_read() {
const i = this._index++;
if (i > this._max)
this.push(null);
else {
const str = String(i);
const buf = Buffer.from(str, 'ascii');
this.push(buf);
}
}
}
Реализация потока типа Duplex
Поток типа 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. -
readableObjectMode<логическое значение> УстанавливаетobjectModeдля стороны чтения потока. Не имеет эффекта, еслиobjectModeустановлено вtrue. По умолчанию:false. -
writableObjectMode<логическое значение> УстанавливаетobjectModeдля стороны записи потока. Не имеет эффекта, еслиobjectModeустановлено вtrue. По умолчанию:false. -
readableHighWaterMark<число> УстанавливаетhighWaterMarkдля стороны чтения потока. Не имеет эффекта, еслиhighWaterMarkуказано. -
writableHighWaterMark<число> УстанавливаетhighWaterMarkдля стороны записи потока. Не имеет эффекта, еслиhighWaterMarkуказано.
-
const { Duplex } = require('stream');
class MyDuplex extends Duplex {
constructor(options) {
super(options);
// ...
}
}
Или, при использовании конструкторов в стиле до ES6:
const { Duplex } = require('stream');
const util = require('util');
function MyDuplex(options) {
if (!(this instanceof MyDuplex))
return new MyDuplex(options);
Duplex.call(this, options);
}
util.inherits(MyDuplex, Duplex);
Или, используя упрощённый подход к конструктору:
const { Duplex } = require('stream');
const myDuplex = new Duplex({
read(size) {
// ...
},
write(chunk, encoding, callback) {
// ...
}
});
Пример потока типа Duplex
Следующий пример демонстрирует простой пример потока типа Duplex , который оборачивает гипотетический объект нижнего уровня, к которому можно записать данные и из которого можно прочитать данные, но с API, несовместимым с потоками Node.js. Следующий пример демонстрирует простой пример потока Duplex, который буферизует входящие данные через интерфейс Writable, который затем считывается через интерфейс Readable.
const { Duplex } = require('stream');
const kSource = Symbol('source');
class MyDuplex extends Duplex {
constructor(source, options) {
super(options);
this[kSource] = source;
}
_write(chunk, encoding, callback) {
// The underlying source only deals with strings
if (Buffer.isBuffer(chunk))
chunk = chunk.toString();
this[kSource].writeSomeData(chunk);
callback();
}
_read(size) {
this[kSource].fetchSomeData(size, (data, encoding) => {
this.push(Buffer.from(data, encoding));
});
}
}
Самая важная особенность потока Duplex заключается в том, что стороны Readable и Writable работают независимо друг от друга, несмотря на совместное существование в одном экземпляре объекта.
Потоки типа Duplex в режиме объектов
Для потоков в режиме объектов Duplex можно установить objectMode только для стороны Readable или Writable с помощью опций readableObjectMode и writableObjectMode соответственно.
Например, в следующем примере создаётся новый поток типа Transform (который является типом потока Duplex), у которого сторона Writable в режиме объектов принимает числа JavaScript, которые преобразуются в шестнадцатеричные строки на стороне Readable.
const { Transform } = require('stream');
// All Transform streams are also Duplex Streams
const myTransform = new Transform({
writableObjectMode: true,
transform(chunk, encoding, callback) {
// Coerce the chunk to a number if necessary
chunk |= 0;
// Transform the chunk into something else.
const data = chunk.toString(16);
// Push the data onto the readable queue.
callback(null, '0'.repeat(data.length % 2) + data);
}
});
myTransform.setEncoding('ascii');
myTransform.on('data', (chunk) => console.log(chunk));
myTransform.write(1);
// Prints: 01
myTransform.write(10);
// Prints: 0a
myTransform.write(100);
// Prints: 64
Реализация потока типа Transform
Поток типа Transform — это поток типа Duplex, где вывод каким-то образом вычисляется из входных данных. Примеры включают потоки zlib или потоки crypto, которые сжимают, шифруют или расшифровывают данные.
Не требуется, чтобы размер вывода был таким же, как у входных данных, число фрагментов или время их появления были одинаковы. Например, поток Hash будет иметь только один фрагмент вывода, который предоставляется при завершении входа. Поток zlib будет генерировать вывод, который может быть значительно меньше или значительно больше, чем его входные данные.
Класс stream.Transform расширяется для реализации потока типа Transform.
Класс stream.Transform прототипически наследуется от stream.Duplex и реализует свои версии методов writable._write() и readable._read(). Пользовательские реализации Transform должны реализовывать метод transform._transform() и могут также реализовывать метод transform._flush().
При работе с потоками Transform необходимо учитывать, что данные, записанные в поток, могут привести к приостановке стороны Writable потока, если вывод на стороне Readable не потребляется.
new stream.Transform([options])
-
options<Объект> Передаётся как аргумент конструкторамWritableиReadable. Также содержит следующие поля:-
transform<Функция> Реализация методаstream._transform(). -
flush<Функция> Реализация методаstream._flush().
-
const { Transform } = require('stream');
class MyTransform extends Transform {
constructor(options) {
super(options);
// ...
}
}
Или, при использовании конструкторов в стиле до ES6:
const { Transform } = require('stream');
const util = require('util');
function MyTransform(options) {
if (!(this instanceof MyTransform))
return new MyTransform(options);
Transform.call(this, options);
}
util.inherits(MyTransform, Transform);
Или, используя упрощённый подход к конструктору:
const { Transform } = require('stream');
const myTransform = new Transform({
transform(chunk, encoding, callback) {
// ...
}
});
События: 'finish' и 'end'
События 'finish' и 'end' исходят от классов stream.Writable и stream.Readable соответственно. Событие 'finish' генерируется после вызова stream.end() и обработки всех фрагментов методом stream._transform(). Событие 'end' генерируется после вывода всех данных, что происходит после вызова обратного вызова в transform._flush().
transform._flush(callback)
-
callback<Функция> Функция обратного вызова (необязательно с аргументом ошибки и данных), которая вызывается при сбросе оставшихся данных.
Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только внутренними методами класса Readable.
В некоторых случаях для преобразования операции может потребоваться сгенерировать дополнительный фрагмент данных в конце потока. Например, поток сжатия zlib сохраняет определённое внутреннее состояние, используемое для оптимального сжатия вывода. Однако при завершении потока эти дополнительные данные необходимо сбросить, чтобы данные сжатия были полными.
Пользовательские реализации Transform могут реализовать метод transform._flush(). Он вызывается, когда больше нет данных для записи, но до вывода события 'end', сигнализирующего о конце потока Readable.
В реализации transform._flush() метод readable.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() обрабатывает записываемые байты, вычисляет выходные данные, а затем передает эти выходные данные в читаемую часть, используя метод readable.push().
Метод transform.push() может быть вызван ноль или более раз для генерации выходных данных из одного входного фрагмента, в зависимости от того, сколько выходных данных должно быть сгенерировано в результате этого фрагмента.
Возможна ситуация, когда из входных данных не генерируется никаких выходных данных.
Функция callback должна вызываться только тогда, когда текущий фрагмент полностью обработан. Первый аргумент, переданный в callback, должен быть объектом Error, если при обработке входных данных произошла ошибка, или null в противном случае. Если второй аргумент передается в callback, он будет передан методу readable.push().
transform.prototype._transform = function(data, encoding, callback) {
this.push(data);
callback();
};
transform.prototype._transform = function(data, encoding, callback) {
callback(null, data);
};
Метод transform._transform() имеет префикс подчеркивания, так как он является внутренним для класса, который его определяет, и его никогда не следует вызывать напрямую программами пользователя.
transform._transform() никогда не вызывается параллельно; потоки реализуют механизм очереди, и для получения следующего фрагмента необходимо вызвать callback — либо синхронно, либо асинхронно.
Класс: stream.PassThrough
Класс stream.PassThrough — тривиальная реализация потока Transform, который просто передает входные байты на выход. Он предназначен в основном для примеров и тестирования, но в некоторых случаях stream.PassThrough может быть полезным строительным блоком для новых типов потоков.
Дополнительные заметки
Совместимость потоков с асинхронными генераторами и итераторами
С поддержкой асинхронных генераторов и итераторов в JavaScript, асинхронные генераторы в настоящее время являются полноценной конструкцией потока на уровне языка.
Ниже приведены некоторые распространенные случаи взаимодействия использования потоков Node.js с асинхронными генераторами и итераторами.
Обработка потоков чтения с помощью асинхронных итераторов
(async function() {
for await (const chunk of readable) {
console.log(chunk);
}
})();
Создание потоков чтения с помощью асинхронных генераторов
Можно создать поток чтения Node.js из асинхронного генератора, используя утилиту Readable.from:
const { Readable } = require('stream');
async function * generate() {
yield 'a';
yield 'b';
yield 'c';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
});
Перенаправление в потоки записи из асинхронных итераторов
При записи в поток записи из асинхронного итератора важно обеспечить правильную обработку обратной связи и ошибок.
const { once } = require('events');
const writeable = fs.createWriteStream('./file');
(async function() {
for await (const chunk of iterator) {
// Handle backpressure on write
if (!writeable.write(value))
await once(writeable, 'drain');
}
writeable.end();
// Ensure completion without errors
await once(writeable, 'finish');
})();
В приведенном выше примере ошибки в потоке записи будут перехвачены и выброшены двумя слушателями once, так как once также обрабатывает события 'error'.
В качестве альтернативы, поток чтения можно обернуть в Readable.from и затем перенаправить через .pipe:
const { once } = require('events');
const writeable = fs.createWriteStream('./file');
(async function() {
const readable = Readable.from(iterator);
readable.pipe(writeable);
// Ensure completion without errors
await once(writeable, 'finish');
})();
Совместимость со старыми версиями Node.js
До Node.js 0.10 интерфейс потока Readable был проще, но также менее мощным и менее полезным.
- Вместо ожидания вызовов метода
stream.read(), события'data'начинали генерироваться сразу. Приложениям, которым нужно было выполнить некоторую работу для определения обработки данных, требовалось хранить считанные данные в буферах, чтобы данные не были потеряны. - Метод
stream.pause()был рекомендательным, а не гарантированным. Это означало, что всё ещё нужно было быть готовым к получению событий'data'даже когда поток был в приостановленном состоянии.
В Node.js 0.10 был добавлен класс Readable. Для обратной совместимости со старыми программами Node.js, потоки Readable переключаются в режим «поточного чтения» при добавлении обработчика события 'data' или при вызове метода stream.resume(). В результате, даже без использования нового метода stream.read() и события 'readable', больше не нужно беспокоиться о потере фрагментов 'data'.
Хотя большинство приложений будут продолжать работать нормально, это вносит особенность в следующих условиях:
- Не добавляется обработчик события
'data'. - Метод
stream.resume()никогда не вызывается. - Поток не перенаправлен ни на какое место назначения записи.
Например, рассмотрим следующий код:
// WARNING! BROKEN!
net.createServer((socket) => {
// we add an 'end' listener, but never consume the data
socket.on('end', () => {
// It will never get here.
socket.end('The message was received but was not processed.\n');
});
}).listen(1337);
До Node.js 0.10 входящие данные сообщения просто отбрасывались. Однако в Node.js 0.10 и выше сокет остается приостановленным навсегда.
Обходным решением в этой ситуации является вызов метода stream.resume() для запуска потока данных:
// Workaround
net.createServer((socket) => {
socket.on('end', () => {
socket.end('The message was received but was not processed.\n');
});
// start the flow of data, discarding it.
socket.resume();
}).listen(1337);
Помимо новых потоков Readable, переключающихся в режим потокового чтения, потоки в стиле до 0.10 могут быть обернуты в класс Readable с помощью метода readable.wrap().
readable.read(0)
В некоторых случаях необходимо инициировать обновление механизмов чтения подлежащего потока чтения, не потребляя никаких данных. В таких случаях можно вызвать readable.read(0), которая всегда вернет null.
Если внутренний буфер чтения находится ниже highWaterMark, и поток в настоящее время не читает, вызов readable.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-v10.x/docs/api/stream.html