Spec-Zone.ru › Node.js 14 LTS

Поток[src]

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

Исходный код: 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: Duplex потоки, которые могут изменять или преобразовывать данные при записи и чтении (например, zlib.createDeflate()).

Кроме того, этот модуль включает вспомогательные функции stream.pipeline(), stream.finished() и stream.Readable.from().

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

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

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

Буферизация

Как потоки Writable, так и Readable будут хранить данные во внутреннем буфере.

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

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

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

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

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

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

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

Механизм внутреннего буферирования является внутренней реализацией и может быть изменен в любое время. Однако для некоторых расширенных реализаций внутренние буферы можно получить с помощью writable.writableBuffer или readable.readableBuffer. Использование этих недокументированных свойств не рекомендуется.

API для потребителей потоков

Практически все приложения Node.js, независимо от сложности, используют потоки каким-либо образом. Ниже приведен пример использования потоков в приложении Node.js, реализующем HTTP-сервер:

const http = require('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
Добавлен в: v0.9.4
Событие: 'close'
История
Версия Изменения
v10.0.0

Добавлен параметр emitClose для указания, будет ли событие 'close' генерироваться при уничтожении.

v0.9.4

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

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

Writable поток всегда генерирует событие 'close', если он создан с параметром emitClose.

Событие: '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', если только параметр autoDestroy не был установлен в false при создании потока.

После 'error', никаких других событий, кроме 'close', не должно быть сгенерировано (включая события 'error').

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

Событие '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'
Добавлен в: v0.9.4
  • 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'
Добавлен в: v0.9.4
  • src Поток-источник <stream.Readable>, который отсоединил этот поток записи

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

Это также генерируется в случае, если этот поток записи 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()
Добавлен в: v0.11.2

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

Действует как операция без действия для потока, который уже был уничтожен.

v8.0.0

Добавлен в: v8.0.0

  • error <Ошибка> Необязательно, ошибка для генерации события 'error'.
  • Возвращает: <this>

Уничтожить поток. Необязательно сгенерировать событие 'error', и сгенерировать событие 'close' (если emitClose не установлен в false). После этого вызова поток записи закрыт, и последующие вызовы write() или end() приведут к ошибке ERR_STREAM_DESTROYED. Это деструктивный и мгновенный способ уничтожения потока. Предыдущие вызовы write() могут не быть завершены и могут вызвать ошибку ERR_STREAM_DESTROYED. Используйте end() вместо destroy, если данные должны быть переданы до закрытия, или подождите события 'drain' перед уничтожением потока.

После вызова destroy() любые дальнейшие вызовы будут игнорироваться, и больше никаких ошибок, кроме тех, которые могут быть сгенерированы _destroy(), не будут генерироваться в виде 'error'.

Разработчики не должны переопределять этот метод, а вместо этого реализовать writable._destroy().

writable.destroyed
Добавлен в: v8.0.0
  • <логическое значение>

Является true после вызова writable.destroy().

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

Вызывается callback, если испущены 'finish' или 'error'.

v10.0.0

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

v8.0.0

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

v0.9.4

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

  • chunk <строка> | <Buffer> | <Uint8Array> | <любой> Дополнительные данные для записи. Для потоков, не работающих в объектном режиме, chunk должно быть строкой, Buffer или Uint8Array. Для потоков в объектном режиме, chunk может быть любым значением JavaScript, кроме null.
  • encoding <строка> Кодировка, если chunk является строкой
  • callback <Функция> Необязательный обработчик для завершения потока или ошибок
  • Возвращает: <this>

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

Вызов метода 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)
История
Версия Изменения
v6.1.0

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

v0.11.15

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

  • encoding <строка> Новая кодировка по умолчанию
  • Возвращает: <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.writable
Добавлен в: v11.4.0
  • <логическое>

Истинно, если безопасно вызвать writable.write(), что означает, что поток не был уничтожен, не получил ошибки или не завершен.

writable.writableEnded
Добавлен в: v12.9.0
  • <логическое>

Истинно после вызова writable.end(). Это свойство не указывает, были ли данные сброшены, для этого используйте writable.writableFinished.

writable.writableCorked
Добавлен в: v13.2.0, v12.16.0
  • <целое число>

