Spec-Zone.ru › Node.js 8 LTS

Поток

Устойчивость: 2 - Стабильно

Поток — это абстрактный интерфейс для работы со струйными данными в Node.js. Модуль stream предоставляет базовый API, который упрощает создание объектов, реализующих интерфейс потока.

Node.js предоставляет множество объектов потоков. Например, запрос к HTTP-серверу и process.stdout являются экземплярами потока.

Потоки могут быть читаемыми, записываемыми или обоими. Все потоки являются экземплярами EventEmitter.

Модуль stream можно получить, используя:

const stream = require('stream');

Хотя важно понимать, как работают потоки, модуль stream сам по себе наиболее полезен для разработчиков, создающих новые типы экземпляров потоков. Разработчикам, которые в основном используют объекты потоков, обычно не нужно напрямую использовать модуль stream.

Организация этого документа

Этот документ разделен на два основных раздела и третий раздел с дополнительными заметками. Первый раздел описывает элементы API потока, необходимые для использования потоков в приложении. Второй раздел описывает элементы API, необходимые для реализации новых типов потоков.

Типы потоков

В Node.js существует четыре основных типа потоков:

  • Читаемый - потоки, из которых можно читать данные (например, fs.createReadStream()).
  • Записываемый - потоки, в которые можно записывать данные (например, fs.createWriteStream()).
  • Дуплексный - потоки, которые являются одновременно читаемыми и записываемыми (например, net.Socket).
  • Преобразующий - дуплексные потоки, которые могут изменять или преобразовывать данные при их записи и чтении (например, zlib.createDeflate()).

Режим объектов

Все потоки, созданные API Node.js, работают исключительно с строками и Buffer (или Uint8Array) объектами. Однако реализация потоков может работать с другими типами значений JavaScript (за исключением null, которое служит специальной цели в потоках). Такие потоки считаются работающими в «режиме объектов».

Экземпляры потоков переключаются в режим объектов с помощью параметра objectMode при создании потока. Попытка переключения уже существующего потока в режим объектов небезопасна.

Буферизация

Как записываемые, так и читаемые потоки будут хранить данные во внутренней буферной памяти, которую можно получить с помощью writable._writableState.getBuffer() или readable._readableState.buffer, соответственно.

Объем потенциально буферизуемых данных зависит от параметра highWaterMark, переданного конструктору потока. Для обычных потоков параметр highWaterMark определяет общее количество байтов. Для потоков, работающих в режиме объектов, параметр highWaterMark определяет общее количество объектов.

Данные буферизуются в читаемых потоках, когда реализация вызывает stream.push(chunk). Если потребитель потока не вызывает stream.read(), данные будут находиться в внутренней очереди до тех пор, пока не будут использованы.

Когда общий размер внутренней буфера чтения достигает порога, указанного параметром highWaterMark, поток временно прекращает чтение данных из базового ресурса до тех пор, пока буферизованные данные не будут использованы (то есть, поток прекратит вызывать внутренний метод readable._read() для заполнения буфера чтения).

Данные буферизуются в записываемых потоках, когда метод writable.write(chunk) вызывается многократно. Пока общий размер внутренней буфера записи ниже порога, заданного параметром highWaterMark, вызовы метода writable.write() будут возвращать true. Как только размер внутренней буфера достигнет или превысит highWaterMark, будет возвращено значение false.

Ключевой целью API stream, особенно метода stream.pipe(), является ограничение буферизации данных приемлемыми уровнями, чтобы источники и пункты назначения с различными скоростями не перегружали доступную память.

Поскольку потоки дуплексные и преобразующие являются одновременно читаемыми и записываемыми, каждый из них поддерживает два отдельных внутренних буфера, используемых для чтения и записи, что позволяет каждой стороне работать независимо друг от друга, сохраняя при этом надлежащий и эффективный поток данных. Например, экземпляры net.Socket являются дуплексными потоками, чья читаемая сторона позволяет потреблять данные, полученные из сокета, а записываемая сторона позволяет записывать данные в сокет. Поскольку данные могут записываться в сокет с большей или меньшей скоростью, чем данные принимаются, важно, чтобы каждая сторона работала (и буферизовала) независимо от другой.

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

Записываемые потоки (например, res в примере) предоставляют методы, такие как write() и end(), которые используются для записи данных в поток.

Читаемые потоки используют API EventEmitter для уведомления кода приложения о том, что данные доступны для чтения из потока. Доступные данные могут быть прочитаны из потока различными способами.

Как записываемые, так и читаемые потоки используют API EventEmitter различными способами для передачи текущего состояния потока.

