Spec-Zone.ru › Node.js 18 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: потоки, которые могут изменять или преобразовывать данные при записи и чтении (например, zlib.createDeflate()).

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

API потоков с обещаниями

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

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

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

stream.pipeline(streams[, options])

Добавлена в: 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>
  • Возвращает: <Promise> Выполняется при завершении конвейера.

Модули CJS

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

Модули MJS

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, передайте его внутри объекта options в качестве последнего аргумента. Когда сигнал прерывается, destroy будет вызван в базовом конвейере с AbortError.

Модули CJS

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

Модули MJS

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 потоков с обещаниями также поддерживает асинхронные генераторы:

Модули CJS

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

Модули MJS

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

Модули CJS

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

Модули MJS

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 потоков с обещаниями предоставляет версию обратного вызова:

stream.finished(stream[, options])

Добавлена в: v15.0.0
  • stream <Stream>
  • options <Object>
    • error <boolean> | <undefined>
    • readable <boolean> | <undefined>
    • writable <boolean> | <undefined>
    • signal: <AbortSignal> | <undefined>
  • Возвращает: <Promise> Выполняется, когда поток больше не читаем или не записывается.

Модули CJS

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.

Модули MJS

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 потоков с обещаниями предоставляет версию обратного вызова:

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

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

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

Буферизация

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

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

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

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

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

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

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

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

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

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

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

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

const server = http.createServer((req, res) => {
  // `req` is an http.IncomingMessage, which is a readable stream.
  // `res` is an http.ServerResponse, which is a writable stream.

  let body = '';
  // Get the data as utf8 strings.
  // If an encoding is not set, Buffer objects will be received.
  req.setEncoding('utf8');

  // Readable streams emit 'data' events once a listener is added.
  req.on('data', (chunk) => {
    body += chunk;
  });

  // The 'end' event indicates that the entire body has been received.
  req.on('end', () => {
    try {
      const data = JSON.parse(body);
      // Write back something interesting to the user:
      res.write(typeof data);
      res.end();
    } catch (er) {
      // uh oh! bad json!
      res.statusCode = 400;
      return res.end(`error: ${er.message}`);
    }
  });
});

server.listen(1337);

// $ curl localhost:1337 -d "{}"
// object
// $ curl localhost:1337 -d "\"foo\""
// string
// $ curl localhost:1337 -d "not json"
// error: Unexpected token 'o', "not json" is not valid JSON copy

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

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

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

Duplex и Transform потоки являются как Writable, так и Readable.

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

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

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

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

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

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

Некоторые из этих примеров фактически являются Duplex потоками, которые реализуют интерфейс Writable.

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

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

const myStream = getWritableStreamSomehow();
myStream.write('some data');
myStream.write('some more data');
myStream.end('done writing data'); copy
Класс: stream.Writable
Добавлен в: v0.9.4
Событие: 'close'
История
Версия Изменения
v10.0.0

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

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

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

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

Событие 'finish' генерируется после вызова метода stream.end() и после того, как все данные были переданы в подлежащую систему.

const writer = getWritableStreamSomehow();
for (let i = 0; i < 100; i++) {
  writer.write(`hello, #${i}!\n`);
}
writer.on('finish', () => {
  console.log('All writes are now complete.');
});
writer.end('This is the end\n'); copy
Событие: 'pipe'
Добавлен в: v0.9.4
  • src Источник-поток <stream.Readable>, который передает данные в этот поток записи

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

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

Событие 'unpipe' генерируется при вызове метода 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'.
  • Возвращает: <this>

