Spec-Zone.ru › Node.js 6 LTS

Поток

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

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

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

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

Модуль stream можно получить с помощью:

const stream = require('stream');

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

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

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

Типы потоков

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

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

Режим работы с объектами

Все потоки, созданные API Node.js, работают исключительно с строками и Buffer объектами. Однако, реализации потоков могут работать с другими типами 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' генерируется, когда для потока Readable вызывается метод stream.unpipe(), удаляя этот потоковый выход Writable из набора получателей.

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])
Добавлен в: v0.9.4
  • chunk <string> | <Buffer> | <any> Необязательные данные для записи. Для потоков, не работающих в режиме объектов, chunk должно быть строкой или Buffer. Для потоков в режиме объектов 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)
Добавлен в: 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.write(chunk[, encoding][, callback])
Добавлен в: v0.9.4
  • chunk <string> | <Buffer> Данные для записи.
  • encoding <string> Кодировка, если chunk является строкой.
  • callback <Function> Обработчик для завершения обработки данного куска данных.
  • Возвращает: <boolean> false если поток ожидает, что вызывающий код подождёт до момента генерации события 'drain', прежде чем продолжить запись дополнительных данных; в противном случае true.

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

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

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

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

Если данные, которые необходимо записать, могут быть сгенерированы или получены по требованию, рекомендуется инкапсулировать логику в Readable и использовать stream.pipe(). Однако, если вызов write() предпочтительнее, можно учесть обратную связь и избежать проблем с памятью, используя событие '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

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

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

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

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

Два Режима

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

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

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

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

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

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

  • Если нет получателей, вызвав метод stream.pause().
  • Если есть получатели, удалив все обработчики события 'data' и все получатели, вызвав метод 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
  • Возвращает: <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.pipe(destination[, options])
Добавлен в: v0.9.4
  • destination <stream.Writable> Пункт назначения для записи данных
  • options <Объект> Параметры соединения
    • 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 никогда не закрываются до выхода процесса Node.js, независимо от заданных параметров.

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() в пользу использования readable.pipe() или события '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)
Добавлен в: v0.9.11
  • chunk <Buffer> | <string> Фрагмент данных для помещения в начало очереди чтения

Метод 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').StringDecoder;
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 <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').OldReader;
const Readable = require('stream').Readable;
const oreader = new OldReader();
const myReader = new Readable().wrap(oreader);

myReader.on('readable', () => {
  myReader.read(); // etc.
});

Потоки Duplex и Transform

Класс: stream.Duplex

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

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

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

  • TCP-сокеты
  • Потоки zlib
  • Потоки crypto

Класс: stream.Transform

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

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

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

  • Потоки zlib
  • Потоки crypto

API для разработчиков потоков

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

Сначала разработчик потока объявит новый класс JavaScript, который расширяет один из четырёх основных классов потоков (stream.Writable, stream.Readable, stream.Duplex, или stream.Transform), убедившись, что они вызывают соответствующий конструктор родительского класса:

const Writable = require('stream').Writable;

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

Затем новый класс потока должен реализовать один или несколько конкретных методов, в зависимости от типа создаваемого потока, как подробно описано в таблице ниже:

Сценарий использования

Класс

Реализуемый(ые) метод(ы)

Только чтение

Readable

_read

Только запись

Writable

_write, _writev

Чтение и запись

Duplex

_read, _write, _writev

Обработка записанных данных, а затем чтение результата

Transform

_transform, _flush

Примечание: Код реализации потока никогда не должен вызывать "публичные" методы потока, предназначенные для использования потребителями (как описано в разделе API для потребителей потоков). Это может привести к нежелательным побочным эффектам в коде приложения, потребляющем поток.

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

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

Например:

const Writable = require('stream').Writable;

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

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

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

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

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

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

Например:

const Writable = require('stream').Writable;

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

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

const Writable = require('stream').Writable;
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').Writable;

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

writable._write(chunk, encoding, callback)

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

writable._writev(chunks, callback)

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

Примечание: Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен реализовываться дочерними классами и вызываться только внутренними методами класса Writable.

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

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

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

Рекомендуется, чтобы ошибки, возникающие во время обработки методов writable._write() и writable._writev(), сообщались путём вызова обратного вызова и передачи ошибки в качестве первого аргумента. Это вызовет событие 'error' для потока Writable. Бросание исключения Error внутри writable._write() может привести к непредсказуемому и непоследовательному поведению в зависимости от того, как используется поток. Использование обратного вызова гарантирует согласованную и предсказуемую обработку ошибок.

const Writable = require('stream').Writable;

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').Writable;

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

Например:

const Readable = require('stream').Readable;

class MyReadable extends Readable {
  constructor(options) {
    // Calls the stream.Readable(options) constructor
    super(options);
    // ...
  }
}

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

const Readable = require('stream').Readable;
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').Readable;

const myReadable = new Readable({
  read(size) {
    // ...
  }
});

readable._read(size)

  • size <number> Количество байтов для асинхронного чтения