Дуплексные и преобразующие потоки являются как записываемыми, так и читаемыми.

Приложениям, которые либо записывают данные в поток, либо потребляют данные из потока, не нужно напрямую реализовывать интерфейсы потоков, и обычно у них нет причин вызывать require('stream').

Разработчики, желающие реализовать новые типы потоков, должны обратиться к разделу API для разработчиков потоков.

Потоки записи

Потоки записи — это абстракция пункта назначения, в который записываются данные.

Примеры записываемых потоков включают:

  • HTTP-запросы (на стороне клиента)
  • HTTP-ответы (на стороне сервера)
  • Потоки записи fs
  • Потоки zlib
  • Потоки crypto
  • TCP-сокеты
  • стандартный ввод дочернего процесса
  • process.stdout, process.stderr

Примечание: Некоторые из этих примеров на самом деле являются дуплексными потоками, реализующими интерфейс записываемого потока.

Все записываемые потоки реализуют интерфейс, определенный классом stream.Writable.

Хотя отдельные экземпляры записываемых потоков могут отличаться по различным параметрам, все записываемые потоки следуют одному и тому же основному шаблону использования, как показано в примере ниже:

const myStream = getWritableStreamSomehow();
myStream.write('some data');
myStream.write('some more data');
myStream.end('done writing data');

Класс: stream.Writable

Добавлен в: v0.9.4
Событие: 'close'
Добавлен в: v0.9.4

Событие 'close' генерируется, когда поток и любые его базовые ресурсы (например, дескриптор файла) были закрыты. Событие указывает, что больше событий не будут генерироваться, и дальнейшие вычисления не будут происходить.

Не все потоки записи будут генерировать событие 'close'.

Событие: 'drain'
Добавлен в: v0.9.4

Если вызов 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'
Добавлен в: v0.9.4
  • <Ошибка>

Событие 'error' генерируется, если во время записи или передачи данных произошла ошибка. Обработчик события получает единственный аргумент Error.

Примечание: поток не закрывается при генерации события 'error'.

Событие: 'finish'
Добавлен в: v0.9.4

Событие '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.error('All writes are now complete.');
});
Событие: 'pipe'
Добавлен в: v0.9.4
  • src <stream.Readable> источник потока, который передается в этот записываемый поток

Событие 'pipe' генерируется, когда на читаемом потоке вызывается метод stream.pipe(), добавляя этот записываемый поток в список пунктов назначения.

const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('pipe', (src) => {
  console.error('something is piping into the writer');
  assert.equal(src, reader);
});
reader.pipe(writer);
Событие: 'unpipe'
Добавлен в: v0.9.4
  • src <stream.Readable> Исходной поток, из которого разорвана эта запись

Событие 'unpipe' генерируется, когда метод stream.unpipe() вызывается на потоке Readable, удаляя эту запись Writable из набора пунктов назначения.

Это событие также генерируется в случае, если этот поток Writable генерирует ошибку, когда поток Readable подключается к нему.

const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
  console.error('Something has stopped piping into the writer.');
  assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer);
writable.cork()
Добавлен в: v0.11.2

Метод writable.cork() принудительно буферизует все записанные данные в памяти. Буферизованные данные будут сброшены, когда вызовут либо метод stream.uncork(), либо stream.end().

Основной целью writable.cork() является предотвращение ситуации, когда запись многих небольших блоков данных в поток не приводит к задержке в внутреннем буфере, что негативно скажется на производительности. В таких ситуациях реализации, которые реализуют метод writable._writev(), могут выполнять буферизованные записи более оптимизированным способом.

См. также: writable.uncork().

writable.end([chunk][, encoding][, callback])
История
Версия Изменения
v8.0.0

Аргумент chunk теперь может быть экземпляром Uint8Array.

v0.9.4

Добавлен в: v0.9.4

  • chunk <string> | <Buffer> | <Uint8Array> | <any> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов, chunk должен быть строкой, Buffer или Uint8Array. Для потоков в режиме объектов chunk может быть любым значением JavaScript, кроме null.
  • encoding <string> Кодировка, если chunk — строка
  • callback <Function> Необязательный обратный вызов, когда поток завершен

Вызов метода writable.end() сигнализирует о том, что больше данных в поток Writable записываться не будет. Необязательные аргументы chunk и encoding позволяют записать один дополнительный блок данных непосредственно перед закрытием потока. Если предоставлена, необязательная функция callback регистрируется в качестве обработчика события 'finish'.

Вызов метода stream.write() после вызова stream.end() вызовет ошибку.

