Spec-Zone.ru › Node.js 22 LTS

Поток[src]

Стабильность: 2 - Стабильный

Исходный код: lib/stream.js

Поток — это абстрактный интерфейс для работы с потоковыми данными в Node.js. Модуль node:stream предоставляет API для реализации интерфейса потоков.

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

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

Чтобы получить доступ к модулю node:stream:

const stream = require('node:stream'); copy

Модуль node:stream полезен для создания новых типов экземпляров потоков. Обычно нет необходимости использовать модуль node:stream для чтения потоков.

Структура этого документа

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

Типы потоков

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

  • Writable: потоки, в которые можно записывать данные (например, fs.createWriteStream()).
  • Readable: потоки, из которых можно читать данные (например, fs.createReadStream()).
  • Duplex: потоки, которые одновременно являются Readable и Writable (например, net.Socket).
  • Transform: потоки Duplex, которые могут изменять или преобразовывать данные при их записи и чтении (например, zlib.createDeflate()).

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

API потоков на основе промисов

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

API stream/promises предоставляет альтернативный набор асинхронных вспомогательных функций для потоков, которые возвращают объекты Promise вместо использования обратных вызовов. Доступ к API можно получить через require('node:stream/promises') или require('node:stream').promises.

stream.pipeline(streams[, options])

stream.pipeline(source[, ...transforms], destination[, options])

История
Версия Изменения
v18.0.0, v17.2.0, v16.14.0

Добавлена опция end, которой можно присвоить значение false, чтобы предотвратить автоматическое закрытие потока назначения при завершении исходного потока.

v15.0.0

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

  • streams <Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]>
  • source <Stream> | <Iterable> | <AsyncIterable> | <Function>
    • Возвращает: <Promise> | <AsyncIterable>
  • ...transforms <Stream> | <Function>
    • source <AsyncIterable>
    • Возвращает: <Promise> | <AsyncIterable>
  • destination <Stream> | <Function>
    • source <AsyncIterable>
    • Возвращает: <Promise> | <AsyncIterable>
  • options <Object> Параметры конвейера
    • signal <AbortSignal>
    • end <boolean> Завершать поток назначения при завершении исходного потока. Потоки преобразования завершаются всегда, даже если это значение равно false. По умолчанию: true.
  • Возвращает: <Promise> Выполняется после завершения конвейера.
CommonJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
const zlib = require('node:zlib');

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

run().catch(console.error);
Модули JavaScript
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';

await pipeline(
  createReadStream('archive.tar'),
  createGzip(),
  createWriteStream('archive.tar.gz'),
);
console.log('Pipeline succeeded.');

Чтобы использовать AbortSignal, передайте его в объекте параметров последним аргументом. При отмене сигнала для базового конвейера будет вызван destroy с объектом AbortError.

CommonJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');
const zlib = require('node:zlib');

async function run() {
  const ac = new AbortController();
  const signal = ac.signal;

  setImmediate(() => ac.abort());
  await pipeline(
    fs.createReadStream('archive.tar'),
    zlib.createGzip(),
    fs.createWriteStream('archive.tar.gz'),
    { signal },
  );
}

run().catch(console.error); // AbortError
Модули JavaScript
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';
import { createGzip } from 'node:zlib';

const ac = new AbortController();
const { signal } = ac;
setImmediate(() => ac.abort());
try {
  await pipeline(
    createReadStream('archive.tar'),
    createGzip(),
    createWriteStream('archive.tar.gz'),
    { signal },
  );
} catch (err) {
  console.error(err); // AbortError
}

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

CommonJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');

async function run() {
  await pipeline(
    fs.createReadStream('lowercase.txt'),
    async function* (source, { signal }) {
      source.setEncoding('utf8');  // Work with strings rather than `Buffer`s.
      for await (const chunk of source) {
        yield await processChunk(chunk, { signal });
      }
    },
    fs.createWriteStream('uppercase.txt'),
  );
  console.log('Pipeline succeeded.');
}

run().catch(console.error);
Модули JavaScript
import { pipeline } from 'node:stream/promises';
import { createReadStream, createWriteStream } from 'node:fs';

await pipeline(
  createReadStream('lowercase.txt'),
  async function* (source, { signal }) {
    source.setEncoding('utf8');  // Work with strings rather than `Buffer`s.
    for await (const chunk of source) {
      yield await processChunk(chunk, { signal });
    }
  },
  createWriteStream('uppercase.txt'),
);
console.log('Pipeline succeeded.');

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

CommonJS
const { pipeline } = require('node:stream/promises');
const fs = require('node:fs');

async function run() {
  await pipeline(
    async function* ({ signal }) {
      await someLongRunningfn({ signal });
      yield 'asd';
    },
    fs.createWriteStream('uppercase.txt'),
  );
  console.log('Pipeline succeeded.');
}

run().catch(console.error);
Модули JavaScript
import { pipeline } from 'node:stream/promises';
import fs from 'node:fs';
await pipeline(
  async function* ({ signal }) {
    await someLongRunningfn({ signal });
    yield 'asd';
  },
  fs.createWriteStream('uppercase.txt'),
);
console.log('Pipeline succeeded.');

API pipeline также предоставляет версию с обратным вызовом:

stream.finished(stream[, options])

История
Версия Изменения
v19.5.0, v18.14.0

Добавлена поддержка ReadableStream и WritableStream.

v19.1.0, v18.13.0

Добавлена опция cleanup.

v15.0.0

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

  • stream <Stream> | <ReadableStream> | <WritableStream> Читаемый и/или записываемый поток/веб-поток.
  • options <Object>
    • error <boolean> | <undefined>
    • readable <boolean> | <undefined>
    • writable <boolean> | <undefined>
    • signal <AbortSignal> | <undefined>
    • cleanup <boolean> | <undefined> Если true, удаляет зарегистрированные этой функцией обработчики событий до выполнения промиса. По умолчанию: false.
  • Возвращает: <Promise> Выполняется, когда поток становится нечитаемым или недоступным для записи.
CommonJS
const { finished } = require('node:stream/promises');
const fs = require('node:fs');

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.
Модули JavaScript
import { finished } from 'node:stream/promises';
import { createReadStream } from 'node:fs';

const rs = createReadStream('archive.tar');

async function run() {
  await finished(rs);
  console.log('Stream is done reading.');
}

run().catch(console.error);
rs.resume(); // Drain the stream.

API finished также предоставляет версию с обратным вызовом.

stream.finished() оставляет висящие обработчики событий (в частности, 'error', 'end', 'finish' и 'close') после выполнения или отклонения возвращённого промиса. Это сделано для того, чтобы неожиданные события 'error' (из-за неправильной реализации потоков) не приводили к неожиданным сбоям. Если такое поведение нежелательно, следует установить для options.cleanup значение true:

await finished(rs, { cleanup: true }); copy

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

Все потоки, создаваемые API Node.js, работают исключительно со строками, объектами <Buffer>, <TypedArray> и <DataView>:

  • Strings и Buffers — наиболее распространённые типы, используемые в потоках.
  • TypedArray и DataView позволяют работать с двоичными данными в таких типах, как Int32Array или Uint8Array. Когда в поток записывается TypedArray или DataView, Node.js обрабатывает необработанные байты.

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

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

Буферизация

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

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

Данные буферизуются в потоках 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('node: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', "not json" is not valid JSON copy

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

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

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

Потоки Duplex и Transform являются одновременно потоками Writable и Readable.

Приложениям, которые записывают данные в поток или считывают их из потока, не требуется напрямую реализовывать интерфейсы потоков, и, как правило, у них нет причин вызывать require('node: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'); copy
Класс: 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);
    }
  }
} copy
Событие: 'error'
Добавлено в: v0.9.4
  • Тип: <Error>

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

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

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

const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
  console.log('Something has stopped piping into the writer.');
  assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer); copy
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'. <Error>
  • Возвращает: <this>

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

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

const myStream = new Writable();

const fooErr = new Error('foo error');
myStream.destroy(fooErr);
myStream.on('error', (fooErr) => console.error(fooErr.message)); // foo error copy
const { Writable } = require('node:stream');

const myStream = new Writable();

myStream.destroy();
myStream.on('error', function wontHappen() {}); copy
const { Writable } = require('node:stream');

const myStream = new Writable();
myStream.destroy();

myStream.write('foo', (error) => console.error(error.code));
// ERR_STREAM_DESTROYED copy

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

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

writable.closed
Добавлено в: v18.0.0
  • Тип: <boolean>

Равно true после генерации события 'close'.

writable.destroyed
Добавлено в: v8.0.0
  • Тип: <boolean>

Равно true после вызова writable.destroy().

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

const myStream = new Writable();