Уничтожить поток. Необязательно сгенерировать событие 'error' и событие 'close' (если параметр emitClose не установлен в false). После этого вызова поток записи завершён, и последующие вызовы write() или end() приведут к ошибке ERR_STREAM_DESTROYED. Это разрушительный и мгновенный способ уничтожения потока. Предыдущие вызовы write() могут не быть завершены и могут вызвать ошибку ERR_STREAM_DESTROYED. Используйте end() вместо destroy, если данные должны быть переданы перед закрытием, или дождитесь события '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])
История
Версия Изменения
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 <строка> | <Буфер> | <Uint8Array> | <любой> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов, chunk должен быть строкой, Buffer или Uint8Array. Для потоков в режиме объектов, chunk может быть любым значением JavaScript, кроме null.
  • encoding <строка> Кодировка, если chunk — строка
  • callback <Функция> Обработчик события завершения потока.
  • Возвращает: <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 <строка> Новая кодировка по умолчанию
  • Возвращает: <this>

Метод writable.setDefaultEncoding() устанавливает кодировку по умолчанию для потока 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
  • <логическое>

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

writable.writableAborted
Добавлен в: v18.0.0
Устойчивость: 1 - Экспериментально
  • <логическое>

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

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

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

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

Число вызовов writable.uncork(), необходимых для полного разбуферизации потока.

writable.errored
Добавлен в: v18.0.0
  • <Ошибка>

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

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

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

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

Возвращает значение highWaterMark установленное при создании Writable.

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

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

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

Является true если буфер потока полон и поток испустит событие 'drain'.

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

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

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

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

v6.0.0

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

v0.9.4

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

const { PassThrough, Writable } = require('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. Для потоков, работающих в режиме объектов, кусок может быть любым значением JavaScript, кроме null.

Событие 'data' испускается всякий раз, когда поток отказывается от владения куском данных потребителю. Это может происходить всякий раз, когда поток переключается в режим потока путём вызова readable.pipe(), readable.resume() или добавления обратного вызова слушателя к событию 'data'. Событие 'data' также будет испускаться всякий раз, когда вызывается метод readable.read() и кусок данных доступен для возврата.

Добавление слушателя события 'data' к потоку, который не был явно приостановлен, переключит поток в режим потока. Данные будут передаваться как только они станут доступны.

Обратный вызов слушателя будет получать кусок данных как строку, если для потока был задан кодирование по умолчанию при помощи метода readable.setEncoding(); в противном случае данные будут переданы как Buffer.

const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
  console.log(`Received ${chunk.length} bytes of data.`);
}); copy
Событие: 'end'
Добавлен в: v0.9.4

Событие 'end' испускается, когда больше нет данных для потребления из потока.

Событие 'end' не будет излучаться, пока данные не будут полностью потреблены. Это можно сделать, переключив поток в режим потоковой передачи или вызвав stream.read() повторно, пока все данные не будут потреблены.

const readable = getReadableStreamSomehow();
readable.on('data', (chunk) => {
  console.log(`Received ${chunk.length} bytes of data.`);
});
readable.on('end', () => {
  console.log('There will be no more data.');
}); copy
Событие: 'error'
Добавлен в: v0.9.4
  • <Error>

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

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

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

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

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

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

v10.0.0

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

v0.9.4

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

Событие 'readable' излучается, когда доступны данные для чтения из потока или когда достигнут конец потока. По сути, событие 'readable' указывает, что в потоке есть новая информация. Если данные доступны, stream.read() вернет эти данные.

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

Работает как no-op в потоке, который уже уничтожен.

v8.0.0

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

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

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

После вызова destroy() любые последующие вызовы будут no-op, и больше никаких ошибок, кроме ошибок от _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. Поток данных будет управляться автоматически, чтобы потоковый Writable поток не перегружался более быстрым потоком Readable.

Следующий пример передаёт все данные из 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

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

Метод 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

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

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

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

Потоки process.stderr и process.stdout Writable никогда не закрываются до выхода процесса 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 GiB.

Метод 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. Фрагменты не конкатенируются. Для потребления всех данных, находящихся в буфере, необходим цикл 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
Добавлена в: v16.8.0
Устойчивость: 1 - Экспериментальная
  • <boolean>

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

readable.readableDidRead
Добавлена в: v16.7.0, v14.18.0
Устойчивость: 1 - Экспериментальная
  • <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() заставляет явным образом приостановленный поток Readable возобновить испускание событий 'data', переключая поток в режим потоковой передачи.

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

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

