Spec-Zone.ru › Node.js 24 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> Завершать поток назначения, когда завершается поток источника. Потоки Transform всегда завершаются, даже если это значение равно 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
  • Поток stdin дочернего процесса
  • 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> исходный поток, который прекратил передачу данных в этот поток для записи методом unpiped

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

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

const writer = getWritableStreamSomehow();
const reader = getReadableStreamSomehow();
writer.on('unpipe', (src) => {
  console.log('Something has stopped piping into the writer.');
  assert.equal(src, reader);
});
reader.pipe(writer);
reader.unpipe(writer); 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, v20.13.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 <string> | <Buffer> | <TypedArray> | <DataView> | <any> Необязательные данные для записи. Для потоков, не работающих в режиме объектов, chunk должно быть значением типа <string>, <Buffer>, <TypedArray> или <DataView>. Для потоков в режиме объектов chunk может быть любым значением JavaScript, кроме null.
  • encoding <string> Кодировка, если chunk является строкой
  • 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
История
Версия Изменения
v24.0.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]()
История
Версия Изменения
v24.2.0

Больше не является экспериментальным.

v22.4.0, v20.16.0

Добавлено в: v22.4.0, v20.16.0

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

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

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

v8.0.0

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

v6.0.0

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

v0.9.4

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

  • chunk <string> | <Buffer> | <TypedArray> | <DataView> | <any> Необязательные данные для записи. Для потоков, не работающих в режиме объектов, chunk должно быть значением типа <string>, <Buffer>, <TypedArray> или <DataView>. Для потоков в режиме объектов chunk может быть любым значением JavaScript, кроме null.
  • encoding <string> | <null> Кодировка, если chunk является строкой. По умолчанию: 'utf8'
  • 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>

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

Функции обратного вызова-обработчика будет передан один объект 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
История
Версия Изменения
v24.0.0

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

v16.8.0

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

  • Тип: <boolean>

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

readable.readableDidRead
История
Версия Изменения
v24.0.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, v20.13.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 КиБ, поскольку параметр highWaterMark не передан в fs.createReadStream().

readable[Symbol.asyncDispose]()
История
Версия Изменения
v24.2.0

Больше не является экспериментальным.

v20.4.0, v18.18.0

Добавлено в: v20.4.0, v18.18.0

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

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

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

v19.1.0, v18.13.0

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

  • stream <Writable> | <Duplex> | <WritableStream> | <TransformStream> | <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(['text passed through', 'composed stream']).compose(splitToWords);
const words = await wordsStream.toArray();

console.log(words); // prints ['text', 'passed', 'through', 'composed', 'stream'] copy

readable.compose(s) эквивалентен stream.compose(readable, s).

Этот метод также позволяет передать <AbortSignal>, который уничтожит объединённый поток при прерывании.

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

readable.iterator([options])
История
Версия Изменения
v24.0.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]

const resolver = new Resolver();

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

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

Из 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, v20.17.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) доступен в потоках <Readable> и <Duplex> как оболочка для этой функции.

stream.isErrored(stream)

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

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

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

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

API признан стабильным.

v16.8.0

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

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

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

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

История
Версия Изменения
v24.14.0

Добавлен параметр 'type' для указания значения 'bytes'.

v24.0.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>
    • type <string> Должно быть 'bytes' или undefined.
  • Возвращает: <ReadableStream>

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

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

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

История
Версия Изменения
v24.0.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[, options])

История
Версия Изменения
v24.14.0

Добавлена опция 'type' для указания значения 'bytes'.

v24.0.0

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

v17.0.0

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

  • streamDuplex <stream.Duplex>
  • options <Object>
    • type <string> Должно быть 'bytes' или undefined.
  • Возвращает: <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

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

Для реализации потока Writable наследуют класс stream.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> Следует ли кодировать string, переданные в stream.write(), в 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, v20.13.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-v24.x/docs/api/stream.html

Spec-Zone.ru

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