Примечание: Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только методами внутреннего класса 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.push(chunk[, encoding])

  • chunk <Buffer> | <null> | <string> Часть данных для помещения в очередь чтения
  • encoding <string> Кодировка строковых частей. Должна быть допустимой кодировкой Buffer, например, 'utf8' или 'ascii'
  • Возвращает: <boolean> true если можно продолжить помещать дополнительные части данных; false в противном случае.

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

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

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

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

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

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

    this._source = getLowlevelSourceObject();

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

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

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

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

Рекомендуется, чтобы ошибки, возникающие во время обработки метода readable._read(), генерировались событием 'error' вместо выбрасывания исключений. Выбрасывание ошибки изнутри readable._read() может привести к неожиданному и несогласованному поведению, в зависимости от того, работает ли поток в режиме потока или приостановки. Использование события 'error' обеспечивает согласованное и предсказуемое обращение с ошибками.

const Readable = require('stream').Readable;

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').Readable;

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 — это поток, который реализует как Readable, так и Writable, например, соединение TCP-соккета.

Поскольку JavaScript не поддерживает множественное наследование, класс stream.Duplex расширяется для реализации потока Duplex (вместо расширения классов stream.Readable и stream.Writable).

Примечание: Класс stream.Duplex прототипически наследуется от stream.Readable и паразитично от stream.Writable, но instanceof будет работать корректно для обоих базовых классов благодаря переопределению Symbol.hasInstance в stream.Writable.

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

new stream.Duplex(options)

  • options <Object> Передаётся в конструкторы как Writable, так и Readable. Также имеет следующие поля:
    • allowHalfOpen <boolean> По умолчанию true. Если установлено в значение false, поток автоматически завершит сторону записи, когда сторона чтения завершится.
    • readableObjectMode <boolean> По умолчанию false. Устанавливает objectMode для стороны чтения потока. Не имеет эффекта, если objectMode равно true.
    • writableObjectMode <boolean> По умолчанию false. Устанавливает objectMode для стороны записи потока. Не имеет эффекта, если objectMode равно true.

Например:

const Duplex = require('stream').Duplex;

class MyDuplex extends Duplex {
  constructor(options) {
    super(options);
    // ...
  }
}

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

const Duplex = require('stream').Duplex;
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').Duplex;

const myDuplex = new Duplex({
  read(size) {
    // ...
  },
  write(chunk, encoding, callback) {
    // ...
  }
});

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

Ниже приведён простой пример дуплексного потока, который оборачивает гипотетический базовый объект источника, в который можно записывать данные и из которого можно читать данные, но использующий API, несовместимый с Node.js потоками. Следующее иллюстрирует простой пример дуплексного потока, который буферизует входящие данные через интерфейс Writable, который считывается обратно через интерфейс Readable.

const Duplex = require('stream').Duplex;
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));
    });
  }
}

Самый важный аспект дуплексного потока заключается в том, что стороны Readable и Writable работают независимо друг от друга, несмотря на совместное существование в одном объекте.

Дуплексные потоки в режиме объектов

Для дуплексных потоков objectMode может быть установлено исключительно для стороны Readable или Writable соответственно с использованием параметров readableObjectMode и writableObjectMode.

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

const Transform = require('stream').Transform;

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

new stream.Transform([options])

  • options <Объект> Передаётся в конструкторы Writable и Readable. Также содержит следующие поля:
    • transform <Функция> Реализация метода stream._transform().
    • flush <Функция> Реализация метода stream._flush().

Например:

const Transform = require('stream').Transform;

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

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

const Transform = require('stream').Transform;
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').Transform;

const myTransform = new Transform({
  transform(chunk, encoding, callback) {
    // ...
  }
});

События: 'finish' и 'end'

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

transform._flush(callback)

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

Примечание: Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только методами внутреннего класса Readable.

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

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

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

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

transform._transform(chunk, encoding, callback)

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

Примечание: Этот метод НЕ должен вызываться кодом приложения напрямую. Он должен быть реализован дочерними классами и вызываться только методами внутреннего класса Readable.

Все реализации потоков Transform должны предоставлять метод _transform(), чтобы принимать входные данные и производить выходные. Реализация transform._transform() обрабатывает записываемые байты, вычисляет выходные данные, а затем передает эти выходные данные в часть readable с помощью метода 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() никогда не вызывается.
  • Поток не перенаправлен ни на один объект writable.

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

// 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, не фактически потребляя данные. В таких случаях можно вызвать readable.read(0), что всегда вернёт null.

Если внутренний буфер чтения находится ниже highWaterMark, и поток в данный момент не выполняет чтение, то вызов stream.read(0) вызовет низкоуровневый вызов stream._read().

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

readable.push('')

Использование readable.push('') не рекомендуется.

Отправка строки нулевой длины или Buffer в поток, который не находится в режиме обработки объектов, имеет интересный побочный эффект. Поскольку это вызов метода readable.push(), вызов завершит процесс чтения. Однако, поскольку аргумент — пустая строка, данные не добавляются в буфер readable, поэтому пользователь не может их получить.

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-v6.x/docs/api/stream.html

Spec-Zone.ru

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