v0.9.11

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

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

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

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

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

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

// Pull off a header delimited by \n\n.
// Use unshift() if we get too much.
// Call the callback with (error, header, stream).
const { StringDecoder } = require('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() в пользовательском потоке). Вызов readable.unshift() с последующим немедленным вызовом stream.push('') сбросит состояние чтения должным образом, однако лучше просто избегать вызова readable.unshift() во время чтения.

readable.wrap(stream)
Добавлена в: v0.9.4
  • stream <Поток> Поток «старого стиля» для чтения
  • Возвращает: <это>

До 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]()
Добавлена в: v18.18.0
Устойчивость: 1 - Экспериментальная

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

readable.compose(stream[, options])
Добавлена в: v18.13.0
Устойчивость: 1 - Экспериментальная
  • stream <Поток> | <Итерируемый> | <AsyncIterable> | <Функция>
  • options <Объект>
    • signal <AbortSignal> позволяет уничтожить поток, если сигнал прерван.
  • Возвращает: <Duplex> составленный поток stream.
import { Readable } from 'node:stream';

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

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

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

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

См. stream.compose для получения дополнительной информации.

readable.iterator([options])
Добавлена в: v16.3.0
Устойчивость: 1 - Экспериментальная
  • options <Объект>
    • destroyOnReturn <boolean> Если установлено в false, вызов return на async итераторе или выход из итерации 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])
Добавлена в: v17.4.0, v16.14.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncFunction> функция для отображения каждого фрагмента в потоке.
    • data <любой> фрагмент данных из потока.
    • options <Объект>
      • signal <AbortSignal> прерывается, если поток уничтожен, что позволяет прервать вызов fn рано.
  • options <Объект>
    • concurrency <число> максимальное одновременное обращение к fn для вызова потока. По умолчанию: 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])
Добавлена в: v17.4.0, v16.14.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncFunction> функция для фильтрации кусков из потока.
    • data <любой> кусок данных из потока.
    • options <Объект>
      • signal <AbortSignal> прерывается, если поток уничтожен, позволяя прервать fn вызов досрочно.
  • options <Объект>
    • concurrency <число> максимальное одновременное обращение к fn для вызова в потоке за раз. По умолчанию: 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
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncFunction> функция для вызова на каждом куске потока.
    • data <любой> кусок данных из потока.
    • options <Объект>
      • signal <AbortSignal> прерывается, если поток уничтожен, позволяя прервать fn вызов досрочно.
  • options <Объект>
    • concurrency <число> максимальное одновременное обращение к fn для вызова в потоке за раз. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток, если сигнал прерван.
  • Возвращает: <Промис> промис, обозначающий завершение потока.

Этот метод позволяет итерировать поток. Для каждого куска в потоке вызывается функция 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
Устойчивость: 1 - Экспериментальная
  • options <Объект>
    • signal <AbortSignal> позволяет отменить операцию toArray, если сигнал прерван.
  • Возвращает: <Промис> промис, содержащий массив с содержимым потока.

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

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

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

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

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

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

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

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

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

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

import { Readable } from 'node:stream';

await Readable.from([1, 2, 3, 4]).take(2).toArray(); // [1, 2] copy
readable.asIndexedPairs([options])
История
Версия Изменения
v18.17.0

Использование метода asIndexedPairs генерирует предупреждение во время выполнения, что оно будет удалено в будущей версии.

v17.5.0

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

Устойчивость: 1 - Экспериментальная
  • options <Объект>
    • signal <AbortSignal> позволяет разрушить поток, если сигнал прерван.
  • Возвращает: <Readable> поток индексированных пар.

Этот метод возвращает новый поток с фрагментами базового потока, спаренными со счётчиком в форме [index, chunk]. Первое значение индекса равно 0, и оно увеличивается на 1 для каждого произведённого фрагмента.

import { Readable } from 'node:stream';