Количество раз, когда необходимо вызвать writable.uncork() для полного разблокирования потока.

writable.writableFinished
Добавлен в: v12.6.0
  • <логическое>

Устанавливается в true непосредственно перед тем, как испустить событие 'finish'.

writable.writableHighWaterMark
Добавлен в: v9.3.0
  • <число>

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

writable.writableLength
Добавлен в: v9.4.0
  • <число>

Это свойство содержит количество байтов (или объектов) в очереди, готовых к записи. Значение предоставляет данные для интроспекции о состоянии highWaterMark.

writable.writableNeedDrain
Добавлен в: v14.17.0
  • <логическое>

Истинно, если буфер потока был заполнен и поток испустит 'drain'.

writable.writableObjectMode
Добавлен в: v12.3.0
  • <логическое>

Геттер для свойства objectMode заданного потока Writable.

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

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

v6.0.0

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

v0.9.4

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

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

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

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

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

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

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

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

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

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

Потоки чтения

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

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

  • Ответы HTTP на стороне клиента
  • Запросы HTTP на стороне сервера
  • Потоки чтения fs
  • Потоки zlib
  • Потоки crypto
  • Сокеты TCP
  • Вывод и стандартный вывод процесса-потомка
  • process.stdin

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

Два режима чтения

Потоки Readable фактически работают в одном из двух режимов: потоковый и приостановленный. Эти режимы отделены от режима обработки объектов. Поток Readable может быть в режиме обработки объектов или не быть, независимо от того, находится ли он в потоковом режиме или режиме приостановки.

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

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

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

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

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

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

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

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

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

Добавление обработчика события 'readable' автоматически останавливает поток, и данные необходимо потреблять через readable.read(). Если обработчик события 'readable' удалён, поток начнёт потоковый режим, если есть обработчик события 'data'.

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

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

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

  • readable.readableFlowing === null
  • readable.readableFlowing === false
  • readable.readableFlowing === true

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

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

const { PassThrough, Writable } = require('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
Добавлен в: v0.9.4
Событие: 'close'
История
Версия Изменения
v10.0.0

Добавлен параметр emitClose для указания, излучается ли 'close' при уничтожении.

v0.9.4

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

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

Поток Readable всегда будет излучать событие 'close', если он был создан с параметром emitClose.

Событие: '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.

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

Событие 'pause' генерируется, когда вызывается stream.pause(), и readableFlowing не false.

Событие: 'readable'
История
Версия Изменения
v10.0.0

Событие 'readable' всегда генерируется в следующем цикле после вызова .push().

v10.0.0

Использование 'readable' требует вызова .read().

v0.9.4

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

Событие '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. Если при удалении 'readable' существуют слушатели 'data', поток начнёт работать, т.е. события 'data' будут генерироваться без вызова .resume().

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

Событие 'resume' генерируется, когда вызывается stream.resume(), и readableFlowing не true.

readable.destroy([error])
История
Версия Изменения
v14.0.0

Действует как команда без операции на потоке, который уже был уничтожен.

v8.0.0

Добавлена в: v8.0.0

  • error <Ошибка> Ошибка, которая будет передана в качестве полезной нагрузки в событии 'error'
  • Возвращает: <this>

Уничтожить поток. При необходимости сгенерировать событие 'error' и событие 'close' (если emitClose не установлено в false). После этого вызова потоковый поток readable освободит все внутренние ресурсы, и последующие вызовы push() будут проигнорированы.

После вызова destroy() все дальнейшие вызовы будут считаться операциями без действия, и дальнейшие ошибки, кроме _destroy(), не могут быть сгенерированы как 'error'.

Реализаторы не должны переопределять этот метод, а вместо этого должны реализовать readable._destroy().

readable.destroyed
Добавлена в: v8.0.0
  • <boolean>

Является ли поток true после вызова readable.destroy().

readable.isPaused()
Добавлена в: v0.11.14
  • Возвращает: <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()
Добавлена в: 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.pause() не имеет эффекта, если есть слушатель события 'readable'.

readable.pipe(destination[, options])
Добавлена в: v0.9.4
  • 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])
Добавлена в: v0.9.4
  • 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
Добавлен в: v11.4.0
  • <boolean>