// write 'hello, ' and then end with 'world!'
const file = fs.createWriteStream('example.txt');
file.write('hello, ');
file.end('world!');
// writing more now is not allowed!
writable.setDefaultEncoding(encoding)
История
Версия Изменения
v6.1.0

Этот метод теперь возвращает ссылку на writable.

v0.11.15

Добавлен в: v0.11.15

  • encoding <string> Новая кодировка по умолчанию
  • Возвращает: <this>

Метод writable.setDefaultEncoding() устанавливает кодировку по умолчанию encoding для потока Writable.

writable.uncork()
Добавлен в: v0.11.2

Метод 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.writableHighWaterMark
Добавлен в: v8.10.0

Возвращает значение highWaterMark , переданное при создании этой Writable.

writable.write(chunk[, encoding][, callback])
История
Версия Изменения
v8.0.0

Аргумент chunk теперь может быть экземпляром Uint8Array.

v6.0.0

Передача null в качестве параметра chunk всегда будет считаться недопустимой теперь, даже в режиме объектов.

v0.9.4

Добавлен в: v0.9.4

  • chunk <string> | <Buffer> | <Uint8Array> | <any> Необязательные данные для записи. Для потоков, не работающих в режиме объектов, chunk должен быть строкой, Buffer или Uint8Array. Для потоков в режиме объектов, chunk может быть любым значением JavaScript, кроме null.
  • encoding <string> Кодировка, если chunk — строка
  • callback <Function> Обратный вызов, когда этот блок данных обработан
  • Возвращает: <boolean> false если поток хочет, чтобы вызывающий код ждал события 'drain', прежде чем продолжить запись дополнительных данных; в противном случае true.

Метод writable.write() записывает данные в поток и вызывает предоставленный callback после того, как данные будут полностью обработаны. Если произошла ошибка, callback может или не может быть вызван с ошибкой в качестве первого аргумента. Для надежного обнаружения ошибок записи добавьте обработчик события 'error'.

Значение возврата — true , если внутренний буфер меньше highWaterMark , настроенного при создании потока после приема chunk. Если возвращается false, дальнейшие попытки записи данных в поток следует приостановить до тех пор, пока не будет отправлено событие 'drain'.

Пока поток не завершает обработку, вызовы write() будут буферизировать chunk, и возвращать false. После того, как все буферизованные блоки будут обработаны (приняты для передачи операционной системой), будет отправлено событие 'drain'. Рекомендуется не отправлять новые блоки до тех пор, пока не будет отправлено событие 'drain' . Вызов write() для потока, который не обрабатывает данные, разрешен, но Node.js будет буферизовать все записанные блоки до достижения максимального использования памяти, после чего произойдёт прерывание. Даже до прерывания, высокое использование памяти повлияет на производительность сборщика мусора и увеличит RSS (который обычно не возвращается системе, даже после того, как память больше не требуется). Поскольку TCP-сокеты могут никогда не завершить обработку данных, если удалённый узел не читает данные, запись в сокет, который не завершает обработку, может привести к удалённо эксплуатируемой уязвимости.

Запись данных, пока поток не обрабатывается, особенно проблематична для потока Transform, потому что потоки Transform приостановлены по умолчанию до подключения или добавления обработчиков событий 'data' или 'readable'.

Если данные для записи могут быть сгенерированы или получены по запросу, рекомендуется инкапсулировать логику в Readable и использовать stream.pipe(). Однако, если вы предпочитаете вызывать write(), можно учесть обратную связь и избежать проблем с памятью, используя событие 'drain':

function write(data, cb) {
  if (!stream.write(data)) {
    stream.once('drain', cb);
  } else {
    process.nextTick(cb);
  }
}

// Wait for cb to be called before doing any other write.
write('hello', () => {
  console.log('write completed, do more writes now');
});

Поток Writable в режиме объектов всегда игнорирует аргумент encoding.

writable.destroy([error])
Добавлен в: v8.0.0
  • Возвращает: <this>

Уничтожить поток и вывести переданную ошибку. После этого вызова поток записи завершён. Реализаторы не должны переопределять этот метод, но вместо этого должны реализовать writable._destroy.

Потоки Readable

Потоки Readable представляют собой абстракцию источника, из которого потребляются данные.

Примеры Readable потоков включают:

  • HTTP-ответы на стороне клиента
  • HTTP-запросы на стороне сервера
  • Потоки чтения fs
  • Потоки zlib
  • Потоки crypto
  • TCP-сокеты
  • Выходной поток child process (stdout) и stderr
  • process.stdin

Все потоки Readable реализуют интерфейс, определенный классом stream.Readable.