const pairs = await Readable.from(['a', 'b', 'c']).asIndexedPairs().toArray();
console.log(pairs); // [[0, 'a'], [1, 'b'], [2, 'c']] copy
readable.reduce(fn[, initial[, options]])
Добавлена в: v17.5.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncFunction> функция редуктора для вызова над каждым фрагментом в потоке.
    • previous <любой> значение, полученное от последнего вызова fn или значение initial если указано, или первый фрагмент потока в противном случае.
    • data <любой> фрагмент данных из потока.
    • options <Объект>
      • signal <AbortSignal> прерывается, если поток разрушен, позволяя прервать вызов fn раньше.
  • initial <любой> начальное значение для использования в редукции.
  • options <Объект>
    • signal <AbortSignal> позволяет разрушить поток, если сигнал прерван.
  • Возвращает: <Обещание> обещание для конечного значения редукции.

Этот метод вызывает 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.

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

  • сокеты TCP
  • потоки zlib
  • потоки crypto
duplex.allowHalfOpen
Добавлен в: v0.9.4
  • <boolean>

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

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

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

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

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

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

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

v8.0.0

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

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

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

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

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

История
Версия Изменения
v18.14.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 <Поток> | <ReadableStream> | <WritableStream>

Поток чтения и/или записи/веб-поток.

  • options <Объект>

    • error <boolean> Если установлено значение false, то вызов emit('error', err) не обрабатывается как завершенный. По умолчанию: true.
    • readable <boolean> При установке в false, обратный вызов будет вызван при завершении потока, даже если поток по-прежнему читается. По умолчанию: true.
    • writable <boolean> При установке в false, обратный вызов будет вызван при завершении потока, даже если поток по-прежнему записывается. По умолчанию: true.
    • signal <AbortSignal> позволяет прервать ожидание завершения потока. Подлежащий поток не будет прерван, если сигнал прерван. Обратный вызов будет вызван с AbortError. Все зарегистрированные слушатели, добавленные этой функцией, также будут удалены.
    • cleanup <boolean> Удалить все зарегистрированные слушатели потока. По умолчанию: false.
  • callback <Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.

  • Возвращает: <Функция> Функция очистки, которая удаляет все зарегистрированные слушатели.

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

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 предоставляет вариант с promise.

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)

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

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

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

v16.9.0

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

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

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

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

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

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

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

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

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

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

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

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

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

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

let res = '';

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

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

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

См. readable.compose(stream) для stream.compose как оператора.

stream.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])

Добавлена в: v17.0.0
Устойчивость: 1 - Экспериментально
  • readableStream <ReadableStream>
  • options <Object>
    • encoding <string>
    • highWaterMark <number>
    • objectMode <boolean>
    • signal <AbortSignal>
  • Возвращает: <stream.Readable>

stream.Readable.isDisturbed(stream)

Добавлена в: v16.8.0
Устойчивость: 1 - Экспериментально
  • stream <stream.Readable> | <ReadableStream>
  • Возвращает: boolean

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

stream.isErrored(stream)

Добавлена в: v17.3.0, v16.14.0
Устойчивость: 1 - Экспериментально
  • stream <Readable> | <Writable> | <Duplex> | <WritableStream> | <ReadableStream>
  • Возвращает: <boolean>

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

stream.isReadable(stream)

Добавлена в: v17.4.0, v16.14.0
Устойчивость: 1 - Экспериментальная
  • stream <Readable> | <Duplex> | <ReadableStream>
  • Возвращает: <boolean>

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

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

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

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

Добавлена в: v17.0.0
Устойчивость: 1 - Экспериментальная
  • writableStream <WritableStream>
  • options <Object>
    • decodeStrings <boolean>
    • highWaterMark <number>
    • objectMode <boolean>
    • signal <AbortSignal>
  • Возвращает: <stream.Writable>

stream.Writable.toWeb(streamWritable)

Добавлена в: v17.0.0
Устойчивость: 1 - Экспериментальная
  • streamWritable <stream.Writable>
  • Возвращает: <WritableStream>

stream.Duplex.from(src)

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