Имеет значение true , если безопасно вызвать readable.read(), что означает, что поток не был уничтожен и не было отправлено событие 'error' или 'end'.

readable.readableEncoding
Добавлен в: v12.7.0
  • <null> | <string>

Получение значения свойства encoding заданного потока Readable. Свойство encoding можно установить с помощью метода readable.setEncoding().

readable.readableEnded
Добавлен в: v12.9.0
  • <boolean>

Принимает значение true при отправке события 'end'.

readable.readableFlowing
Добавлен в: v9.4.0
  • <boolean>

Это свойство отражает текущее состояние потока Readable , как описано в разделе Три состояния.

readable.readableHighWaterMark
Добавлен в: v9.3.0
  • <number>

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

readable.readableLength
Добавлен в: v9.4.0
  • <number>

Это свойство содержит количество байтов (или объектов) в очереди, готовых для чтения. Значение предоставляет данные для интроспекции состояния потока highWaterMark.

readable.readableObjectMode
Добавлен в: v12.3.0
  • <boolean>

Получение значения свойства objectMode заданного потока Readable.

readable.resume()
История
Версия Изменения
v10.0.0

Метод resume() не имеет эффекта, если существует обработчик события 'readable'.

v0.9.4

Добавлен в: 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.resume() не имеет эффекта, если существует обработчик события 'readable'.

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> Необязательный конкретный поток для разрыва соединения
  • Возвращает: <this>

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

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

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

const fs = require('fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt',
// but only for the first second.
readable.pipe(writable);
setTimeout(() => {
  console.log('Stop writing to file.txt.');
  readable.unpipe(writable);
  console.log('Manually close the file stream.');
  writable.end();
}, 1000);
readable.unshift(chunk[, encoding])
История
Версия Изменения
v8.0.0

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

v0.9.11

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

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

Передача chunk в качестве null сигнализирует об окончании потока (EOF) и ведет себя так же, как readable.push(null), после чего больше данных записать нельзя. Сигнал EOF помещается в конец буфера, и любые буферизованные данные всё ещё будут сброшены.

Метод readable.unshift() возвращает фрагмент данных в внутренний буфер. Это полезно в определенных ситуациях, когда поток потребляется кодом, которому необходимо "отменить потребление" некоторого количества данных, которые он оптимистично извлек из источника, чтобы данные могли быть переданы какой-либо другой стороне.

Метод stream.unshift(chunk) не может быть вызван после отправки события 'end', в противном случае будет возбуждена ошибка выполнения.

Разработчики, использующие stream.unshift(), часто должны рассмотреть возможность переключения на использование потока Transform вместо него. Дополнительную информацию см. в разделе API для разработчиков потоков.

// Pull off a header delimited by \n\n.
// Use unshift() if we get too much.
// Call the callback with (error, header, stream).
const { StringDecoder } = require('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)
Добавлен в: v0.9.4
  • stream <Поток> Поток чтения "старого стиля"
  • Возвращает: <this>

До версии 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]()
История
Версия Изменения
v11.14.0

Поддержка Symbol.asyncIterator больше не является экспериментальной.

v10.0.0

Добавлен в: v10.0.0

  • Возвращает: <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
История
Версия Изменения
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])
История
Версия Изменения
v14.0.0

Действует как no-op для потока, который уже был уничтожен.

v8.0.0

Добавлен в: v8.0.0

  • error <Ошибка>
  • Возвращает: <this>

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

После того, как destroy() был вызван, любые дальнейшие вызовы будут бесполезны, и больше не будет генерироваться никаких ошибок, кроме ошибок от _destroy(), как 'error'.

stream.finished(stream[, options], callback)

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

finished(stream, cb) будет ждать события 'close' перед вызовом обратного вызова. Реализация пытается обнаружить устаревшие потоки и применять это поведение только к потокам, которые, как ожидается, генерируют 'close'.

v14.0.0

Генерация 'close' до 'end' в потоке Readable вызовет ошибку ERR_STREAM_PREMATURE_CLOSE.

v14.0.0

Обратный вызов будет вызван для потоков, которые уже завершены до вызова finished(stream, cb).

v10.0.0

