Поток[src]
Исходный код: lib/stream.js
Поток — это абстрактный интерфейс для работы со потоковой передачей данных в Node.js. Модуль stream предоставляет API для реализации интерфейса потока.
Node.js предоставляет множество объектов потока. Например, запрос к HTTP-серверу и process.stdout являются экземплярами потоков.
Потоки могут быть читаемыми, записываемыми или обоими. Все потоки являются экземплярами EventEmitter.
Для доступа к модулю stream:
const stream = require('stream'); Модуль stream полезен для создания новых типов экземпляров потоков. Обычно нет необходимости использовать модуль stream для потребления потоков.
Организация данного документа
Этот документ содержит два основных раздела и третий раздел для примечаний. Первый раздел объясняет, как использовать существующие потоки в приложении. Второй раздел объясняет, как создавать новые типы потоков.
Типы потоков
В Node.js существует четыре основных типа потоков:
-
Writable: потоки, в которые можно записывать данные (например,fs.createWriteStream()). -
Readable: потоки, из которых можно читать данные (например,fs.createReadStream()). -
Duplex: потоки, которые являются какReadable, так иWritable(например,net.Socket). -
Transform: потоки, которые могут изменять или преобразовывать данные при записи и чтении (например,zlib.createDeflate()).
Кроме того, этот модуль включает служебные функции stream.pipeline(), stream.finished() и stream.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(), является ограничение буферизации данных приемлемыми уровнями, чтобы источники и пункты назначения с различной скоростью не перегружали доступную память.
Параметр highWaterMark — это порог, а не ограничение: он определяет количество данных, которое поток буферизует, прежде чем перестанет запрашивать больше данных. Он не навязывает строгого ограничения на использование памяти в целом. Конкретные реализации потоков могут выбрать навязывание более строгих ограничений, но это необязательно.
Поскольку 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' , если опция autoDestroy не была установлена в значение true при создании потока.
Событие: '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'); Событие: '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' генерируется при вызове метода 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); 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' перед уничтожением потока. Разработчики не должны переопределять этот метод, но вместо этого реализовать writable._destroy().
writable.destroyed
Истинно после вызова writable.destroy().
writable.end([chunk[, encoding]][, callback])
-
chunk<строка> | <Буфер> | <Uint8Array> | <любое значение> Необязательные данные для записи. Для потоков, которые не работают в режиме объектов,chunkдолжно быть строкой,BufferилиUint8Array. Для потоков в режиме объектов,chunkможет быть любым значением JavaScript, кромеnull. -
encoding<Функция> Необязательный обратный вызов, когда поток завершается - Возвращает: <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
Истинно, если безопасно вызвать writable.write().
writable.writableEnded
Истинно после вызова writable.end(). Эта свойство не указывает, были ли данные переданы; для этого используйте writable.writableFinished.
writable.writableCorked
Количество вызовов writable.uncork(), необходимых для полного разблокирования потока.
writable.writableFinished
Устанавливается в значение true непосредственно перед генерацией события 'finish'.
writable.writableHighWaterMark
Возвращает значение highWaterMark , переданное при создании данного потока Writable.
writable.writableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение предоставляет данные для интроспекции состояния очереди highWaterMark.
writable.writableObjectMode
Метод получения свойства objectMode заданного потока Writable.
writable.write(chunk[, encoding][, callback])
-
chunk<string> | <Buffer> | <Uint8Array> | <any> Дополнительные данные для записи. Для потоков, не работающих в режиме работы с объектами,chunkдолжен быть строкой,BufferилиUint8Array. Для потоков, работающих в режиме работы с объектами,chunkможет быть любым значением JavaScript, кромеnull. -
encoding<string> Кодировка, еслиchunk— строка. По умолчанию:'utf8' -
callback<Function> Обратный вызов, когда этот фрагмент данных буферизуется - Возвращает: <boolean>
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 могут никогда не завершить обработку, если удалённый узел не считывает данные, запись в сокет, который не завершает обработку, может привести к удалённо эксплуатируемой уязвимости.
Запись данных, пока поток не завершает обработку, особенно проблематична для 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.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.
Событие: 'pause'
Событие 'pause' генерируется, когда вызывается stream.pause() и readableFlowing не false.
Событие: '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().
Событие: 'resume'
Событие 'resume' генерируется, когда вызывается stream.resume() и readableFlowing не true.
readable.destroy([error])
-
error<Ошибка> Ошибка, которая будет передана в качестве полезной нагрузки в событии'error' - Возвращает: <this>
Уничтожить поток. Необязательно, выпустить событие 'error' и выпустить событие 'close' (если emitClose установлено в false). После этого вызова читабельный поток освободит все внутренние ресурсы, и последующие вызовы push() будут игнорироваться. Реализаторы не должны переопределять этот метод, а вместо этого реализовать readable._destroy().
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
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.Writable> Назначение для записи данных -
options<Объект> Параметры конвейера-
end<boolean> Завершить писатель, когда читатель закончит. По умолчанию:true.
-
- Возвращает: <stream.Writable> Назначение, позволяющее цепочку конвейеров, если это
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<число> Необязательный аргумент для указания количества считываемых данных. - Возвращает: <строка> | <Буфер> | <null> | <любой>
Метод readable.read() извлекает данные из внутреннего буфера и возвращает их. Если данных для чтения нет, возвращается null. По умолчанию данные возвращаются в виде объекта Buffer, если не было указано кодирование с помощью метода readable.setEncoding() или поток работает в режиме объектов.
Необязательный аргумент size указывает конкретное количество байтов для чтения. Если не доступно size байтов, возвращается null, *за исключением* случая, когда поток завершен, в этом случае возвращаются все оставшиеся данные из внутреннего буфера.
Если аргумент size не указан, возвращаются все данные, содержащиеся во внутреннем буфере.
Аргумент size должен быть меньше или равен 1 ГБ.
Метод readable.read() должен вызываться только для потоков Readable в режиме паузы. В режиме потокового чтения метод readable.read() вызывается автоматически до тех пор, пока внутренний буфер не будет полностью опустошен.
const readable = getReadableStreamSomehow();
// 'readable' may be triggered multiple times as data is buffered in
readable.on('readable', () => {
let chunk;
console.log('Stream is readable (new data received in buffer)');
// Use a loop to make sure we read all currently available data
while (null !== (chunk = readable.read())) {
console.log(`Read ${chunk.length} bytes of data...`);
}
});
// 'end' will be triggered once when there is no more data available
readable.on('end', () => {
console.log('Reached end of stream.');
}); Каждый вызов 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('');
}); Поток Readable в режиме объектов всегда возвращает один элемент из вызова readable.read(size), независимо от значения аргумента size.
Если метод readable.read() возвращает фрагмент данных, то будет также выброшено событие 'data'.
Вызов stream.read([size]) после того, как было выброшено событие 'end', вернёт null. Ошибка во время выполнения не будет поднята.
readable.readable
Является true, если безопасно вызвать readable.read().
readable.readableEncoding
Геттер для свойства encoding данного потока Readable. Свойство encoding можно установить с помощью метода readable.setEncoding().
readable.readableEnded
Становится true, когда выброшено событие 'end'.
readable.readableFlowing
Это свойство отражает текущее состояние потока Readable как описано в разделе Три состояния.
readable.readableHighWaterMark
Возвращает значение highWaterMark , переданное при создании этого потока Readable.
readable.readableLength
Это свойство содержит количество байтов (или объектов) в очереди, готовых для чтения. Значение предоставляет данные для интроспекции относительно состояния потока highWaterMark.
readable.readableObjectMode
Геттер для свойства objectMode данного потока Readable.
readable.resume()
- Возвращает: <текущий>
Метод 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> Необязательный конкретный поток для отсоединения. - Возвращает: <текущий>
Метод 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[, encoding])
-
chunk<Буфер> | <Uint8Array> | <строка> | <null> | <любой> Фрагмент данных для вставки в очередь чтения. Для потоков, не работающих в режиме объектов,chunkдолжен быть строкой,Buffer,Uint8Arrayилиnull. Для потоков в режиме объектовchunkможет быть любым значением JavaScript. -
encoding<строка> Кодировка строковых фрагментов. Должна быть допустимой кодировкойBuffer, такой как'utf8'или'ascii'.
Передача chunk в качестве null сигнализирует об окончании потока (EOF) и ведет себя так же, как readable.push(null), после чего больше данных нельзя записать. Сигнал EOF помещается в конец буфера, и любые буферизованные данные все равно будут выгружены.
Метод readable.unshift() помещает фрагмент данных обратно во внутренний буфер. Это полезно в определенных ситуациях, когда поток потребляется кодом, которому нужно «отменить потребление» некоторого количества данных, которые он оптимистично извлек из источника, чтобы эти данные можно было передать другой стороне.
Метод stream.unshift(chunk) не может быть вызван после того, как событие 'end' было излучено, в противном случае будет брошено исключение времени выполнения.
Разработчики, использующие stream.unshift(), часто должны рассмотреть возможность переключения на использование потока Transform вместо этого. Дополнительная информация доступна в разделе API для разработчиков потоков.
// Pull off a header delimited by \n\n.
// Use unshift() if we get too much.
// Call the callback with (error, header, stream).
const { StringDecoder } = require('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 chunk of readable) {
data += chunk;
}
console.log(data);
}
print(fs.createReadStream('file')).catch(console.error); Если цикл завершается с 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'. После этого вызова поток transform высвободит все внутренние ресурсы. Разработчики не должны переопределять этот метод, а вместо этого реализовывать readable._destroy(). По умолчанию реализация _destroy() для Transform также излучает событие 'close', если emitClose не установлено в false.
stream.finished(stream[, options], callback)
-
stream<Поток> Чтение и/или запись потока. -
options<Объект>-
error<логическое значение> Если установлено вfalse, вызовemit('error', err)не рассматривается как завершённый. По умолчанию:true. -
readable<логическое значение> Если установлено вfalse, обратный вызов будет вызван при завершении потока, даже если поток по-прежнему может читаться. По умолчанию:true. -
writable<логическое значение> Если установлено вfalse, обратный вызов будет вызван при завершении потока, даже если поток по-прежнему может записывать данные. По умолчанию:true.
-
-
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 также допускает применение promisify.
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.finished() оставляет висящие обработчики событий (в частности, 'error', 'end', 'finish' и 'close') после вызова callback. Причина в том, чтобы неожиданные 'error' события (из-за неправильной реализации потоков) не приводили к неожиданным сбоям. Если это поведение нежелательно, необходимо вызвать возвращённую функцию очистки в обратном вызове:
const cleanup = finished(rs, (err) => {
cleanup();
// ...
}); stream.pipeline(...streams, callback)
-
...streams<Поток> Два или более потоков для соединения. -
callback<Функция> Вызывается, когда соединение потоков завершается.-
err<Ошибка>
-
Метод модуля для соединения потоков, перенаправляющий ошибки, должным образом очищающий и предоставляющий обратный вызов при завершении соединения.
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 также допускает применение promisify.
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); stream.pipeline() вызовет stream.destroy(err) для всех потоков, за исключением:
-
Readableпотоков, которые излучили события'end'или'close'. -
Writableпотоков, которые излучили события'finish'или'close'.
stream.pipeline() оставляет висящие обработчики событий в потоках после вызова callback. В случае повторного использования потоков после сбоя это может привести к утечкам обработчиков событий и необработанным ошибкам.
stream.Readable.from(iterable, [options])
-
iterable<Итерируемый объект> Объект, реализующий протокол итерированияSymbol.asyncIteratorилиSymbol.iterator. Излучает событие 'error', если передано значение null. -
options<Объект> Параметры, передаваемые вnew stream.Readable([options]). По умолчаниюReadable.from()установитoptions.objectModeвtrue, если явно не отключено, установивoptions.objectModeвfalse. - Возвращает: <stream.Readable>
Утилитарный метод для создания потоков чтения из итераторов.
const { Readable } = require('stream');
async function * generate() {
yield 'hello';
yield 'streams';
}
const readable = Readable.from(generate());
readable.on('data', (chunk) => {
console.log(chunk);
}); Вызов Readable.from(string) или Readable.from(buffer) не будет перебирать строки или буферы для соответствия семантике других потоков по причинам производительности.
API для разработчиков потоков
API модуля stream разработан для того, чтобы сделать реализацию потоков с использованием прототипного наследования JavaScript лёгкой.
Сначала разработчик потока объявит новый класс JavaScript, который расширяет один из четырёх основных классов потоков (stream.Writable, stream.Readable, stream.Duplex, или stream.Transform), убедившись, что они вызывают конструктор соответствующего родительского класса:
const { Writable } = require('stream');
class MyWritable extends Writable {
constructor({ highWaterMark, ...options }) {
super({
highWaterMark,
autoDestroy: true,
emitClose: true
});
// ...
}
} При расширении потоков, помните, какие параметры пользователь может и должен предоставить перед передачей их базовому конструктору. Например, если реализация делает предположения относительно параметров 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('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<Буфер> | <строка> | <любой> ДанныеBufferдля записи, преобразованные изstring, переданного вstream.write(). Если параметр потокаdecodeStringsравенfalseили поток работает в объектном режиме, фрагмент не преобразуется и будет таким, каким был передан вstream.write(). -
encoding<строка> Если фрагмент является строкой, тоencoding— это кодировка символов этой строки. Если фрагмент являетсяBuffer, или если поток работает в объектном режиме,encodingможет быть проигнорировано. -
callback<Функция> Вызовите эту функцию (при необходимости с аргументом ошибки) при завершении обработки предоставленного фрагмента.
Все реализации потоков Writable должны предоставлять метод writable._write() и/или writable._writev() для отправки данных в базовый ресурс.
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<Object[]> Части для записи. Каждая часть имеет следующий формат:{ chunk: ..., encoding: ... }. -
callback<Функция> Функция обратного вызова (по желанию с аргументом ошибки), которая вызывается по завершении обработки переданных частей.
Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только внутренними методами класса Writable.
Метод writable._writev() может быть реализован дополнительно или как альтернатива методу writable._write() в реализациях потоков, способных обрабатывать несколько частей данных одновременно. Если реализован и если есть буферизованные данные из предыдущих записей, будет вызван метод _writev() вместо _write().
Метод writable._writev() имеет префикс подчёркивания, потому что он внутренний для класса, который его определяет, и никогда не должен вызываться напрямую пользовательскими программами.
writable._destroy(err, callback)
-
err<Ошибка> Возможная ошибка. -
callback<Функция> Функция обратного вызова, принимающая необязательный аргумент ошибки.
Метод _destroy() вызывается методом writable.destroy(). Он может быть переопределён дочерними классами, но его нельзя вызывать напрямую.
writable._final(callback)
-
callback<Функция> Вызовите эту функцию (по желанию с аргументом ошибки) по завершении записи всех оставшихся данных.
Метод _final() нельзя вызывать напрямую. Он может быть реализован дочерними классами и, если реализован, будет вызываться только внутренними методами класса Writable.
Эта необязательная функция будет вызвана перед закрытием потока, откладывая событие 'finish' до вызова callback. Это полезно для закрытия ресурсов или записи буферизованных данных перед завершением потока.
Ошибки при записи
Ошибки, возникающие во время обработки методов writable._write(), writable._writev() и writable._final(), должны быть переданы вызовом обратного вызова и передачей ошибки в качестве первого аргумента. Бросание Error внутри этих методов или ручное отправление события 'error' приводит к неопределённому поведению.
Если поток 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. -
emitClose<логическое значение> Указывает, должен ли поток отправлять событие'close'после уничтожения. По умолчанию:true. -
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<Buffer> | <Uint8Array> | <строка> | <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._read().
Для потоков, не работающих в режиме объектов, если параметр chunk метода readable.push() имеет значение undefined, он будет интерпретирован как пустая строка или буфер. Подробнее см. readable.push('').
Ошибки при чтении
Ошибки, возникшие во время обработки метода readable._read(), должны быть обработаны с помощью метода readable.destroy(err). Выброс Error внутри метода readable._read() или ручное вызов события 'error' приводит к неопределённому поведению.
const { Readable } = require('stream');
const myReadable = new Readable({
read(size) {
const err = checkSomeErrorCondition();
if (err) {
this.destroy(err);
} else {
// 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 — это поток, который реализует как 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<логическое> Устанавливает, должен ли поток быть читаемым. По умолчанию:true. -
writable<логическое> Устанавливает, должен ли поток быть записываемым. По умолчанию: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 потока, который оборачивает гипотетический низкоуровневый объект, в который можно писать и из которого можно читать данные, хотя 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 можно установить исключительно для стороны чтения или записи, используя опции readableObjectMode и writableObjectMode соответственно.
Например, в следующем примере создаётся новый Transform поток (который является типом Duplex потока), у которого сторона чтения в режиме объектов принимает числа JavaScript, которые преобразуются в шестнадцатеричные строки на стороне записи.
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 — это дуплексный поток Duplex, где выходные данные вычисляются каким-либо образом из входных данных. Примеры включают потоки zlib или crypto, которые сжимают, шифруют или дешифруют данные.
Нет требования, чтобы выходные данные имели тот же размер, количество блоков или появлялись одновременно с входными данными. Например, поток Hash будет иметь только один блок выходных данных, который предоставляется при завершении входных данных. Поток zlib будет производить выходные данные, которые либо намного меньше, либо намного больше, чем входные.
Класс stream.Transform расширяется для реализации потока преобразования Transform.
Класс stream.Transform прототипически наследуется от stream.Duplex и реализует собственные версии методов writable._write() и readable._read(). Пользовательские реализации Transform должны реализовывать метод transform._transform() и могут также реализовывать метод transform._flush().
Следует соблюдать осторожность при использовании потоков преобразования, так как данные, записанные в поток, могут привести к приостановке стороны записи потока, если выходные данные на стороне чтения не обрабатываются.
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(), метод 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);
}; Метод 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 finished = util.promisify(stream.finished);
const writable = fs.createWriteStream('./file');
function drain(writable) {
if (writable.destroyed) {
return Promise.reject(new Error('premature close'));
}
return Promise.race([
once(writable, 'drain'),
once(writable, 'close')
.then(() => Promise.reject(new Error('premature close')))
]);
}
async function pump(iterable, writable) {
for await (const chunk of iterable) {
// Handle backpressure on write().
if (!writable.write(chunk)) {
await drain(writable);
}
}
writable.end();
}
(async function() {
// Ensure completion without errors.
await Promise.all([
pump(iterable, writable),
finished(writable)
]);
})(); В приведенном примере ошибки в write() будут перехвачены и выброшены слушателем once() события 'drain', так как once() также обрабатывает событие 'error'. Чтобы гарантировать завершение потока записи без ошибок, безопаснее использовать метод finished() как в примере, вместо использования слушателя для события 'finish'. В некоторых случаях событие 'error' может быть сгенерировано потоком записи после 'finish', и поскольку once() освободит обработчик 'error' при обработке события 'finish', это может привести к необработанной ошибке.
В качестве альтернативы, поток чтения можно обернуть с помощью Readable.from() и затем передать через .pipe():
const finished = util.promisify(stream.finished);
const writable = fs.createWriteStream('./file');
(async function() {
const readable = Readable.from(iterable);
readable.pipe(writable);
// Ensure completion without errors.
await finished(writable);
})(); Или, используя stream.pipeline() для перенаправления потоков:
const pipeline = util.promisify(stream.pipeline);
const writable = fs.createWriteStream('./file');
(async function() {
const readable = Readable.from(iterable);
await pipeline(readable, writable);
})(); Совместимость с более старыми версиями 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, и поток в данный момент не читает, то вызов 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-v12.x/docs/api/stream.html