Утилитарный метод для создания потоков типа duplex.

  • 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])

Добавлена в: v17.0.0
Устойчивость: 1 - Экспериментальная
  • pair <Object>
    • readable <ReadableStream>
    • writable <WritableStream>
  • options <Object>
    • allowHalfOpen <boolean>
    • decodeStrings <boolean>
    • encoding <string>
    • highWaterMark <number>
    • objectMode <boolean>
    • signal <AbortSignal>
  • Возвращает: <stream.Duplex>

Модули MJS

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

Модули CJS

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)

Добавлена в: v17.0.0
Устойчивость: 1 - Экспериментальная
  • streamDuplex <stream.Duplex>
  • Возвращает: <Объект>
    • readable <ReadableStream>
    • writable <WritableStream>

Модули MJS

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

Модули CJS

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)

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

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

v15.4.0

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

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

Поток, к которому нужно прикрепить сигнал.

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

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

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)

Добавлена в: v18.17.0
  • objectMode <логическое>
  • Возвращает: <целое>

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

stream.setDefaultHighWaterMark(objectMode, value)

Добавлена в: v18.17.0
  • objectMode <логическое>
  • value <целое> значение highWaterMark

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

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

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

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

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

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

При расширении потоков учитывайте, какие параметры пользователь может и должен предоставить перед передачей их в базовый конструктор. Например, если реализация делает предположения относительно параметров autoDestroy и emitClose, не позволяйте пользователю их переопределять. Будьте явными в отношении передаваемых параметров, вместо неявной передачи всех параметров.

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

Сценарий использования Класс Реализуемые методы
Только чтение Readable _read()
Только запись Writable _write(), _writev(), _final()
Чтение и запись Duplex _read(), _write(), _writev(), _final()
Обработка записанных данных, затем чтение результата Transform _transform(), _flush(), _final()

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

Избегайте переопределения публичных методов, таких как write(), end(), cork(), uncork(), read() и destroy(), или отправки внутренних событий, таких как 'error', 'data', 'end', 'finish' и 'close' через .emit().

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

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

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

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

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

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

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

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

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

поддержка передачи AbortSignal.

v14.0.0

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

v11.2.0, v10.16.0

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

v10.0.0

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

  • options <Объект>
    • highWaterMark <число> Уровень буфера, когда stream.write() начинает возвращать false. По умолчанию: 16384 (16 Кбайт), или 16 для потоков objectMode.
    • decodeStrings <логическое> Нужно ли кодировать string, переданные в stream.write(), в Buffer (с кодировкой, указанной в вызове stream.write()) перед передачей в stream._write(). Другие типы данных не преобразуются (т.е. Buffer не декодируются в string). Установка в false предотвратит преобразование string в Buffer. По умолчанию: true.
    • defaultEncoding <строка> Кодировка по умолчанию, используемая, когда кодировка не указана в качестве аргумента в stream.write(). По умолчанию: 'utf8'.
    • objectMode <логическое> Является ли stream.write(anyObj) допустимой операцией. Если установлено, можно записывать значения JavaScript, отличные от строк, Buffer или Uint8Array, если это поддерживается реализацией потока. По умолчанию: false.
    • emitClose <логическое> Выводить ли 'close' после уничтожения потока. По умолчанию: true.
    • и т.д.
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 <Функция> Вызовите эту функцию (при необходимости с аргументом ошибки), когда поток завершил инициализацию.

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

поддержка передачи объекта AbortSignal.

v14.0.0

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