Добавлен в: v10.0.0

  • stream <Поток> Поток для чтения и/или записи.
  • options <Объект>
    • error <boolean> Если установлено в false, то вызов emit('error', err) не считается завершенным. По умолчанию: true.
    • readable <boolean> При установке в false, обратный вызов будет вызван при завершении потока, даже если поток все еще может быть читаемым. По умолчанию: true.
    • writable <boolean> При установке в 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 также поддается промификатции;

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(source[, ...transforms], destination, callback)

stream.pipeline(streams, callback)

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

pipeline(..., cb) будет ждать события 'close' перед вызовом обратного вызова. Реализация пытается обнаружить устаревшие потоки и применять это поведение только к потокам, которые, как ожидается, генерируют 'close'.

v13.10.0

Добавлена поддержка асинхронных генераторов.

v10.0.0

Добавлен в: v10.0.0

  • streams <Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]>
  • source <Stream> | <Iterable> | <AsyncIterable> | <Function>
    • Возвращает: <Iterable> | <AsyncIterable>
  • ...transforms <Stream> | <Function>
    • source <AsyncIterable>
    • Возвращает: <AsyncIterable>
  • destination <Stream> | <Function>
    • source <AsyncIterable>
    • Возвращает: <AsyncIterable> | <Promise>
  • callback <Function> Вызывается, когда конвейер полностью завершен.
    • err <Error>
    • val Возвращаемое значение Promise от destination.
  • Возвращает: <Stream>

Метод модуля для перенаправления данных между потоками и генераторами, передающий ошибки и должным образом очищающий ресурсы, а также предоставляющий обратный вызов по завершении конвейера.

const { pipeline } = require('stream');
const fs = require('fs');
const zlib = require('zlib');

// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.

// A pipeline to gzip a potentially huge tar file efficiently:

pipeline(
  fs.createReadStream('archive.tar'),
  zlib.createGzip(),
  fs.createWriteStream('archive.tar.gz'),
  (err) => {
    if (err) {
      console.error('Pipeline failed.', err);
    } else {
      console.log('Pipeline succeeded.');
    }
  }
);

API pipeline также поддерживает промисификацию:

const pipeline = util.promisify(stream.pipeline);

async function run() {
  await pipeline(
    fs.createReadStream('archive.tar'),
    zlib.createGzip(),
    fs.createWriteStream('archive.tar.gz')
  );
  console.log('Pipeline succeeded.');
}

run().catch(console.error);

API pipeline также поддерживает асинхронные генераторы:

const pipeline = util.promisify(stream.pipeline);
const fs = require('fs');

async function run() {
  await pipeline(
    fs.createReadStream('lowercase.txt'),
    async function* (source) {
      source.setEncoding('utf8');  // Work with strings rather than `Buffer`s.
      for await (const chunk of source) {
        yield chunk.toUpperCase();
      }
    },
    fs.createWriteStream('uppercase.txt')
  );
  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])

Добавлен в: v12.3.0, v10.17.0
  • iterable <Iterable> Объект, реализующий протокол итерации Symbol.asyncIterator или Symbol.iterator. Если передан null, генерируется событие 'error'.
  • options <Object> Опции, предоставляемые для 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 и 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(). Это может нарушить текущие и будущие инварианты потока, что приведет к проблемам с поведением и/или совместимостью с другими потоками, утилитами для потоков и ожиданиями пользователей.

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

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

Во многих простых случаях можно создать поток, не прибегая к наследованию. Это можно сделать, напрямую создав экземпляры объектов 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])
История изменений
Версия Изменения
v14.0.0

Изменено значение параметра autoDestroy по умолчанию на true.

v11.2.0, v10.16.0

Добавлен параметр autoDestroy для автоматического destroy() потока при генерации события 'finish' или ошибки.

v10.0.0

Добавлен параметр emitClose для указания, нужно ли генерировать событие 'close' при уничтожении.

  • 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() на потоке после завершения. По умолчанию: true.
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)
История изменений
Версия Изменения
v12.11.0