Два режима

Потоки Readable эффективно работают в одном из двух режимов: «потоковый» и «приостановленный».

В режиме «потоковый» данные считываются из базовой системы автоматически и предоставляются приложению как можно быстрее с помощью событий через интерфейс EventEmitter.

В режиме «приостановленный» метод stream.read() необходимо вызывать явно для чтения фрагментов данных из потока.

Все потоки Readable начинаются в режиме «приостановленный», но могут быть переключены в режим «потоковый» одним из следующих способов:

  • Добавление обработчика событий 'data'.
  • Вызов метода stream.resume().
  • Вызов метода stream.pipe() для отправки данных в поток Writable.

Readable может вернуться в режим «приостановленный» одним из следующих способов:

  • Если нет пунктов назначения для передачи, вызовом метода stream.pause().
  • Если есть пункты назначения для передачи, удалив все пункты назначения. Несколько пунктов назначения можно удалить, вызвав метод stream.unpipe().

Важно помнить, что поток Readable не будет генерировать данные до тех пор, пока не будет предоставлен механизм для потребления или игнорирования этих данных. Если механизм потребления отключён или удалён, поток Readable попытается остановить генерацию данных.

Примечание: по причинам обратной совместимости удаление обработчиков событий 'data' не автоматически приостановит поток. Кроме того, если существуют подключенные пункты назначения, вызов stream.pause() не гарантирует, что поток останется приостановленным после того, как эти пункты назначения обработают все данные и запросят больше данных.

Примечание: если поток Readable переключается в режим «потоковый», а потребителей для обработки данных нет, эти данные будут потеряны. Это может произойти, например, когда вызывается метод readable.resume() без установленного обработчика для события 'data', или когда обработчик события 'data' удаляется из потока.

Три состояния

«Два режима» работы потока Readable — это упрощённая абстракция для более сложного внутреннего управления состоянием, происходящего внутри реализации потока Readable.

Конкретно, в любой момент времени каждый Readable находится в одном из трёх возможных состояний:

  • readable._readableState.flowing = null
  • readable._readableState.flowing = false
  • readable._readableState.flowing = true

Когда readable._readableState.flowing находится в состоянии null, механизм потребления данных потока не предоставлен, поэтому поток не будет генерировать свои данные. В этом состоянии подключение обработчика для события 'data', вызов метода readable.pipe() или readable.resume() переведут readable._readableState.flowing в состояние true, заставив Readable начать активно генерировать события по мере генерации данных.

Вызов readable.pause(), readable.unpipe() или получение «обратной задержки» приведут к установке состояния readable._readableState.flowing в false, временно приостановив поток событий, но не приостановив генерацию данных. В этом состоянии подключение обработчика для события 'data' не заставит readable._readableState.flowing переключиться в состояние true.

const { PassThrough, Writable } = require('stream');
const pass = new PassThrough();
const writable = new Writable();

pass.pipe(writable);
pass.unpipe(writable);
// flowing 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 'data' being emitted

Пока readable._readableState.flowing находится в состоянии false, данные могут накапливаться во внутренней буферной памяти потока.

Выбор одного метода

API потоков Readable эволюционировал во многих версиях Node.js и предоставляет несколько методов потребления данных потока. В целом, разработчики должны выбрать один метод потребления данных и никогда не использовать несколько методов для потребления данных из одного потока.

Для большинства пользователей рекомендуется использовать метод readable.pipe(), так как он реализован для обеспечения наилучшего способа потребления данных потока. Разработчики, которым требуется более точный контроль над передачей и генерацией данных, могут использовать API EventEmitter и readable.pause()/readable.resume().

Класс: stream.Readable

Добавлен в: v0.9.4
Событие: 'close'
Добавлен в: v0.9.4

Событие 'close' генерируется, когда поток и все его базовые ресурсы (например, дескриптор файла) закрыты. Это событие указывает, что больше событий не будет генерироваться и дальнейших вычислений не будет происходить.

Не все потоки Readable генерируют событие 'close'.

Событие: 'data'
Добавлен в: v0.9.4
  • 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'
Добавлен в: v0.9.4

Событие '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'
Добавлен в: v0.9.4
  • <Ошибка>

Событие 'error' может быть сгенерировано реализацией Readable в любой момент. Как правило, это может произойти, если базовый поток не может сгенерировать данные из-за внутренней ошибки, или когда реализация потока пытается передать неверный фрагмент данных.

Обратный вызов обработчика получит объект Error.

Событие: 'readable'
Добавлен в: v0.9.4