console.log(myStream.destroyed); // false
myStream.destroy();
console.log(myStream.destroyed); // true copy
writable.end([chunk[, encoding]][, callback])
История
Версия Изменения
v22.0.0

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

v15.0.0

Вызов callback выполняется до 'finish' или при ошибке.

v14.0.0

Вызов callback выполняется при генерации 'finish' или 'error'.

v10.0.0

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

v8.0.0

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

v0.9.4

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

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

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

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

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

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

v0.11.15

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

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

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

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

Метод writable.uncork() сбрасывает все данные, буферизованные с момента вызова stream.cork().

При использовании writable.cork() и writable.uncork() для управления буферизацией записей в поток откладывайте вызовы writable.uncork() с помощью process.nextTick(). Это позволяет объединять все вызовы writable.write(), происходящие в течение одной фазы цикла событий Node.js.

stream.cork();
stream.write('some ');
stream.write('data ');
process.nextTick(() => stream.uncork()); copy

Если метод 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();
}); copy

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

writable.writable
Добавлено в: v11.4.0
  • Тип: <boolean>

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

writable.writableAborted
История
Версия Изменения
v22.17.0

API помечен как стабильный.

v18.0.0, v16.17.0

Добавлено в: v18.0.0, v16.17.0

  • Тип: <boolean>

Возвращает признак того, что поток был уничтожен или в нём произошла ошибка до генерации 'finish'.

writable.writableEnded
Добавлено в: v12.9.0
  • Тип: <boolean>

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

writable.writableCorked
Добавлено в: v13.2.0, v12.16.0
  • Тип: <integer>

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

writable.errored
Добавлено в: v18.0.0
  • Тип: <Error>

Возвращает ошибку, если поток был уничтожен с ошибкой.

writable.writableFinished
Добавлено в: v12.6.0
  • Тип: <boolean>

Устанавливается в true непосредственно перед генерацией события 'finish'.

writable.writableHighWaterMark
Добавлено в: v9.3.0
  • Тип: <number>

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

writable.writableLength
Добавлено в: v9.4.0
  • Тип: <number>

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

writable.writableNeedDrain
Добавлено в: v15.2.0, v14.17.0
  • Тип: <boolean>

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

writable.writableObjectMode
Добавлено в: v12.3.0
  • Тип: <boolean>

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

writable[Symbol.asyncDispose]()
Добавлено в: v22.4.0
Стабильность: 1 - Экспериментальный

Вызывает writable.destroy() с объектом AbortError и возвращает промис, который выполняется после завершения потока.

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

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

v8.0.0

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

v6.0.0

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

v0.9.4

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

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

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