_write() необязательно при наличии _writev().

  • chunk <Буфер> | <строка> | <любой тип> Данные для записи, преобразованные из string , переданного в stream.write(). Если параметр decodeStrings потока равен false, или поток работает в режиме объектов, данные не преобразуются и будут такими, как были переданы в stream.write().
  • encoding <строка> Если данные – строка, то encoding – кодировка символов этой строки. Если данные являются Buffer, или если поток работает в режиме объектов, encoding может быть проигнорировано.
  • callback <Функция> Вызовите эту функцию (при необходимости с аргументом ошибки) при завершении обработки предоставленных данных.

Все реализации потоков Writable должны предоставлять метод writable._write() и/или writable._writev() для отправки данных на подлежащий ресурс.

Transform потоки предоставляют собственную реализацию writable._write().

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

Функция callback должна вызываться синхронно внутри writable._write() или асинхронно (т.е. в другом цикле) для сигнализации об успешном завершении записи или о возникновении ошибки. Первым аргументом, передаваемым в callback, должен быть объект Error, если вызов завершился ошибкой, или null, если запись прошла успешно.

Все вызовы writable.write(), которые происходят между вызовом writable._write() и вызовом callback, приведут к буферизации записанных данных. Когда вызывается callback, поток может генерировать событие 'drain'. Если реализация потока способна обрабатывать несколько блоков данных одновременно, должен быть реализован метод writable._writev().

Если свойство decodeStrings явно установлено в false в опциях конструктора, то chunk останется тем же объектом, который передается в .write(), и может быть строкой, а не объектом Buffer. Это поддерживает реализации с оптимизированной обработкой определенных кодировок строк. В этом случае аргумент encoding укажет кодировку символов строки. В противном случае аргумент encoding можно безопасно игнорировать.

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

writable._writev(chunks, callback)
  • chunks <Массив объектов> Данные для записи. Значение представляет собой массив объектов <Объект>, каждый из которых представляет собой отдельный блок данных для записи. Свойства этих объектов:
    • chunk <Буфер> | <строка> Экземпляр буфера или строка, содержащая данные для записи. Буфер chunk будет строкой, если Writable был создан с опцией decodeStrings установленной в значение false, а в write() была передана строка.
    • encoding <строка> Кодировка символов chunk. Если chunk является Buffer, то кодировка encoding будет 'buffer'.
  • callback <Функция> Функция обратного вызова (необязательно с аргументом ошибки), которая вызывается после завершения обработки предоставленных блоков.

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

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

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

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

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

writable._final(callback)
Добавлена в: v8.0.0
  • 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:

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

Изменение значения по умолчанию опции autoDestroy на true.

v11.2.0, v10.16.0

Добавление опции autoDestroy для автоматического destroy() потока при генерации события 'end' или ошибках.

  • 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() по завершении. По умолчанию: true.
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)
Добавлена в: v0.9.4
  • 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)
Добавлен в: 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 <строка> Кодировка строковых чанков. Должна быть допустимой кодировкой Buffer, например 'utf8' или 'ascii'.
  • Возвращает: <логическое значение> true если можно продолжить передачу дополнительных чанков данных; false в противном случае.

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

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

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

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

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

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

    this._source = getLowLevelSourceObject();

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

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

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

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

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

Потоки дуплексного режима с объектами

Для потоков Duplex режим объектов можно установить исключительно для стороны Readable или Writable соответственно с помощью опций readableObjectMode и writableObjectMode.

Например, в следующем примере создается новый поток Transform (который является типом потока Duplex), у которого сторона режима объектов Writable принимает числа JavaScript, которые преобразуются в шестнадцатеричные строки на стороне Readable.

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

// All Transform streams are also Duplex Streams.
const myTransform = new Transform({
  writableObjectMode: true,

  transform(chunk, encoding, callback) {
    // Coerce the chunk to a number if necessary.
    chunk |= 0;

    // Transform the chunk into something else.
    const data = chunk.toString(16);

    // Push the data onto the readable queue.
    callback(null, '0'.repeat(data.length % 2) + data);
  }
});

myTransform.setEncoding('ascii');
myTransform.on('data', (chunk) => console.log(chunk));

myTransform.write(1);
// Prints: 01
myTransform.write(10);
// Prints: 0a
myTransform.write(100);
// Prints: 64

Реализация потока преобразования