Событие 'readable' генерируется, когда данные доступны для чтения из потока. В некоторых случаях подключение обработчика для события 'readable' вызовет считывание определённого количества данных во внутренний буфер.

const readable = getReadableStreamSomehow();
readable.on('readable', () => {
  // there is some data to read now
});

Событие '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.isPaused()
Добавлен в: v0.11.14
  • Возвращает: <логическое>

Метод 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()
Добавлен в: v0.9.4
  • Возвращает: <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.pipe(destination[, options])
Добавлен в: v0.9.4
  • destination <stream.Writable> Место назначения для записи данных
  • options <Object> Параметры конвейера
    • end <boolean> Завершить запись при завершении чтения. По умолчанию: true.

Метод readable.pipe() присоединяет поток Writable к readable, заставляя его автоматически перейти в режим потоковой передачи и передать все данные присоединённому потоку Writable. Поток данных будет автоматически управляться, чтобы поток назначения Writable не перегружался более быстрым потоком Readable.

Следующий пример передает все данные из readable в файл с именем file.txt:

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 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.readableHighWaterMark
Добавлена в: v8.10.0

Возвращает значение highWaterMark , переданное при создании этого Readable.

readable.read([size])
Добавлена в: v0.9.4
  • size <number> Необязательный аргумент для указания количества данных для чтения.
  • Возвращает: <string> | <Buffer> | <null>

Метод readable.read() извлекает данные из внутреннего буфера и возвращает их. Если доступны данные для чтения, возвращается null. По умолчанию данные возвращаются как объект Buffer, если не указано кодирование с помощью метода readable.setEncoding() или поток не работает в режиме объектов.

Необязательный аргумент size задаёт конкретное количество байтов для чтения. Если size байта недоступны для чтения, будет возвращено null, за исключением случаев, когда поток завершён, в этом случае возвращаются все оставшиеся данные во внутреннем буфере.

Если аргумент size не указан, возвращаются все данные, содержащиеся во внутреннем буфере.