v11.2.0, v10.16.0

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

  • options <Объект>
    • highWaterMark <число> Максимальное количество байтов, хранимых во внутреннем буфере перед прекращением чтения из базового ресурса. По умолчанию: 16384 (16 КБ) или 16 для потоков типа objectMode.
    • encoding <строка> Если указано, буферы будут декодированы в строки с использованием указанного кодирования. По умолчанию: null.
    • objectMode <логическое_значение> Следует ли этому потоку вести себя как потоку объектов. Это означает, что stream.read(n) возвращает одно значение вместо Buffer размера n. По умолчанию: false.
    • emitClose <логическое_значение> Следует ли потоку генерировать 'close' после уничтожения. По умолчанию: true.
    • read <Функция> Реализация метода stream._read().
    • destroy <Функция> Реализация метода stream._destroy().
    • construct <Функция> Реализация метода stream._construct().
    • autoDestroy <логическое_значение> Следует ли этому потоку автоматически вызывать .destroy() сам после завершения. По умолчанию: true.
    • signal <СигналПрерывания> Сигнал, представляющий возможную отмену.
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 <Функция> Вызовите эту функцию (по желанию с аргументом ошибки) при завершении инициализации потока.

Метод _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 <число> Количество байтов для асинхронного чтения

Эта функция НЕ ДОЛЖНА вызываться кодом приложения напрямую. Она должна реализовываться дочерними классами и вызываться только внутренними методами класса 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 <Ошибка> Возможная ошибка.
  • callback <Функция> Функция обратного вызова, которая принимает необязательный аргумент ошибки.

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

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

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

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

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

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

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

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

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

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

    this._source = getLowLevelSourceObject();

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

    // When the source ends, push the EOF-signaling `null` chunk.
    this._source.onend = () => {
      this.push(null);
    };
  }
  // _read() will be called when the stream wants to pull more data in.
  // The advisory size argument is ignored in this case.
  _read(size) {
    this._source.readStart();
  }
} 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 не поддерживает множественное наследование, класс stream.Duplex расширяется для реализации потока Duplex (вместо расширения классов stream.Readable и stream.Writable).

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

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

new stream.Duplex(options)
История
Версия Изменения
v8.4.0

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

  • options <Объект> Передаётся в конструкторы Writable и Readable. Также содержит следующие поля:
    • allowHalfOpen <логическое значение> Если установлено в значение false, поток автоматически завершит запись, когда закончится чтение. По умолчанию: true.
    • readable <логическое значение> Устанавливает, будет ли поток Duplex читаемым. По умолчанию: true.
    • writable <логическое значение> Устанавливает, будет ли поток Duplex записываемым. По умолчанию: true.
    • readableObjectMode <логическое значение> Устанавливает objectMode для стороны чтения потока. Не имеет эффекта, если objectMode равно true. По умолчанию: false.
    • writableObjectMode <логическое значение> Устанавливает objectMode для стороны записи потока. Не имеет эффекта, если objectMode предоставлено. По умолчанию: false.
    • readableHighWaterMark <число> Устанавливает highWaterMark для стороны чтения потока. Не имеет эффекта, если highWaterMark предоставлено.
    • writableHighWaterMark <число> Устанавливает 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 может производить выходные данные, которые существенно меньше или больше, чем входные.

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

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

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

new stream.Transform([options])
  • options <Объект> Передаётся в конструкторы Writable и Readable и содержит следующие поля:
    • transform <Функция> Реализация метода stream._transform().
    • flush <Функция> Реализация метода stream._flush().
const { Transform } = require('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 <Функция> Обратный вызов (возможно, с аргументом ошибки и данными), вызываемый при сбросе оставшихся данных.

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

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

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

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

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

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

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

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

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

Возможна ситуация, когда из данных входного куска не генерируется никаких выходных данных.

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

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

transform.prototype._transform = function(data, encoding, callback) {
  callback(null, data);
}; 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('') не рекомендуется.

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

highWaterMark расхождение после вызова readable.setEncoding()

Использование readable.setEncoding() изменит поведение потока highWaterMark в режиме не-объектов.

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

Это не проблема в распространённых случаях с latin1 или ascii. Но следует учитывать это поведение при работе со строками, которые могут содержать символы с несколькими байтами.

© Joyent, Inc. and other Node contributors
Licensed under the MIT License.
Node.js is a trademark of Joyent, Inc. and is used with its permission.
We are not endorsed by or affiliated with Joyent.
https://nodejs.org/dist/latest-v18.x/docs/api/stream.html

Spec-Zone.ru

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