Возвращаемое значение равно true, если после принятия chunk внутренний буфер меньше значения highWaterMark, заданного при создании потока. Если возвращено 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.');
}); copy

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

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

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

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

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

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

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

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

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

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

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

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

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('node: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()); });
// readableFlowing is still false.
pass.write('ok');  // Will not emit 'data'.
pass.resume();     // Must be called to make stream emit 'data'.
// readableFlowing is now true. copy

Пока readable.readableFlowing имеет значение false, данные могут накапливаться во внутреннем буфере потока.

Выберите один способ работы с API

API потока Readable развивался в разных версиях Node.js и предоставляет несколько способов потребления данных потока. В целом разработчикам следует выбрать один способ потребления данных и не следует использовать несколько способов для потребления данных из одного потока. В частности, сочетание on('data'), on('readable'), pipe() или асинхронных итераторов может привести к неожиданному поведению.

Класс: 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> | <string> | <any> Фрагмент данных. Для потоков, работающих не в объектном режиме, фрагментом будет строка или 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.`);
}); copy
Событие: '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.');
}); copy
Событие: 'error'
Добавлено в: v0.9.4
  • Тип: <Error>

Событие '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' испускается, когда в потоке есть данные для чтения, объём которых не превышает настроенный предел буфера (state.highWaterMark). Фактически оно указывает на то, что в буфере потока появилась новая информация. Если в буфере есть данные, их можно получить вызовом stream.read(). Кроме того, событие 'readable' может испускаться и при достижении конца потока.

const readable = getReadableStreamSomehow();
readable.on('readable', function() {
  // There is some data to read now.
  let data;

  while ((data = this.read()) !== null) {
    console.log(data);
  }
}); copy

Если достигнут конец потока, вызов stream.read() вернёт null и вызовет событие 'end'. Это также верно, если данных для чтения не было вовсе. Например, в следующем примере файл foo.txt пуст:

const fs = require('node:fs');
const rr = fs.createReadStream('foo.txt');
rr.on('readable', () => {
  console.log(`readable: ${rr.read()}`);
});
rr.on('end', () => {
  console.log('end');
}); copy

При запуске этого скрипта вывод будет следующим:

$ node test.js
readable: null
end copy

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

В целом механизмы событий 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> Ошибка, которая будет передана в качестве данных события 'error'
  • Возвращает: <this>

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

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

Разработчикам не следует переопределять этот метод; вместо этого следует реализовать readable._destroy().

readable.closed
Добавлено в: v18.0.0
  • Тип: <boolean>

Равно true после испускания 'close'.

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 copy
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);
}); copy

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

readable.pipe(destination[, options])
Добавлено в: v0.9.4
  • destination <stream.Writable> Назначение для записи данных
  • options <Object> Параметры канала
    • end <boolean> Завершать записывающий поток после завершения читающего. По умолчанию: true.
  • Возвращает: <stream.Writable> Назначение, позволяющее выстраивать цепочки каналов, если это поток Duplex или Transform

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

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

const fs = require('node:fs');
const readable = getReadableStreamSomehow();
const writable = fs.createWriteStream('file.txt');
// All the data from readable goes into 'file.txt'.
readable.pipe(writable); copy

К одному потоку Readable можно подключить несколько потоков Writable.

Метод readable.pipe() возвращает ссылку на поток назначения, позволяя создавать цепочки потоков, соединённых каналами:

const fs = require('node:fs');
const zlib = require('node:zlib');
const r = fs.createReadStream('file.txt');
const z = zlib.createGzip();
const w = fs.createWriteStream('file.txt.gz');
r.pipe(z).pipe(w); copy

По умолчанию для потока назначения Writable вызывается stream.end(), когда исходный поток Readable испускает событие 'end', после чего в назначение больше нельзя записывать данные. Чтобы отключить такое поведение по умолчанию и оставить поток назначения открытым, параметр end можно передать со значением false:

reader.pipe(writer, { end: false });
reader.on('end', () => {
  writer.end('Goodbye\n');
}); copy

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

Потоки Writable process.stderr и process.stdout не закрываются до завершения процесса Node.js независимо от указанных параметров.

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

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

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

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

Значение аргумента size должно быть не больше 1 ГиБ.

Метод readable.read() следует вызывать только для потоков Readable, работающих в режиме паузы. В режиме передачи данных readable.read() вызывается автоматически, пока внутренний буфер не будет полностью опустошён.

const readable = getReadableStreamSomehow();

// 'readable' 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.');
}); copy

Каждый вызов readable.read() возвращает фрагмент данных или null, означающий, что в данный момент данных для чтения больше нет. Эти фрагменты не объединяются автоматически. Поскольку один вызов read() не возвращает все данные, может потребоваться цикл 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('');
}); copy

Поток 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.readableAborted
История
Версия Изменения
v22.17.0

API объявлен стабильным.

v16.8.0

Добавлено в: v16.8.0

  • Тип: <boolean>

Возвращает, был ли поток уничтожен или в нём произошла ошибка до испускания 'end'.

readable.readableDidRead
История
Версия Изменения
v22.17.0

API объявлен стабильным.

v16.7.0, v14.18.0

Добавлено в: v16.7.0, v14.18.0

  • Тип: <boolean>

Возвращает, было ли испущено 'data'.

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

Геттер свойства encoding указанного потока Readable. Свойство encoding можно задать с помощью метода readable.setEncoding().

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

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

readable.errored
Добавлено в: v18.0.0
  • Тип: <Error>

Возвращает ошибку, если поток был уничтожен из-за ошибки.

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

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

getReadableStreamSomehow()
  .resume()
  .on('end', () => {
    console.log('Reached the end, but did not read anything.');
  }); copy

Метод 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);
}); copy
readable.unpipe([destination])
Добавлено в: v0.9.4
  • destination <stream.Writable> Необязательный поток, который требуется отсоединить от канала
  • Возвращает: <this>

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

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

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

const fs = require('node: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); copy
readable.unshift(chunk[, encoding])
История
Версия Изменения
v22.0.0

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

v8.0.0

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

v0.9.11

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

  • chunk <Buffer> | <TypedArray> | <DataView> | <string> | <null> | <any> Фрагмент данных, добавляемый в начало очереди чтения. Для потоков, работающих не в объектном режиме, chunk должен быть значением типа <string>, <Buffer>, <TypedArray>, <DataView> или 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('node: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.includes('\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);
        return;
      }
      // Still reading the header.
      header += str;
    }
  }
} copy

В отличие от stream.push(chunk), stream.unshift(chunk) не завершает процесс чтения сбросом внутреннего состояния чтения потока. Это может привести к неожиданным результатам, если readable.unshift() вызван во время чтения (например, из реализации stream._read() пользовательского потока). Немедленный вызов stream.push('') после readable.unshift() правильно сбросит состояние чтения, однако лучше вообще не вызывать readable.unshift() во время чтения.

readable.wrap(stream)
Добавлено в: v0.9.4
  • stream <Stream> Читаемый поток «старого образца»
  • Возвращает: <this>

До Node.js 0.10 потоки не реализовывали весь API модуля node:stream в его нынешнем виде. (Дополнительные сведения см. в разделе Совместимость.)

При использовании старой библиотеки Node.js, испускающей события 'data' и имеющей лишь рекомендательный метод stream.pause(), метод readable.wrap() позволяет создать поток Readable, использующий старый поток в качестве источника данных.

Метод readable.wrap() требуется редко, но он предоставлен для удобства взаимодействия со старыми приложениями и библиотеками Node.js.

const { OldReader } = require('./old-api-module.js');
const { Readable } = require('node:stream');
const oreader = new OldReader();
const myReader = new Readable().wrap(oreader);

myReader.on('readable', () => {
  myReader.read(); // etc.
}); copy
readable[Symbol.asyncIterator]()
История
Версия Изменения
v11.14.0

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

v10.0.0

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

  • Возвращает: <AsyncIterator> для полного чтения потока.
const fs = require('node: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); copy

Если цикл завершается с помощью break, return или throw, поток будет уничтожен. Иными словами, итерация по потоку приводит к его полному чтению. Поток будет считываться фрагментами размером, равным параметру highWaterMark. В приведённом выше примере кода данные будут считаны одним фрагментом, если размер файла меньше 64 КиБ, поскольку для fs.createReadStream() не указан параметр highWaterMark.

readable[Symbol.asyncDispose]()
Добавлено в: v20.4.0, v18.18.0
Стабильность: 1 — Экспериментальный

Вызывает readable.destroy() с AbortError и возвращает промис, который выполняется после завершения потока.

readable.compose(stream[, options])
История
Версия Изменения
v22.17.0

API объявлен стабильным.

v19.1.0, v18.13.0

Добавлено в: v19.1.0, v18.13.0

  • stream <Stream> | <Iterable> | <AsyncIterable> | <Function>
  • options <Object>
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Duplex> поток, объединённый с потоком stream.
import { Readable } from 'node:stream';

async function* splitToWords(source) {
  for await (const chunk of source) {
    const words = String(chunk).split(' ');

    for (const word of words) {
      yield word;
    }
  }
}

const wordsStream = Readable.from(['this is', 'compose as operator']).compose(splitToWords);
const words = await wordsStream.toArray();

console.log(words); // prints ['this', 'is', 'compose', 'as', 'operator'] copy

Дополнительные сведения см. в разделе stream.compose.

readable.iterator([options])
История
Версия Изменения
v22.17.0

API объявлен стабильным.

v16.3.0

Добавлено в: v16.3.0

  • options <Object>
    • destroyOnReturn <boolean> Если задано значение false, вызов return для асинхронного итератора или выход из итерации for await...of с помощью break, return или throw не приведёт к уничтожению потока. По умолчанию: true.
  • Возвращает: <AsyncIterator> для чтения потока.

Итератор, созданный этим методом, позволяет отменить уничтожение потока, если цикл for await...of завершается с помощью return, break или throw, либо указать, что итератор должен уничтожить поток, если во время итерации в потоке возникла ошибка.

const { Readable } = require('node:stream');

async function printIterator(readable) {
  for await (const chunk of readable.iterator({ destroyOnReturn: false })) {
    console.log(chunk); // 1
    break;
  }

  console.log(readable.destroyed); // false

  for await (const chunk of readable.iterator({ destroyOnReturn: false })) {
    console.log(chunk); // Will print 2 and then 3
  }

  console.log(readable.destroyed); // True, stream was totally consumed
}

async function printSymbolAsyncIterator(readable) {
  for await (const chunk of readable) {
    console.log(chunk); // 1
    break;
  }

  console.log(readable.destroyed); // true
}

async function showBoth() {
  await printIterator(Readable.from([1, 2, 3]));
  await printSymbolAsyncIterator(Readable.from([1, 2, 3]));
}

showBoth(); copy
readable.map(fn[, options])
История
Версия Изменения
v20.7.0, v18.19.0

В параметры добавлен highWaterMark.

v17.4.0, v16.14.0

Добавлено в: v17.4.0, v16.14.0

Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция, применяемая к каждому фрагменту потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • highWaterMark <number> количество элементов, буферизуемых в ожидании обработки сопоставленных элементов пользователем. По умолчанию: concurrency * 2 - 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Readable> поток, преобразованный функцией fn.

Этот метод позволяет преобразовать поток. Функция fn вызывается для каждого фрагмента потока. Если функция fn возвращает промис, этот промис будет обработан с помощью await перед передачей результата в выходной поток.

import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';

// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).map((x) => x * 2)) {
  console.log(chunk); // 2, 4, 6, 8
}
// With an asynchronous mapper, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
  'nodejs.org',
  'openjsf.org',
  'www.linuxfoundation.org',
]).map((domain) => resolver.resolve4(domain), { concurrency: 2 });
for await (const result of dnsResults) {
  console.log(result); // Logs the DNS result of resolver.resolve4.
} copy
readable.filter(fn[, options])
История
Версия Изменения
v20.7.0, v18.19.0

В параметры добавлен highWaterMark.

v17.4.0, v16.14.0

Добавлено в: v17.4.0, v16.14.0

Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция для фильтрации фрагментов потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • highWaterMark <number> количество элементов, буферизуемых в ожидании обработки отфильтрованных элементов пользователем. По умолчанию: concurrency * 2 - 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Readable> поток, отфильтрованный предикатом fn.

Этот метод позволяет фильтровать поток. Для каждого фрагмента потока вызывается функция fn; если она возвращает истинное значение, фрагмент передаётся в выходной поток. Если функция fn возвращает промис, этот промис будет обработан с помощью await.

import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';

// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
  console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
  'nodejs.org',
  'openjsf.org',
  'www.linuxfoundation.org',
]).filter(async (domain) => {
  const { address } = await resolver.resolve4(domain, { ttl: true });
  return address.ttl > 60;
}, { concurrency: 2 });
for await (const result of dnsResults) {
  // Logs domains with more than 60 seconds on the resolved dns record.
  console.log(result);
} copy
readable.forEach(fn[, options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Promise> промис, выполняющийся после завершения потока.

Этот метод позволяет перебирать поток. Для каждого фрагмента потока вызывается функция fn. Если функция fn возвращает промис, этот промис будет обработан с помощью await.

Этот метод отличается от циклов for await...of тем, что позволяет при необходимости обрабатывать фрагменты параллельно. Кроме того, итерацию forEach можно остановить только посредством передачи параметра signal и прерывания связанного с ним AbortController, тогда как for await...of можно остановить с помощью break или return. В обоих случаях поток будет уничтожен.

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

import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';

// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
  console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
  'nodejs.org',
  'openjsf.org',
  'www.linuxfoundation.org',
]).map(async (domain) => {
  const { address } = await resolver.resolve4(domain, { ttl: true });
  return address;
}, { concurrency: 2 });
await dnsResults.forEach((result) => {
  // Logs result, similar to `for await (const result of dnsResults)`
  console.log(result);
});
console.log('done'); // Stream has finished copy
readable.toArray([options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • options <Object>
    • signal <AbortSignal> позволяет отменить операцию toArray при прерывании сигнала.
  • Возвращает: <Promise> промис, содержащий массив с содержимым потока.

Этот метод позволяет легко получить содержимое потока.

Поскольку этот метод считывает весь поток в память, он сводит на нет преимущества потоков. Он предназначен для совместимости и удобства, а не для основного способа чтения потоков.

import { Readable } from 'node:stream';
import { Resolver } from 'node:dns/promises';

await Readable.from([1, 2, 3, 4]).toArray(); // [1, 2, 3, 4]

// Make dns queries concurrently using .map and collect
// the results into an array using toArray
const dnsResults = await Readable.from([
  'nodejs.org',
  'openjsf.org',
  'www.linuxfoundation.org',
]).map(async (domain) => {
  const { address } = await resolver.resolve4(domain, { ttl: true });
  return address;
}, { concurrency: 2 }).toArray(); copy
readable.some(fn[, options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Promise> промис, результатом которого будет true, если fn вернуло истинное значение хотя бы для одного из фрагментов.

Этот метод похож на Array.prototype.some и вызывает fn для каждого фрагмента потока, пока ожидаемое возвращаемое значение не станет true (или любым другим истинным значением). Как только ожидаемое возвращаемое значение вызова fn для фрагмента становится истинным, поток уничтожается, а промис выполняется со значением true. Если ни один из вызовов fn для фрагментов не возвращает истинное значение, промис выполняется со значением false.

import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';

// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).some((x) => x > 2); // true
await Readable.from([1, 2, 3, 4]).some((x) => x < 0); // false

// With an asynchronous predicate, making at most 2 file checks at a time.
const anyBigFile = await Readable.from([
  'file1',
  'file2',
  'file3',
]).some(async (fileName) => {
  const stats = await stat(fileName);
  return stats.size > 1024 * 1024;
}, { concurrency: 2 });
console.log(anyBigFile); // `true` if any file in the list is bigger than 1MB
console.log('done'); // Stream has finished copy
readable.find(fn[, options])
Добавлено в: v17.5.0, v16.17.0
Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Promise> промис, результатом которого будет первый фрагмент, для которого fn вернуло истинное значение, либо undefined, если элемент не найден.

Этот метод похож на Array.prototype.find и вызывает fn для каждого фрагмента потока, чтобы найти фрагмент, для которого fn имеет истинное значение. Как только ожидаемое возвращаемое значение вызова fn становится истинным, поток уничтожается, а промис выполняется со значением, для которого fn вернуло истинное значение. Если все вызовы fn для фрагментов возвращают ложное значение, промис выполняется со значением undefined.

import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';

// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).find((x) => x > 2); // 3
await Readable.from([1, 2, 3, 4]).find((x) => x > 0); // 1
await Readable.from([1, 2, 3, 4]).find((x) => x > 10); // undefined

// With an asynchronous predicate, making at most 2 file checks at a time.
const foundBigFile = await Readable.from([
  'file1',
  'file2',
  'file3',
]).find(async (fileName) => {
  const stats = await stat(fileName);
  return stats.size > 1024 * 1024;
}, { concurrency: 2 });
console.log(foundBigFile); // File name of large file, if any file in the list is bigger than 1MB
console.log('done'); // Stream has finished copy
readable.every(fn[, options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция, вызываемая для каждого фрагмента потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Promise> промис, результатом которого будет true, если fn вернуло истинное значение для всех фрагментов.

Этот метод похож на Array.prototype.every и вызывает fn для каждого фрагмента потока, чтобы проверить, являются ли все ожидаемые возвращаемые значения истинными для fn. Как только ожидаемое возвращаемое значение вызова fn для фрагмента становится ложным, поток уничтожается, а промис выполняется со значением false. Если все вызовы fn для фрагментов возвращают истинное значение, промис выполняется со значением true.

import { Readable } from 'node:stream';
import { stat } from 'node:fs/promises';

// With a synchronous predicate.
await Readable.from([1, 2, 3, 4]).every((x) => x > 2); // false
await Readable.from([1, 2, 3, 4]).every((x) => x > 0); // true

// With an asynchronous predicate, making at most 2 file checks at a time.
const allBigFiles = await Readable.from([
  'file1',
  'file2',
  'file3',
]).every(async (fileName) => {
  const stats = await stat(fileName);
  return stats.size > 1024 * 1024;
}, { concurrency: 2 });
// `true` if all files in the list are bigger than 1MiB
console.log(allBigFiles);
console.log('done'); // Stream has finished copy
readable.flatMap(fn[, options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncGeneratorFunction> | <AsyncFunction> функция, применяемая к каждому фрагменту потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • options <Object>
    • concurrency <number> максимальное число одновременных вызовов fn для обработки потока. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Readable> поток, преобразованный функцией fn с разворачиванием результата.

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

Функция fn может возвращать поток, другой итерируемый объект или асинхронный итерируемый объект; результирующие потоки будут объединены (развёрнуты) в возвращаемый поток.

import { Readable } from 'node:stream';
import { createReadStream } from 'node:fs';

// With a synchronous mapper.
for await (const chunk of Readable.from([1, 2, 3, 4]).flatMap((x) => [x, x])) {
  console.log(chunk); // 1, 1, 2, 2, 3, 3, 4, 4
}
// With an asynchronous mapper, combine the contents of 4 files
const concatResult = Readable.from([
  './1.mjs',
  './2.mjs',
  './3.mjs',
  './4.mjs',
]).flatMap((fileName) => createReadStream(fileName));
for await (const result of concatResult) {
  // This will contain the contents (all chunks) of all 4 files
  console.log(result);
} copy
readable.drop(limit[, options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • limit <number> количество фрагментов, отбрасываемых из читаемого потока.
  • options <Object>
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Readable> поток, из которого удалены limit фрагментов.

Этот метод возвращает новый поток, отбросив первые limit фрагментов.

import { Readable } from 'node:stream';

await Readable.from([1, 2, 3, 4]).drop(2).toArray(); // [3, 4] copy
readable.take(limit[, options])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • limit <number> количество фрагментов, извлекаемых из читаемого потока.
  • options <Object>
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Readable> поток, содержащий limit извлечённых фрагментов.

Этот метод возвращает новый поток с первыми limit фрагментами.

import { Readable } from 'node:stream';

await Readable.from([1, 2, 3, 4]).take(2).toArray(); // [1, 2] copy
readable.reduce(fn[, initial[, options]])
Добавлено в: v17.5.0, v16.15.0
Стабильность: 1 — Экспериментальный
  • fn <Function> | <AsyncFunction> функция-редуктор, вызываемая для каждого фрагмента потока.
    • previous <any> значение, полученное при последнем вызове fn, или значение initial, если оно указано, либо первый фрагмент потока.
    • data <any> фрагмент данных из потока.
    • options <Object>
      • signal <AbortSignal> прерывается при уничтожении потока, что позволяет досрочно прервать вызов fn.
  • initial <any> начальное значение для редукции.
  • options <Object>
    • signal <AbortSignal> позволяет уничтожить поток при прерывании сигнала.
  • Возвращает: <Promise> промис с итоговым значением редукции.

Этот метод последовательно вызывает fn для каждого фрагмента потока, передавая результат вычисления для предыдущего элемента. Он возвращает промис с итоговым значением редукции.

Если начальное значение initial не задано, в качестве начального значения используется первый фрагмент потока. Если поток пуст, промис отклоняется с ошибкой TypeError, у которой свойство кода имеет значение ERR_INVALID_ARGS.

import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';

const directoryPath = './src';
const filesInDir = await readdir(directoryPath);

const folderSize = await Readable.from(filesInDir)
  .reduce(async (totalSize, file) => {
    const { size } = await stat(join(directoryPath, file));
    return totalSize + size;
  }, 0);

console.log(folderSize); copy

Функция-редуктор перебирает элементы потока по одному, поэтому параметр concurrency и параллельная обработка отсутствуют. Чтобы выполнить reduce параллельно, можно вынести асинхронную функцию в метод readable.map.

import { Readable } from 'node:stream';
import { readdir, stat } from 'node:fs/promises';
import { join } from 'node:path';

const directoryPath = './src';
const filesInDir = await readdir(directoryPath);

const folderSize = await Readable.from(filesInDir)
  .map((file) => stat(join(directoryPath, file)), { concurrency: 2 })
  .reduce((totalSize, { size }) => totalSize + size, 0);

console.log(folderSize); copy

Двунаправленные потоки и потоки преобразования

Класс: stream.Duplex
История
Версия Изменения
v6.8.0

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

v0.9.4

Добавлено в версии v0.9.4

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

Примеры потоков Duplex:

  • TCP-сокеты
  • потоки zlib
  • криптографические потоки
duplex.allowHalfOpen
Добавлено в версии v0.9.4
  • Тип: <boolean>

Если false, то поток автоматически завершит записываемую сторону при завершении читаемой стороны. Изначально значение задаётся параметром конструктора allowHalfOpen, по умолчанию равным true.

Это значение можно изменить вручную, чтобы изменить поведение полузакрытого режима существующего экземпляра потока Duplex, но сделать это необходимо до вызова события 'end'.

Класс: stream.Transform
Добавлено в версии v0.9.4

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

Примеры потоков Transform:

  • потоки zlib
  • криптографические потоки
transform.destroy([error])
История
Версия Изменения
v14.0.0

Ничего не делает, если поток уже уничтожен.

v8.0.0

Добавлено в версии v8.0.0

  • error <Error>
  • Возвращает: <this>

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

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

stream.duplexPair([options])
Добавлено в версии v22.6.0
  • options <Object> Значение, передаваемое обоим конструкторам Duplex для задания таких параметров, как буферизация.
  • Возвращает: <Array> из двух экземпляров Duplex.

Вспомогательная функция duplexPair возвращает массив из двух элементов, каждый из которых является потоком Duplex, соединённым с другой стороной:

const [ sideA, sideB ] = duplexPair(); copy

Всё, что записывается в один поток, становится доступным для чтения из другого. Это поведение аналогично сетевому соединению, в котором данные, записанные клиентом, становятся доступными для чтения сервером, и наоборот.

Двунаправленные потоки симметричны: любой из них можно использовать без различий в поведении.

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

История
Версия Изменения
v19.5.0

Добавлена поддержка ReadableStream и WritableStream.

v15.11.0

Добавлен параметр signal.

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 <Stream> | <ReadableStream> | <WritableStream> Читаемый и/или записываемый поток/веб-поток.
  • options <Object>
    • error <boolean> Если задано значение false, вызов emit('error', err) не считается завершением. По умолчанию: true.
    • readable <boolean> Если задано значение false, обратный вызов будет вызван при завершении потока, даже если поток всё ещё может быть доступен для чтения. По умолчанию: true.
    • writable <boolean> Если задано значение false, обратный вызов будет вызван при завершении потока, даже если в поток всё ещё можно записывать. По умолчанию: true.
    • signal <AbortSignal> позволяет прервать ожидание завершения потока. Базовый поток не будет прерван, если сигнал будет прерван. Обратный вызов будет вызван с ошибкой AbortError. Все зарегистрированные этой функцией обработчики также будут удалены.
  • callback <Function> Функция обратного вызова, принимающая необязательный аргумент ошибки.
  • Возвращает: <Function> Функция очистки, удаляющая все зарегистрированные обработчики.

Функция для получения уведомления о том, что поток больше не доступен для чтения или записи, либо в нём произошла ошибка или преждевременное закрытие.

const { finished } = require('node:stream');
const fs = require('node:fs');

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. copy

Особенно полезна при обработке ошибок, когда поток уничтожается преждевременно (например, при прерванном HTTP-запросе) и не генерирует 'end' или 'finish'.

API finished предоставляет версию на основе промисов.

stream.finished() оставляет необработанные слушатели событий (в частности, 'error', 'end', 'finish' и 'close') после вызова callback. Это сделано для того, чтобы неожиданные события 'error' (из-за некорректной реализации потока) не приводили к неожиданным сбоям. Если такое поведение нежелательно, в обратном вызове необходимо вызвать возвращённую функцию очистки:

const cleanup = finished(rs, (err) => {
  cleanup();
  // ...
}); copy

stream.pipeline(source[, ...transforms], destination, callback)

stream.pipeline(streams, callback)

История
Версия Изменения
v19.7.0, v18.16.0

Добавлена поддержка веб-потоков.

v18.0.0

Передача недопустимого обратного вызова в аргумент callback теперь вызывает ERR_INVALID_ARG_TYPE вместо ERR_INVALID_CALLBACK.

v14.0.0

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

v13.10.0

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

v10.0.0

Добавлено в версии v10.0.0

  • streams <Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]>
  • source <Stream> | <Iterable> | <AsyncIterable> | <Function> | <ReadableStream>
    • Возвращает: <Iterable> | <AsyncIterable>
  • ...transforms <Stream> | <Function> | <TransformStream>
    • source <AsyncIterable>
    • Возвращает: <AsyncIterable>
  • destination <Stream> | <Function> | <WritableStream>
    • source <AsyncIterable>
    • Возвращает: <AsyncIterable> | <Promise>
  • callback <Function> Вызывается после полного завершения конвейера.
    • err <Error>
    • val Разрешённое значение Promise, возвращённого функцией destination.
  • Возвращает: <Stream>

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

const { pipeline } = require('node:stream');
const fs = require('node:fs');
const zlib = require('node: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.');
    }
  },
); copy

API pipeline предоставляет версию на основе промисов.

stream.pipeline() вызовет stream.destroy(err) для всех потоков, кроме:

  • Readable потоков, которые сгенерировали 'end' или 'close'.
  • Writable потоков, которые сгенерировали 'finish' или 'close'.

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

stream.pipeline() закрывает все потоки при возникновении ошибки. Использование IncomingRequest с pipeline может привести к неожиданному поведению: сокет будет уничтожен без отправки ожидаемого ответа. См. пример ниже:

const fs = require('node:fs');
const http = require('node:http');
const { pipeline } = require('node:stream');

const server = http.createServer((req, res) => {
  const fileStream = fs.createReadStream('./fileNotExist.txt');
  pipeline(fileStream, res, (err) => {
    if (err) {
      console.log(err); // No such file
      // this message can't be sent once `pipeline` already destroyed the socket
      return res.end('error!!!');
    }
  });
}); copy

stream.compose(...streams)

История
Версия Изменения
v21.1.0, v20.10.0

Добавлена поддержка класса stream.

v19.8.0, v18.16.0

Добавлена поддержка веб-потоков.

v16.9.0

Добавлено в версии v16.9.0

Стабильность: 1 — stream.compose является экспериментальным.
  • streams <Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> | <Duplex[]> | <Function>
  • Возвращает: <stream.Duplex>

Объединяет два или более потоков в поток Duplex, который записывает данные в первый поток и читает их из последнего. Каждый переданный поток соединяется со следующим с помощью stream.pipeline. При возникновении ошибки в любом из потоков уничтожаются все потоки, включая внешний поток Duplex.

Поскольку stream.compose возвращает новый поток, который, в свою очередь, можно (и следует) передавать в другие потоки, эта функция позволяет создавать композиции. В отличие от этого, при передаче потоков в stream.pipeline первый поток обычно является читаемым, а последний — записываемым, образуя замкнутую цепь.

Если передан Function, он должен быть фабричным методом, принимающим source Iterable.

import { compose, Transform } from 'node:stream';

const removeSpaces = new Transform({
  transform(chunk, encoding, callback) {
    callback(null, String(chunk).replace(' ', ''));
  },
});

async function* toUpper(source) {
  for await (const chunk of source) {
    yield String(chunk).toUpperCase();
  }
}

let res = '';
for await (const buf of compose(removeSpaces, toUpper).end('hello world')) {
  res += buf;
}

console.log(res); // prints 'HELLOWORLD' copy

stream.compose можно использовать для преобразования асинхронных итерируемых объектов, генераторов и функций в потоки.

  • AsyncIterable преобразует объект в читаемый Duplex. Не может выдавать null.
  • AsyncGeneratorFunction преобразует объект в читаемый/записываемый поток преобразования Duplex. Первым параметром должен принимать исходный AsyncIterable. Не может выдавать null.
  • AsyncFunction преобразует объект в записываемый Duplex. Должен возвращать либо null, либо undefined.
import { compose } from 'node:stream';
import { finished } from 'node:stream/promises';

// Convert AsyncIterable into readable Duplex.
const s1 = compose(async function*() {
  yield 'Hello';
  yield 'World';
}());

// Convert AsyncGenerator into transform Duplex.
const s2 = compose(async function*(source) {
  for await (const chunk of source) {
    yield String(chunk).toUpperCase();
  }
});

let res = '';

// Convert AsyncFunction into writable Duplex.
const s3 = compose(async function(source) {
  for await (const chunk of source) {
    res += chunk;
  }
});

await finished(compose(s1, s2, s3));

console.log(res); // prints 'HELLOWORLD' copy

См. readable.compose(stream), чтобы использовать stream.compose в качестве оператора.

stream.isErrored(stream)

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.3.0, v16.14.0

Добавлено в версиях v17.3.0 и v16.14.0

  • stream <Readable> | <Writable> | <Duplex> | <WritableStream> | <ReadableStream>
  • Возвращает: <boolean>

Возвращает значение, указывающее, произошла ли в потоке ошибка.

stream.isReadable(stream)

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.4.0, v16.14.0

Добавлено в версиях v17.4.0 и v16.14.0

  • stream <Readable> | <Duplex> | <ReadableStream>
  • Возвращает: <boolean> | <null> — возвращает null только в том случае, если stream не является допустимым Readable, Duplex или ReadableStream.

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

stream.isWritable(stream)

  • stream <Writable> | <Duplex> | <WritableStream>
  • Возвращает: <boolean> | <null> — возвращает null только в том случае, если stream не является допустимым Writable, Duplex или WritableStream.

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

stream.Readable.from(iterable[, options])

Добавлено в версиях v12.3.0 и v10.17.0
  • iterable <Iterable> Объект, реализующий протокол перебора Symbol.asyncIterator или Symbol.iterator. Генерирует событие 'error', если передано значение null.
  • options <Object> Параметры, передаваемые в new stream.Readable([options]). По умолчанию Readable.from() устанавливает options.objectMode в значение true, если это явно не отменено установкой options.objectMode в значение false.
  • Возвращает: <stream.Readable>

Вспомогательный метод для создания читаемых потоков из итераторов.

const { Readable } = require('node:stream');

async function * generate() {
  yield 'hello';
  yield 'streams';
}

const readable = Readable.from(generate());

readable.on('data', (chunk) => {
  console.log(chunk);
}); copy

В целях повышения производительности вызов Readable.from(string) или Readable.from(buffer) не приводит к перебору строк или буферов в соответствии с семантикой других потоков.

Если в качестве аргумента передан объект Iterable, содержащий промисы, это может привести к необработанному отклонению промиса.

const { Readable } = require('node:stream');

Readable.from([
  new Promise((resolve) => setTimeout(resolve('1'), 1500)),
  new Promise((_, reject) => setTimeout(reject(new Error('2')), 1000)), // Unhandled rejection
]); copy

stream.Readable.fromWeb(readableStream[, options])

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.0.0

Добавлено в версии v17.0.0

  • readableStream <ReadableStream>
  • options <Object>
    • encoding <string>
    • highWaterMark <number>
    • objectMode <boolean>
    • signal <AbortSignal>
  • Возвращает: <stream.Readable>

stream.Readable.isDisturbed(stream)

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v16.8.0

Добавлено в версии v16.8.0

  • stream <stream.Readable> | <ReadableStream>
  • Возвращает: boolean

Возвращает значение, указывающее, считывались ли данные из потока или был ли он отменён.

stream.Readable.toWeb(streamReadable[, options])

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v18.7.0

Добавлена поддержка параметров стратегии для Readable.

v17.0.0

Добавлено в версии v17.0.0

  • streamReadable <stream.Readable>
  • options <Object>
    • strategy <Object>
      • highWaterMark <number> Максимальный размер внутренней очереди (созданного ReadableStream), после которого при чтении из указанного stream.Readable начинает применяться обратное давление. Если значение не задано, оно будет взято из указанного stream.Readable.
      • size <Function> Функция, вычисляющая размер указанного фрагмента данных. Если значение не задано, размер всех фрагментов будет 1.
        • chunk <any>
        • Возвращает: <number>
  • Возвращает: <ReadableStream>

stream.Writable.fromWeb(writableStream[, options])

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.0.0

Добавлено в версии v17.0.0

  • writableStream <WritableStream>
  • options <Object>
    • decodeStrings <boolean>
    • highWaterMark <number>
    • objectMode <boolean>
    • signal <AbortSignal>
  • Возвращает: <stream.Writable>

stream.Writable.toWeb(streamWritable)

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.0.0

Добавлено в версии v17.0.0

  • streamWritable <stream.Writable>
  • Возвращает: <WritableStream>

stream.Duplex.from(src)

История
Версия Изменения
v19.5.0, v18.17.0

Аргумент src теперь может быть ReadableStream или WritableStream.

v16.8.0

Добавлено в: v16.8.0

  • src <Stream> | <Blob> | <ArrayBuffer> | <string> | <Iterable> | <AsyncIterable> | <AsyncGeneratorFunction> | <AsyncFunction> | <Promise> | <Object> | <ReadableStream> | <WritableStream>

Вспомогательный метод для создания двунаправленных потоков.

  • Stream преобразует поток для записи в записываемый Duplex, а поток для чтения — в Duplex.
  • Blob преобразует в читаемый Duplex.
  • string преобразует в читаемый Duplex.
  • ArrayBuffer преобразует в читаемый Duplex.
  • AsyncIterable преобразует в читаемый Duplex. Не может выдавать null.
  • AsyncGeneratorFunction преобразует в преобразующий Duplex для чтения и записи. Первым параметром должен принимать исходный AsyncIterable. Не может выдавать null.
  • AsyncFunction преобразует в записываемый Duplex. Должен возвращать либо null, либо undefined
  • Object ({ writable, readable }) преобразует readable и writable в Stream, а затем объединяет их в Duplex, где Duplex будет записывать в writable и читать из readable.
  • Promise преобразует в читаемый Duplex. Значение null игнорируется.
  • ReadableStream преобразует в читаемый Duplex.
  • WritableStream преобразует в записываемый Duplex.
  • Возвращает: <stream.Duplex>

Если в качестве аргумента передан объект Iterable, содержащий промисы, это может привести к необработанному отклонению промиса.

const { Duplex } = require('node:stream');

Duplex.from([
  new Promise((resolve) => setTimeout(resolve('1'), 1500)),
  new Promise((_, reject) => setTimeout(reject(new Error('2')), 1000)), // Unhandled rejection
]); copy

stream.Duplex.fromWeb(pair[, options])

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.0.0

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

  • pair <Object>
    • readable <ReadableStream>
    • writable <WritableStream>
  • options <Object>
    • allowHalfOpen <boolean>
    • decodeStrings <boolean>
    • encoding <string>
    • highWaterMark <number>
    • objectMode <boolean>
    • signal <AbortSignal>
  • Возвращает: <stream.Duplex>
Модули JavaScript
import { Duplex } from 'node:stream';
import {
  ReadableStream,
  WritableStream,
} from 'node:stream/web';

const readable = new ReadableStream({
  start(controller) {
    controller.enqueue('world');
  },
});

const writable = new WritableStream({
  write(chunk) {
    console.log('writable', chunk);
  },
});

const pair = {
  readable,
  writable,
};
const duplex = Duplex.fromWeb(pair, { encoding: 'utf8', objectMode: true });

duplex.write('hello');

for await (const chunk of duplex) {
  console.log('readable', chunk);
}
CommonJS
const { Duplex } = require('node:stream');
const {
  ReadableStream,
  WritableStream,
} = require('node:stream/web');

const readable = new ReadableStream({
  start(controller) {
    controller.enqueue('world');
  },
});

const writable = new WritableStream({
  write(chunk) {
    console.log('writable', chunk);
  },
});

const pair = {
  readable,
  writable,
};
const duplex = Duplex.fromWeb(pair, { encoding: 'utf8', objectMode: true });

duplex.write('hello');
duplex.once('readable', () => console.log('readable', duplex.read()));

stream.Duplex.toWeb(streamDuplex)

История
Версия Изменения
v22.17.0

API объявлен стабильным.

v17.0.0

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

  • streamDuplex <stream.Duplex>
  • Возвращает: <Object>
    • readable <ReadableStream>
    • writable <WritableStream>
Модули JavaScript
import { Duplex } from 'node:stream';

const duplex = Duplex({
  objectMode: true,
  read() {
    this.push('world');
    this.push(null);
  },
  write(chunk, encoding, callback) {
    console.log('writable', chunk);
    callback();
  },
});

const { readable, writable } = Duplex.toWeb(duplex);
writable.getWriter().write('hello');

const { value } = await readable.getReader().read();
console.log('readable', value);
CommonJS
const { Duplex } = require('node:stream');

const duplex = Duplex({
  objectMode: true,
  read() {
    this.push('world');
    this.push(null);
  },
  write(chunk, encoding, callback) {
    console.log('writable', chunk);
    callback();
  },
});

const { readable, writable } = Duplex.toWeb(duplex);
writable.getWriter().write('hello');

readable.getReader().read().then((result) => {
  console.log('readable', result.value);
});

stream.addAbortSignal(signal, stream)

История
Версия Изменения
v19.7.0, v18.16.0

Добавлена поддержка ReadableStream и WritableStream.

v15.4.0

Добавлено в: v15.4.0

  • signal <AbortSignal> Сигнал, обозначающий возможность отмены
  • stream <Stream> | <ReadableStream> | <WritableStream> Поток, к которому нужно привязать сигнал.

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

Вызов abort для AbortController, соответствующего переданному AbortSignal, будет работать так же, как вызов .destroy(new AbortError()) для потока, а для веб-потоков — controller.error(new AbortError()).

const fs = require('node:fs');

const controller = new AbortController();
const read = addAbortSignal(
  controller.signal,
  fs.createReadStream(('object.json')),
);
// Later, abort the operation closing the stream
controller.abort(); copy

Или используя AbortSignal с читаемым потоком в качестве асинхронного итерируемого объекта:

const controller = new AbortController();
setTimeout(() => controller.abort(), 10_000); // set a timeout
const stream = addAbortSignal(
  controller.signal,
  fs.createReadStream(('object.json')),
);
(async () => {
  try {
    for await (const chunk of stream) {
      await process(chunk);
    }
  } catch (e) {
    if (e.name === 'AbortError') {
      // The operation was cancelled
    } else {
      throw e;
    }
  }
})(); copy

Или используя AbortSignal с ReadableStream:

const controller = new AbortController();
const rs = new ReadableStream({
  start(controller) {
    controller.enqueue('hello');
    controller.enqueue('world');
    controller.close();
  },
});

addAbortSignal(controller.signal, rs);

finished(rs, (err) => {
  if (err) {
    if (err.name === 'AbortError') {
      // The operation was cancelled
    }
  }
});

const reader = rs.getReader();

reader.read().then(({ value, done }) => {
  console.log(value); // hello
  console.log(done); // false
  controller.abort();
}); copy

stream.getDefaultHighWaterMark(objectMode)

Добавлено в: v19.9.0, v18.17.0
  • objectMode <boolean>
  • Возвращает: <integer>

Возвращает значение highWaterMark по умолчанию, используемое потоками. По умолчанию это 65536 (64 КиБ) или 16 для objectMode.

stream.setDefaultHighWaterMark(objectMode, value)

Добавлено в: v19.9.0, v18.17.0
  • objectMode <boolean>
  • value <integer> Значение highWaterMark

Устанавливает значение highWaterMark по умолчанию, используемое потоками.

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

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

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

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

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

При расширении потоков учитывайте, какие параметры пользователь может и должен передавать, прежде чем передавать их базовому конструктору. Например, если реализация предполагает определённые значения параметров 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('node:stream');

const myWritable = new Writable({
  construct(callback) {
    // Initialize state and load resources...
  },
  write(chunk, encoding, callback) {
    // ...
  },
  destroy() {
    // Free resources...
  },
}); copy

Реализация записываемого потока

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

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

new stream.Writable([options])
История
Версия Изменения
v22.0.0

Увеличено значение highWaterMark по умолчанию.

v15.5.0

Добавлена поддержка передачи AbortSignal.

v14.0.0

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

v11.2.0, v10.16.0

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

v10.0.0

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

  • options <Object>
    • highWaterMark <number> Размер буфера, при котором stream.write() начинает возвращать false. По умолчанию: 65536 (64 КиБ) или 16 для потоков objectMode.
    • decodeStrings <boolean> Определяет, следует ли кодировать переданные в stream.write() значения string в значения Buffer (с кодировкой, указанной при вызове stream.write()) перед передачей в stream._write(). Другие типы данных не преобразуются (то есть значения Buffer не декодируются в значения string). Значение false предотвращает преобразование значений string. По умолчанию: true.
    • defaultEncoding <string> Кодировка по умолчанию, используемая, если при вызове stream.write() не указана кодировка. По умолчанию: 'utf8'.
    • objectMode <boolean> Определяет, является ли stream.write(anyObj) допустимой операцией. Если параметр задан, становится возможным записывать значения JavaScript, отличные от строк, <Buffer>, <TypedArray> или <DataView>, если это поддерживается реализацией потока. По умолчанию: false.
    • emitClose <boolean> Определяет, должен ли поток генерировать 'close' после уничтожения. По умолчанию: true.
    • write <Function> Реализация метода stream._write().
    • writev <Function> Реализация метода stream._writev().
    • destroy <Function> Реализация метода stream._destroy().
    • final <Function> Реализация метода stream._final().
    • construct <Function> Реализация метода stream._construct().
    • autoDestroy <boolean> Определяет, должен ли поток автоматически вызывать .destroy() после завершения. По умолчанию: true.
    • signal <AbortSignal> Сигнал, обозначающий возможность отмены.
const { Writable } = require('node:stream');

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

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

const { Writable } = require('node:stream');
const util = require('node:util');

function MyWritable(options) {
  if (!(this instanceof MyWritable))
    return new MyWritable(options);
  Writable.call(this, options);
}
util.inherits(MyWritable, Writable); copy

Или с использованием упрощённого конструктора:

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

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

Вызов abort для AbortController, соответствующего переданному AbortSignal, будет работать так же, как вызов .destroy(new AbortError()) для записываемого потока.

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

const controller = new AbortController();
const myWritable = new Writable({
  write(chunk, encoding, callback) {
    // ...
  },
  writev(chunks, callback) {
    // ...
  },
  signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort(); copy
writable._construct(callback)
Добавлено в: v15.0.0
  • callback <Function> Вызовите эту функцию (необязательно с аргументом ошибки), когда инициализация потока завершится.

Метод _construct() НЕЛЬЗЯ вызывать напрямую. Его могут реализовать дочерние классы; в таком случае он будет вызываться только внутренними методами класса Writable.

Эта необязательная функция будет вызвана в следующем тике после возврата конструктора потока, откладывая вызовы _write(), _final() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной подготовки ресурсов до начала использования потока.

const { Writable } = require('node:stream');
const fs = require('node:fs');

class WriteStream extends Writable {
  constructor(filename) {
    super();
    this.filename = filename;
    this.fd = null;
  }
  _construct(callback) {
    fs.open(this.filename, 'w', (err, fd) => {
      if (err) {
        callback(err);
      } else {
        this.fd = fd;
        callback();
      }
    });
  }
  _write(chunk, encoding, callback) {
    fs.write(this.fd, chunk, callback);
  }
  _destroy(err, callback) {
    if (this.fd) {
      fs.close(this.fd, (er) => callback(er || err));
    } else {
      callback(err);
    }
  }
} copy
writable._write(chunk, encoding, callback)
История
Версия Изменения
v12.11.0

_write() становится необязательным при передаче _writev().

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

Все реализации потоков 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 <Object[]> Записываемые данные. Это массив объектов <Object>, каждый из которых представляет отдельный фрагмент данных для записи. Свойства этих объектов:
    • chunk <Buffer> | <string> Экземпляр буфера или строка, содержащая записываемые данные. Значение chunk будет строкой, если Writable создан с параметром decodeStrings, равным false, и в write() передана строка.
    • encoding <string> Кодировка символов chunk. Если chunk является Buffer, значение encoding будет 'buffer'.
  • callback <Function> Функция обратного вызова (необязательно с аргументом ошибки), которая будет вызвана после завершения обработки переданных фрагментов.

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

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

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

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

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

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

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

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

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

Об ошибках, возникающих при выполнении методов writable._write(), writable._writev() и writable._final(), необходимо сообщать, вызывая функцию обратного вызова и передавая ошибку в качестве первого аргумента. Выбрасывание Error внутри этих методов или ручная генерация события 'error' приводят к неопределённому поведению.

Если поток Readable передаёт данные в поток Writable, то при возникновении ошибки в Writable поток Readable будет отключён от конвейера.

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

const myWritable = new Writable({
  write(chunk, encoding, callback) {
    if (chunk.toString().indexOf('a') >= 0) {
      callback(new Error('chunk is invalid'));
    } else {
      callback();
    }
  },
}); copy
Пример записываемого потока

Ниже показана довольно простая (и в некоторой степени бесполезная) реализация пользовательского потока Writable. Хотя этот конкретный экземпляр потока Writable не имеет особой практической пользы, пример демонстрирует все необходимые элементы пользовательского экземпляра потока Writable:

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

class MyWritable extends Writable {
  _write(chunk, encoding, callback) {
    if (chunk.toString().indexOf('a') >= 0) {
      callback(new Error('chunk is invalid'));
    } else {
      callback();
    }
  }
} copy
Декодирование буферов в записываемом потоке

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

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

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

Реализация потока для чтения

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

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

new stream.Readable([options])
История
Версия Изменения
v22.0.0

увеличено значение highWaterMark по умолчанию.

v15.5.0

добавлена поддержка передачи AbortSignal.

v14.0.0

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

v11.2.0, v10.16.0

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

  • options <Object>
    • highWaterMark <number> Максимальное количество байтов, которое можно хранить во внутреннем буфере до прекращения чтения из нижележащего ресурса. По умолчанию: 65536 (64 КиБ) или 16 для потоков objectMode.
    • encoding <string> Если задано, буферы будут декодироваться в строки с использованием указанной кодировки. По умолчанию: null.
    • objectMode <boolean> Должен ли этот поток работать как поток объектов. Это означает, что stream.read(n) возвращает одно значение, а не Buffer размером n. По умолчанию: false.
    • emitClose <boolean> Должен ли поток генерировать событие 'close' после уничтожения. По умолчанию: true.
    • read <Function> Реализация метода stream._read().
    • destroy <Function> Реализация метода stream._destroy().
    • construct <Function> Реализация метода stream._construct().
    • autoDestroy <boolean> Должен ли этот поток автоматически вызывать .destroy() для себя после завершения. По умолчанию: true.
    • signal <AbortSignal> Сигнал, представляющий возможную отмену.
const { Readable } = require('node:stream');

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

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

const { Readable } = require('node:stream');
const util = require('node:util');

function MyReadable(options) {
  if (!(this instanceof MyReadable))
    return new MyReadable(options);
  Readable.call(this, options);
}
util.inherits(MyReadable, Readable); copy

Или при использовании упрощенного подхода к созданию конструктора:

const { Readable } = require('node:stream');

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

Вызов abort для AbortController, соответствующего переданному AbortSignal, будет вести себя так же, как вызов .destroy(new AbortError()) для созданного потока для чтения.

const { Readable } = require('node:stream');
const controller = new AbortController();
const read = new Readable({
  read(size) {
    // ...
  },
  signal: controller.signal,
});
// Later, abort the operation closing the stream
controller.abort(); copy
readable._construct(callback)
Добавлено в: v15.0.0
  • callback <Function> Вызовите эту функцию (при необходимости с аргументом ошибки), когда инициализация потока завершена.

Метод _construct() НЕ ДОЛЖЕН вызываться напрямую. Он может быть реализован дочерними классами; в этом случае его будут вызывать только внутренние методы класса Readable.

Эта необязательная функция будет запланирована конструктором потока на следующий такт, откладывая вызовы _read() и _destroy() до вызова callback. Это полезно для инициализации состояния или асинхронной инициализации ресурсов до того, как поток можно будет использовать.

const { Readable } = require('node:stream');
const fs = require('node:fs');

class ReadStream extends Readable {
  constructor(filename) {
    super();
    this.filename = filename;
    this.fd = null;
  }
  _construct(callback) {
    fs.open(this.filename, (err, fd) => {
      if (err) {
        callback(err);
      } else {
        this.fd = fd;
        callback();
      }
    });
  }
  _read(n) {
    const buf = Buffer.alloc(n);
    fs.read(this.fd, buf, 0, n, null, (err, bytesRead) => {
      if (err) {
        this.destroy(err);
      } else {
        this.push(bytesRead > 0 ? buf.slice(0, bytesRead) : null);
      }
    });
  }
  _destroy(err, callback) {
    if (this.fd) {
      fs.close(this.fd, (er) => callback(er || err));
    } else {
      callback(err);
    }
  }
} copy
readable._read(size)
Добавлено в: v0.9.4
  • size <number> Количество байтов для асинхронного чтения

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

Все реализации потока Readable должны предоставлять реализацию метода readable._read() для получения данных из нижележащего ресурса.

При вызове readable._read(), если из ресурса доступны данные, реализация должна начать помещать их в очередь чтения с помощью метода this.push(dataChunk). _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 <Error> Возможная ошибка.
  • callback <Function> Функция обратного вызова, принимающая необязательный аргумент ошибки.

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

readable.push(chunk[, encoding])
История
Версия Изменения
v22.0.0

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

v8.0.0

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

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

Если chunk имеет тип <Buffer>, <TypedArray>, <DataView> или <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();
  }
} copy

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

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

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

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

const { Readable } = require('node:stream');

const myReadable = new Readable({
  read(size) {
    const err = checkSomeErrorCondition();
    if (err) {
      this.destroy(err);
    } else {
      // Do some work.
    }
  },
}); copy
Пример потока-счетчика

Ниже приведен простой пример потока Readable, который последовательно выдает числа от 1 до 1 000 000, а затем завершается.

const { Readable } = require('node: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);
    }
  }
} copy

Реализация дуплексного потока

Duplex — это поток, реализующий одновременно Readable и Writable, например TCP-соединение через сокет.

Поскольку JavaScript не поддерживает множественное наследование, для реализации потока Duplex расширяется класс stream.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 <Object> Передается конструкторам Writable и Readable. Также содержит следующие поля:
    • allowHalfOpen <boolean> Если задано значение false, сторона для записи потока автоматически завершится при завершении стороны для чтения. По умолчанию: true.
    • readable <boolean> Определяет, должна ли сторона Duplex быть доступна для чтения. По умолчанию: true.
    • writable <boolean> Определяет, должна ли сторона Duplex быть доступна для записи. По умолчанию: true.
    • readableObjectMode <boolean> Задает objectMode для стороны потока, доступной для чтения. Не влияет на поток, если objectMode имеет значение true. По умолчанию: false.
    • writableObjectMode <boolean> Задает objectMode для стороны потока, доступной для записи. Не влияет на поток, если objectMode имеет значение true. По умолчанию: false.
    • readableHighWaterMark <number> Задает highWaterMark для стороны потока, доступной для чтения. Не влияет на поток, если указан highWaterMark.
    • writableHighWaterMark <number> Задает highWaterMark для стороны потока, доступной для записи. Не влияет на поток, если указан highWaterMark.
const { Duplex } = require('node:stream');

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

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

const { Duplex } = require('node:stream');
const util = require('node:util');

function MyDuplex(options) {
  if (!(this instanceof MyDuplex))
    return new MyDuplex(options);
  Duplex.call(this, options);
}
util.inherits(MyDuplex, Duplex); copy

Или при использовании упрощенного подхода к созданию конструктора:

const { Duplex } = require('node:stream');

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

При использовании конвейера:

const { Transform, pipeline } = require('node:stream');
const fs = require('node:fs');

pipeline(
  fs.createReadStream('object.json')
    .setEncoding('utf8'),
  new Transform({
    decodeStrings: false, // Accept string input rather than Buffers
    construct(callback) {
      this.data = '';
      callback();
    },
    transform(chunk, encoding, callback) {
      this.data += chunk;
      callback();
    },
    flush(callback) {
      try {
        // Make sure is valid json.
        JSON.parse(this.data);
        this.push(this.data);
        callback();
      } catch (err) {
        callback(err);
      }
    },
  }),
  fs.createWriteStream('valid-object.json'),
  (err) => {
    if (err) {
      console.error('failed', err);
    } else {
      console.log('completed');
    }
  },
); copy
Пример дуплексного потока

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

const { Duplex } = require('node: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));
    });
  }
} copy

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

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

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

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

const { Transform } = require('node: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 copy

Реализация преобразующего потока

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

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

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

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

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

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

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

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

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

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

Или при использовании упрощенного подхода к созданию конструктора:

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

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

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

Событие: 'finish'

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

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

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

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

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

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

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

transform._transform(chunk, encoding, callback)
  • chunk <Buffer> | <string> | <any> Фрагмент Buffer для преобразования, полученный из string, переданного в stream.write(). Если параметр потока decodeStrings имеет значение false или поток работает в режиме объектов, фрагмент не преобразуется и будет иметь то же значение, что и переданное в stream.write().
  • encoding <string> Если фрагмент является строкой, это тип кодировки. Если фрагмент является буфером, это специальное значение 'buffer'. В этом случае игнорируйте его.
  • callback <Function> Функция обратного вызова (при необходимости с аргументом ошибки и данными), вызываемая после обработки переданного 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);
}; copy

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

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

Создание потоков для чтения с помощью асинхронных генераторов

Поток для чтения Node.js можно создать из асинхронного генератора с помощью вспомогательного метода Readable.from():

const { Readable } = require('node:stream');

const ac = new AbortController();
const signal = ac.signal;

async function * generate() {
  yield 'a';
  await someLongRunningFn({ signal });
  yield 'b';
  yield 'c';
}

const readable = Readable.from(generate());
readable.on('close', () => {
  ac.abort();
});

readable.on('data', (chunk) => {
  console.log(chunk);
}); copy
Передача данных в потоки для записи из асинхронных итераторов

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

const fs = require('node:fs');
const { pipeline } = require('node:stream');
const { pipeline: pipelinePromise } = require('node:stream/promises');

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

const ac = new AbortController();
const signal = ac.signal;

const iterator = createIterator({ signal });

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

// Promise Pattern
pipelinePromise(iterator, writable)
  .then((value) => {
    console.log(value, 'value returned');
  })
  .catch((err) => {
    console.error(err);
    ac.abort();
  }); copy

Совместимость со старыми версиями 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); copy

До 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); copy

Помимо переключения новых потоков 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('') не рекомендуется.

Передача строки <string>, объекта <Buffer>, объекта <TypedArray> или объекта <DataView> нулевой длины в поток, работающий не в объектном режиме, имеет интересный побочный эффект. Поскольку это вызов 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-v22.x/docs/api/stream.html

Spec-Zone.ru

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