Метод 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.`);
  }
});

Поток Readable в режиме объектов всегда возвращает один элемент при вызове readable.read(size), независимо от значения аргумента size.

Примечание: Если метод readable.read() возвращает фрагмент данных, также будет генерироваться событие 'data'.

Примечание: Вызов stream.read([size]) после генерации события 'end' вернёт null. Ошибка во время выполнения не будет поднята.

readable.resume()
Добавлена в: v0.9.4
  • Возвращает: <this>

Метод readable.resume() вызывает возобновление явно приостановленного потока Readable, возобновляя генерацию событий 'data' и перевод потока в режим потоковой передачи.

Метод readable.resume() можно использовать для полного потребления данных из потока без фактической обработки этих данных, как показано в следующем примере:

getReadableStreamSomehow()
  .resume()
  .on('end', () => {
    console.log('Reached the end, but did not read anything.');
  });
readable.setEncoding(encoding)
Добавлена в: v0.9.4
  • encoding <string> Кодировка для использования.
  • Возвращает: <this>

Метод 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])
Добавлена в: v0.9.4
  • destination <stream.Writable> Необязательный конкретный поток для отключения конвейера

Метод readable.unpipe() отсоединяет поток Writable, ранее присоединённый с помощью метода stream.pipe().

Если destination не указан, то все конвейеры отсоединяются.

Если destination указан, но для него не настроен конвейер, то метод ничего не делает.

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)
История
Версия Изменения
v8.0.0

Аргумент chunk теперь может быть экземпляром Uint8Array.

v0.9.11

Добавлена в: v0.9.11

  • 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() в пользовательском потоке). Следующий вызов stream.push('') после вызова readable.unshift() сбросит состояние чтения должным образом, однако лучше просто избегать вызова readable.unshift() во время чтения.

readable.wrap(stream)
Добавлена в: v0.9.4
  • stream <Stream> Поток чтения "старого стиля"

Версии Node.js до v0.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.destroy([error])
Добавлен в: v8.0.0

Уничтожить поток и вывести 'error'. После этого вызова читаемый поток освободит все внутренние ресурсы. Реализаторы не должны переопределять этот метод, а вместо этого реализовать readable._destroy.

Потоки Duplex и Transform

Класс: stream.Duplex

История
Версия Изменения
v6.8.0

Экземпляры Duplex теперь возвращают true при проверке instanceof stream.Writable.

v0.9.4

Добавлен в: v0.9.4

Потоки Duplex — это потоки, которые реализуют как интерфейс Readable, так и Writable.

Примеры потоков Duplex включают:

  • сокеты TCP
  • потоки zlib
  • потоки crypto

Класс: stream.Transform

Добавлен в: v0.9.4

Потоки Transform — это потоки Duplex, где вывод каким-то образом связан с вводом. Как и все потоки Duplex, потоки Transform реализуют как интерфейс Readable, так и Writable.

Примеры потоков Transform включают:

  • потоки zlib
  • потоки crypto
transform.destroy([error])
Добавлен в: v8.0.0

Уничтожить поток и вывести 'error'. После этого вызова поток transform освободит все внутренние ресурсы. Реализаторы не должны переопределять этот метод, а вместо этого реализовать readable._destroy. По умолчанию реализация _destroy для Transform также выводит 'close'.

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 для потребителей потоков). Это может привести к нежелательным побочным эффектам в прикладном коде, использующем поток.

Упрощенное создание

Добавлен в: v1.2.0

Во многих простых случаях можно создать поток, не используя наследование. Это можно сделать, напрямую создав экземпляры объектов stream.Writable, stream.Readable, stream.Duplex или stream.Transform и передав соответствующие методы в качестве параметров конструктора.

Например:

const { Writable } = require('stream');

const myWritable = new Writable({
  write(chunk, encoding, callback) {
    // ...
  }
});

Реализация потока Writable

Класс stream.Writable расширяется для реализации потока Writable.

Пользовательские потоки Writable обязательно должны вызывать конструктор new stream.Writable([options]) и реализовать метод writable._write().

Конструктор: new stream.Writable([options])

  • options <Объект>
    • highWaterMark <число> Уровень буфера, когда stream.write() начинает возвращать false. По умолчанию: 16384 (16 кб) или 16 для потоков objectMode.
    • decodeStrings <логическое значение> Нужно ли декодировать строки в буферы перед передачей их в stream._write(). По умолчанию: true.
    • objectMode <логическое значение> Является ли stream.write(anyObj) валидной операцией. При установке становится возможной запись значений JavaScript, отличных от строки, Buffer или Uint8Array (если это поддерживается реализацией потока). По умолчанию: false.
    • write <Функция> Реализация метода stream._write().
    • writev <Функция> Реализация метода stream._writev().
    • destroy <Функция> Реализация метода stream._destroy().
    • final <Функция> Реализация метода stream._final().

Например:

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 <Буфер> | <строка> | <любое значение> Буфер для записи. Всегда будет буфером, если не установлено опцию decodeStrings на false, или если поток работает в режиме объектов.
  • 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 установлено в параметрах конструктора, тогда chunk может быть строкой, а не буфером, и encoding будет указывать кодировку символов строки. Это для поддержки реализаций, которые имеют оптимизированную обработку определённых кодировок строк. Если свойство decodeStrings явно установлено в false, аргумент encoding можно безопасно проигнорировать, и chunk останется тем же объектом, что и переданный в .write().

Метод writable._write() имеет префикс подчеркивания, потому что он внутренний для класса, который его определяет, и его никогда не следует вызывать напрямую программами пользователя.

END_OF_DOCUMENT_MARKER

writable._writev(chunks, callback)

  • chunks <Массив> Набор фрагментов для записи. Каждый фрагмент имеет следующий формат: { chunk: ..., encoding: ... }.
  • callback <Функция> Функция обратного вызова (возможно, с аргументом ошибки), которая вызывается после завершения обработки переданных фрагментов.

Примечание: к этой функции ПОЛНОСТЬЮ ЗАПРЕЩЕНО обращаться из прикладного кода напрямую. Она должна быть реализована дочерними классами и вызываться только методами внутреннего класса Writable.

Метод writable._writev() может быть реализован дополнительно к методу writable._write() в реализациях потоков, способных обрабатывать сразу несколько фрагментов данных. Если реализован, этот метод будет вызван со всеми фрагментами данных, в настоящее время буферизованными в очереди записи.

Метод writable._writev() имеет префикс подчеркивания, потому что он внутренний для класса, который его определяет, и никогда не должен вызываться непосредственно программами пользователя.

writable._destroy(err, callback)

Добавлена в: v8.0.0
  • err <Ошибка> Возможная ошибка.
  • callback <Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.

Метод _destroy() вызывается методом writable.destroy(). Он может быть переопределён дочерними классами, но не должен вызываться напрямую.

writable._final(callback)

Добавлена в: v8.0.0
  • callback <Функция> Вызовите эту функцию (возможно, с аргументом ошибки), когда завершена запись всех оставшихся данных.

Метод _final() не должен вызываться напрямую. Он может быть реализован дочерними классами и, если это сделано, будет вызван только методами внутреннего класса Writable.

Эта необязательная функция будет вызвана перед закрытием потока, откладывая событие finish до вызова callback. Это полезно для закрытия ресурсов или записи буферизованных данных перед окончанием потока.

Ошибки при записи

Рекомендуется, чтобы ошибки, возникающие во время обработки методов writable._write() и writable._writev(), сообщались путём вызова обратного вызова и передачи ошибки в качестве первого аргумента. Это вызовет событие 'error' , которое будет излучаться Writable. Бросание ошибки внутри 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 не обладает какой-либо реальной полезностью, пример демонстрирует каждый из необходимых элементов экземпляра пользовательского потока Writable:

const { Writable } = require('stream');

class MyWritable extends Writable {
  constructor(options) {
    super(options);
    // ...
  }

  _write(chunk, encoding, callback) {
    if (chunk.toString().indexOf('a') >= 0) {
      callback(new Error('chunk is invalid'));
    } else {
      callback();
    }
  }
}

Декодирование буферов в потоке Writable

Декодирование буферов — это распространённая задача, например, при использовании преобразователей, вход которых является строкой. Это нетривиальный процесс при использовании кодировки символов с несколькими байтами, такой как UTF-8. Следующий пример показывает, как декодировать строки с несколькими байтами, используя StringDecoder и Writable.

const { Writable } = require('stream');
const { StringDecoder } = require('string_decoder');

class StringWritable extends Writable {
  constructor(options) {
    super(options);
    const state = this._writableState;
    this._decoder = new StringDecoder(state.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: €

Реализация потока Readable

Класс 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) возвращает единственное значение вместо буфера размером n. По умолчанию: false.
    • read <Функция> Реализация метода stream._read().
    • destroy <Функция> Реализация метода stream._destroy().

Например:

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().

Аргумент size является рекомендательным. Для реализаций, где «чтение» — это однократная операция, возвращающая данные, можно использовать аргумент size для определения того, сколько данных извлечь. Другие реализации могут игнорировать этот аргумент и просто предоставлять данные по мере их появления. Нет необходимости «ждать», пока size байтов станет доступно, прежде чем вызывать stream.push(chunk).

Метод readable._read() имеет префикс подчеркивания, потому что он внутренний для класса, который его определяет, и никогда не должен вызываться непосредственно программами пользователя.

readable._destroy(err, callback)

Добавлена в: v8.0.0
  • err <Ошибка> Возможная ошибка.
  • callback <Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.

Метод _destroy() вызывается методом readable.destroy(). Он может быть переопределен дочерними классами, но не должен вызываться напрямую.

readable.push(chunk[, encoding])

История
Версия Изменения
v8.0.0

Аргумент chunk теперь может быть экземпляром Uint8Array.

  • chunk <Буфер> | <Uint8Array> | <Строка> | <null> | <любое> Фрагмент данных для помещения в очередь чтения. Для потоков, не работающих в режиме объектов, chunk должен быть строкой, Buffer или Uint8Array. Для потоков в режиме объектов, chunk может быть любым значением JavaScript.
  • encoding <Строка> Кодировка фрагментов строк. Должна быть допустимой кодировкой буфера, такой как 'utf8' или 'ascii'
  • Возвращает: <Булево> true если дополнительные фрагменты данных могут быть помещены; false в противном случае.

Когда chunk является Buffer, Uint8Array или string, данные chunk будут добавлены в внутреннюю очередь для потребления пользователями потока. Передача chunk в качестве null сигнализирует об окончании потока (EOF), после чего больше данных записывать нельзя.

Когда Readable работает в режиме приостановки, добавленные данные с помощью readable.push() могут быть прочитаны вызовом метода readable.read(), когда срабатывает событие 'readable'.

Когда Readable работает в режиме потока, данные, добавленные с помощью readable.push(), будут переданы с помощью события 'data'.

Метод readable.push() разработан для максимальной гибкости. Например, при обёртке источника более низкого уровня, предоставляющего какой-либо механизм приостановки/возобновления и обратный вызов для данных, низкоуровневый источник можно обернуть в пользовательском экземпляре Readable, как показано в следующем примере:

// source is an object with readStop() and readStart() methods,
// and an `ondata` member that gets called when it has data, and
// an `onend` member that gets called when the data is over.

class SourceWrapper extends Readable {
  constructor(options) {
    super(options);

    this._source = getLowlevelSourceObject();

    // Every time there's data, push it into the internal buffer.
    this._source.ondata = (chunk) => {
      // if push() returns false, then stop reading from source
      if (!this.push(chunk))
        this._source.readStop();
    };

    // When the source ends, push the EOF-signaling `null` chunk
    this._source.onend = () => {
      this.push(null);
    };
  }
  // _read will be called when the stream wants to pull more data in
  // the advisory size argument is ignored in this case.
  _read(size) {
    this._source.readStart();
  }
}

Примечание: Метод readable.push() предназначен только для реализации Readable и только изнутри метода readable._read().

Для потоков, не работающих в режиме объектов, если параметр chunk метода readable.push() равен undefined, он будет обработан как пустая строка или буфер. Подробнее см. readable.push('').

Ошибки при чтении

Рекомендуется, чтобы ошибки, возникающие во время обработки метода readable._read(), передавались с помощью события '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 = '' + 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)

История
Версия Изменения
v8.4.0

Теперь поддерживаются параметры readableHighWaterMark и writableHighWaterMark.

  • 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 потока является то, что стороны чтения и записи работают независимо друг от друга, несмотря на совместное существование в одном экземпляре объекта.

Потоки Duplex в режиме объектов

Для Duplex потоков, objectMode может быть установлено исключительно для стороны чтения или записи, используя параметры 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 потоков необходимо соблюдать осторожность, так как данные, записанные в поток, могут привести к приостановке стороны записи потока, если выходные данные на стороне чтения не потребляются.

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.

END_OF_DOCUMENT_MARKER

В некоторых случаях операция преобразования может потребоваться для вывода дополнительных данных в конце потока. Например, поток сжатия zlib будет хранить объем внутренней информации, используемой для оптимального сжатия выходных данных. Однако, когда поток завершается, эти дополнительные данные необходимо сбросить, чтобы сжатые данные были полными.

Реализации пользовательского преобразования преобразования могут реализовать метод transform._flush(). Он будет вызван, когда больше нет данных для чтения, но до того, как будет испущен событие 'end', сигнализирующее о завершении потока Readable.

В реализации transform._flush(), метод readable.push() может быть вызван ноль или более раз, по мере необходимости. Функция callback должна быть вызвана, когда операция сброса завершена.

Метод transform._flush() имеет префикс подчеркивания, потому что он является внутренним для класса, который его определяет, и его никогда не следует вызывать напрямую программами пользователя.

transform._transform(chunk, encoding, callback)

  • chunk <Буфер> | <строка> | <любой> Обрабатываемый фрагмент. Всегда будет буфером, если опция decodeStrings не была установлена в false, или поток работает в режиме объектов.
  • encoding <строка> Если фрагмент является строкой, то это тип кодирования. Если фрагмент — буфер, то это специальное значение — 'buffer', игнорируйте его в этом случае.
  • callback <Функция> Функция обратного вызова (необязательно с аргументом ошибки и данными), которая должна быть вызвана после обработки предоставленного chunk.

Примечание: эту функцию НЕЛЬЗЯ вызывать кодом приложения напрямую. Она должна быть реализована производными классами и вызываться только внутренними методами класса Readable.

Все реализации потоков преобразования должны предоставить метод _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 может быть полезен как строительный блок для новых типов потоков.

Дополнительные заметки

Совместимость со старыми версиями Node.js

В версиях Node.js до v0.10 интерфейс потока Readable был проще, но также менее мощным и менее полезным.

  • Вместо ожидания вызовов метода stream.read(), события 'data' начинали испускаться немедленно. Приложениям, которым потребовалось бы выполнить некоторую работу, чтобы решить, как обработать данные, приходилось сохранять прочитанные данные в буферы, чтобы данные не потерялись.
  • Метод stream.pause() был рекомендательным, а не гарантированным. Это означало, что все равно нужно было быть готовым к получению событий 'data' даже когда поток был в приостановленном состоянии.

В Node.js v0.10 был добавлен класс Readable. Для обеспечения обратной совместимости со старыми программами Node.js потоки Readable переключаются в «режим потока» при добавлении обработчика события 'data' или при вызове метода stream.resume(). В результате даже без использования нового метода stream.read() и события 'readable' больше не нужно беспокоиться о потере фрагментов 'data'.

Хотя большинство приложений будут продолжать работать нормально, это создает исключительную ситуацию в следующих условиях:

  • Не добавляется обработчик события 'data'.
  • Метод stream.resume() никогда не вызывается.
  • Поток не направлен ни на какой записываемый конечный пункт.

Например, рассмотрим следующий код:

// WARNING!  BROKEN!
net.createServer((socket) => {

  // we add an 'end' method, 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 до v0.10 данные входящего сообщения просто отбрасывались. Однако в Node.js v0.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 в режим потока, потоки в стиле до v0.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-v8.x/docs/api/stream.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API