Поток Transform — это поток Duplex, где выходные данные вычисляются каким-либо образом из входных данных. Примерами являются потоки zlib или crypto, которые сжимают, шифруют или дешифруют данные.

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

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

Класс stream.Transform прототипически наследуется от stream.Duplex и реализует свои собственные версии методов writable._write() и readable._read(). Реализации пользовательских Transform обязаны реализовывать метод transform._transform() и могут также реализовывать метод transform._flush().

При использовании потоков Transform необходимо учитывать, что данные, записанные в поток, могут привести к приостановке стороны Writable потока, если выходные данные на стороне Readable не обрабатываются.

new stream.Transform([options])
  • options <Объект> Передается как в конструкторы Writable и Readable. Также имеет следующие поля:
    • transform <Функция> Реализация метода stream._transform().
    • flush <Функция> Реализация метода stream._flush().
const { Transform } = require('stream');

class MyTransform extends Transform {
  constructor(options) {
    super(options);
    // ...
  }
}

Или, при использовании конструкторов в стиле до ES6:

const { Transform } = require('stream');
const util = require('util');

function MyTransform(options) {
  if (!(this instanceof MyTransform))
    return new MyTransform(options);
  Transform.call(this, options);
}
util.inherits(MyTransform, Transform);

Или, используя упрощенный подход к конструктору:

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

const myTransform = new Transform({
  transform(chunk, encoding, callback) {
    // ...
  }
});
Событие: 'end'

Событие 'end' исходит из класса stream.Readable. Событие 'end' генерируется после того, как все данные будут выведены, что происходит после вызова обратного вызова в transform._flush(). В случае ошибки, 'end' не должно генерироваться.

Событие: 'finish'

Событие 'finish' исходит из класса stream.Writable. Событие 'finish' генерируется после вызова stream.end() и обработки всех фрагментов методом stream._transform(). В случае ошибки, 'finish' не должно генерироваться.

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

Эта функция НЕ ДОЛЖНА вызываться кодом приложения напрямую. Она должна реализовываться дочерними классами и вызываться только внутренними методами класса Readable.

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

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

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

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

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

Эта функция НЕ ДОЛЖНА вызываться кодом приложения напрямую. Она должна реализовываться дочерними классами и вызываться только внутренними методами класса Readable.

Все реализации потоков Transform должны предоставлять метод _transform() для приема входных данных и создания выходных данных. Реализация transform._transform() обрабатывает записываемые байты, вычисляет вывод, а затем передает этот вывод в читаемую часть с помощью метода transform.push().

Метод transform.push() может быть вызван ноль или более раз для генерации вывода из одного входного фрагмента в зависимости от того, сколько вывода нужно сгенерировать в результате этого фрагмента.

Возможен случай, когда из каких-либо входных данных не будет сгенерировано никакого вывода.

Функция callback должна вызываться только после полного потребления текущего фрагмента. Первый аргумент, переданный в callback, должен быть объектом Error, если при обработке входных данных произошла ошибка, или null в противном случае. Если во второй аргумент функции callback передано значение, оно будет передано методу transform.push(). Другими словами, следующие варианты эквивалентны:

transform.prototype._transform = function(data, encoding, callback) {
  this.push(data);
  callback();
};

transform.prototype._transform = function(data, encoding, callback) {
  callback(null, data);
};

Метод 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);
});
Перенаправление в потоки записи из асинхронных итераторов

При записи в поток записи из асинхронного итератора необходимо обеспечить правильную обработку обратного давления и ошибок. stream.pipeline() абстрагирует обработку обратного давления и ошибок, связанных с обратным давлением.

const { pipeline } = require('stream');
const util = require('util');
const fs = require('fs');

const writable = fs.createWriteStream('./file');

// Callback Pattern
pipeline(iterator, writable, (err, value) => {
  if (err) {
    console.error(err);
  } else {
    console.log(value, 'value returned');
  }
});

// Promise Pattern
const pipelinePromise = util.promisify(pipeline);
pipelinePromise(iterator, writable)
  .then((value) => {
    console.log(value, 'value returned');
  })
  .catch(console.error);

Совместимость со старыми версиями 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('') не рекомендуется.

Передача строки длиной 0 байт, 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-v14.x/docs/api/stream.html

Spec-Zone.ru

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