Spec-Zone.ru › Node.js 20 LTS

Поток[src]

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

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

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

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

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

Для доступа к модулю node:stream:

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

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

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

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

Типы потоков

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

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

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

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

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

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

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

stream.pipeline(streams[, options])

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

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

v15.0.0

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

  • streams <Поток[]> | <Итерируемый массив[]> | <Асинхронно итерируемый массив[]> | <Функция[]>
  • source <Поток> | <Итерируемый> | <Асинхронно итерируемый> | <Функция>
    • Возвращает: <Обещание> | <Асинхронно итерируемый>
  • ...transforms <Поток> | <Функция>
    • source <Асинхронно итерируемый>
    • Возвращает: <Обещание> | <Асинхронно итерируемый>
  • destination <Поток> | <Функция>
    • source <Асинхронно итерируемый>
    • Возвращает: <Обещание> | <Асинхронно итерируемый>
  • options <Объект> Параметры конвейера
    • signal <Объект отмены>
    • end <логическое> Завершить целевой поток при завершении исходного. Потоки преобразования всегда завершаются, даже если это значение false. По умолчанию: true.
  • Возвращает: <Обещание> Выполняется, когда конвейер завершен.

Модули 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])

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

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

v15.0.0

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

  • stream <Поток> | <Поток чтения> | <Поток записи> Поток чтения и/или записи/веб-поток.
  • options <Объект>
    • error <логическое> | <неопределённо>
    • readable <логическое> | <неопределённо>
    • writable <логическое> | <неопределённо>
    • signal: <Объект отмены> | <неопределённо>
  • Возвращает: <Обещание> Выполняется, когда поток больше не читается или не записывается.

Модули 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, работают исключительно со строками, <Буфер>, <Массивом типов> и <Представлением данных>:

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

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

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

Буферизация

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

server.listen(1337);

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

v0.9.4

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

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

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

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

Если вызов stream.write(chunk) возвращает false, событие 'drain' будет излучено, когда будет уместно возобновить запись данных в поток.

// Write the data to the supplied writable stream one million times.
// Be attentive to back-pressure.
function writeOneMillionTimes(writer, data, encoding, callback) {
  let i = 1000000;
  write();
  function write() {
    let ok = true;
    do {
      i--;
      if (i === 0) {
        // Last time!
        writer.write(data, encoding, callback);
      } else {
        // See if we should continue, or wait.
        // Don't pass the callback, because we're not done yet.
        ok = writer.write(data, encoding);
      }
    } while (i > 0 && ok);
    if (i > 0) {
      // Had to stop early!
      // Write some more once it drains.
      writer.once('drain', write);
    }
  }
} copy
Событие: 'error'
Добавлен в: v0.9.4
  • <Ошибка>

Событие 'error' излучается, если при записи или передаче данных возникла ошибка. Обработчик вызывается с одним аргументом типа Error.

Поток закрывается, когда излучается событие 'error', если параметр 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, удаляя этот поток записи из списка его назначений.

Также излучается, если этот поток записи генерирует ошибку, когда в него перенаправляется поток 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])
История
Версия Изменения
v20.13.0

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

v15.0.0

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

v14.0.0

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

v10.0.0

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

v8.0.0

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

v0.9.4

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

  • chunk <строка> | <Buffer> | <TypedArray> | <DataView> | <любой> Дополнительные данные для записи. Для потоков, не работающих в режиме объектов, chunk должен быть <строкой>, <Buffer>, <TypedArray> или <DataView>. Для потоков в режиме объектов, 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() устанавливает кодировку по умолчанию encoding для потока Writable.

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

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

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

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

Если метод writable.cork() вызывается несколько раз для потока, необходимо вызвать такое же количество раз метод writable.uncork() для сброса буферизованных данных.

stream.cork();
stream.write('some ');
stream.cork();
stream.write('data ');
process.nextTick(() => {
  stream.uncork();
  // The data will not be flushed until uncork() is called a second time.
  stream.uncork();
}); copy

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

writable.writable
Добавлен в: v11.4.0
  • <булево>

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

writable.writableAborted
Добавлен в: v18.0.0, v16.17.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])
История
Версия Изменения
v20.13.0

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

v8.0.0

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

v6.0.0

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

v0.9.4

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

  • chunk <строка> | <Буфер> | <TypedArray> | <DataView> | <любое> Данные для записи. Для потоков, не работающих в режиме объектов, chunk должно быть строкой <string>, буфером <Buffer>, <TypedArray> или <DataView>. Для потоков в режиме объектов, 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.

Потоковые потоки для чтения

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

v8.0.0

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

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

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

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

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

readable.closed
Добавлен в: v18.0.0
  • <логическое значение>

Является true после того, как сгенерировано 'close'.

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

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

readable.isPaused()
Добавлен в: v0.11.14
  • Возвращает: <логическое значение>

Метод 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
  • Возвращает: <текущий объект>

Метод 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 <Объект> Параметры для соединения
    • end <логическое значение> Завершить писатель, когда читатель завершит. По умолчанию: 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 <число> Необязательный аргумент для указания объёма считываемых данных.
  • Возвращает: <строка> | <Буфер> | <null> | <любой>

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

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

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

Аргумент size должен быть меньше или равен 1 Гигабайту.

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

const readable = getReadableStreamSomehow();

// 'readable' may be triggered multiple times as data is buffered in
readable.on('readable', () => {
  let chunk;
  console.log('Stream is readable (new data received in buffer)');
  // Use a loop to make sure we read all currently available data
  while (null !== (chunk = readable.read())) {
    console.log(`Read ${chunk.length} bytes of data...`);
  }
});

// 'end' will be triggered once when there is no more data available
readable.on('end', () => {
  console.log('Reached end of stream.');
}); 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
  • <логическое>

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

readable.readableAborted
Добавлен в: v16.8.0
Стабильность: 1 - Экспериментально
  • <логическое>

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

readable.readableDidRead
Добавлен в: v16.7.0, v14.18.0
Стабильность: 1 - Экспериментально
  • <логическое>

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

readable.readableEncoding
Добавлен в: v12.7.0
  • <null> | <строка>

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

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

Становится true при выводе события 'end'.

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

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

readable.readableFlowing
Добавлен в: v9.4.0
  • <логическое>

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

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

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

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

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

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

Геттер для свойства 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 <строка> Кодировка для использования.
  • Возвращает: <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])
История
Версия Изменения
v20.13.0

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

v8.0.0

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

v0.9.11

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

  • chunk <Буфер> | <Массив типов> | <DataView> | <строка> | <null> | <любой> Чанк данных для добавления в начало очереди чтения. Для потоков, не работающих в режиме объектов, chunk должен быть <строкой>, <буфером>, <массивом типов>, <DataView> или null. Для потоков в режиме объектов chunk может быть любым значением JavaScript.
  • encoding <строка> Кодировка строковых чанков. Должна быть допустимой кодировкой 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 <Поток> Поток чтения "старого стиля"
  • Возвращает: <this>

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

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

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

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

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

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

v10.0.0

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

  • Возвращает: <AsyncIterator> для полного потребления потока.
const fs = require('node:fs');

async function print(readable) {
  readable.setEncoding('utf8');
  let data = '';
  for await (const chunk of readable) {
    data += chunk;
  }
  console.log(data);
}

print(fs.createReadStream('file')).catch(console.error); copy

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

readable[Symbol.asyncDispose]()
Добавлена в: v20.4.0
Уровень стабильности: 1 - Экспериментальная

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

readable.compose(stream[, options])
Добавлена в: v19.1.0, v18.13.0
Уровень стабильности: 1 - Экспериментальная
  • stream <Поток> | <Итерируемый объект> | <Асинхронно итерируемый объект> | <Функция>
  • options <Объект>
    • signal <AbortSignal> позволяет уничтожить поток, если сигнал прерван.
  • Возвращает: <Дуплексный> поток, составленный из потока 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 <логическое> Если установлено в false, вызов return на асинхронном итераторе или выход из итерации for await...of с использованием break, return, или throw не уничтожит поток. По умолчанию: true.
  • Возвращает: <AsyncIterator> для потребления потока.

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

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

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

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

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

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

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

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

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

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

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

v17.4.0, v16.14.0

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

Уровень стабильности: 1 - Экспериментальная
  • fn <Функция> | <AsyncFunction> функция для обработки каждого фрагмента в потоке.
    • data <любой тип> фрагмент данных из потока.
    • options <Объект>
      • signal <AbortSignal> прерывается, если поток разрушен, что позволяет прервать вызов fn раньше.
  • options <Объект>
    • concurrency <число> максимальное одновременное выполнение вызова fn для потока. По умолчанию: 1.
    • highWaterMark <число> количество элементов для буферизации, пока ожидается обработка сопоставленных элементов пользователем. По умолчанию: concurrency * 2 - 1.
    • signal <AbortSignal> позволяет разрушить поток, если сигнал прерван.
  • Возвращает: <Чтение> поток, отображённый с помощью функции fn.

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

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

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

добавлен highWaterMark в опциях.

v17.4.0, v16.14.0

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

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

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

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

// With a synchronous predicate.
for await (const chunk of Readable.from([1, 2, 3, 4]).filter((x) => x > 2)) {
  console.log(chunk); // 3, 4
}
// With an asynchronous predicate, making at most 2 queries at a time.
const resolver = new Resolver();
const dnsResults = Readable.from([
  'nodejs.org',
  'openjsf.org',
  'www.linuxfoundation.org',
]).filter(async (domain) => {
  const { address } = await resolver.resolve4(domain, { ttl: true });
  return address.ttl > 60;
}, { concurrency: 2 });
for await (const result of dnsResults) {
  // Logs domains with more than 60 seconds on the resolved dns record.
  console.log(result);
} copy
readable.forEach(fn[, options])
Добавлен в: v17.5.0, v16.15.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <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, v16.15.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, v16.15.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, v16.17.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncФункция> функция, которая вызывается для каждого фрагмента потока.
    • 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, v16.15.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncФункция> функция, которая вызывается для каждого фрагмента потока.
    • 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, v16.15.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncГенераторФункция> | <AsyncФункция> функция для отображения каждого фрагмента в потоке.
    • 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, v16.15.0
Устойчивость: 1 - Экспериментальная
  • limit <число> количество фрагментов, которые нужно выбросить из readable.
  • 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, v16.15.0
Устойчивость: 1 - Экспериментальная
  • limit <число> количество фрагментов для взятия из readable.
  • 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])
История
Версия Изменения
v20.3.0

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

v17.5.0, v16.15.0

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

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

Этот метод возвращает новый поток с частями исходного потока, соединёнными с счётчиком в формате [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, v16.15.0
Устойчивость: 1 - Экспериментальный
  • fn <Функция> | <АсинхроннаяФункция> функция-редуктор для вызова каждой части потока.
    • 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
  • <логическое>

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

Это можно вручную изменить, чтобы изменить поведение полуоткрытого потока типа «дуплекс» существующего экземпляра потока, но это необходимо сделать до выдачи события '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 <Ошибка>
  • Возвращает: <текущий>

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

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

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

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

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

v15.11.0

Была добавлена опция signal.

v14.0.0

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

v14.0.0

Выдача события 'close' до события 'end' в потоке типа «дуплекс» Readable вызовет ошибку ERR_STREAM_PREMATURE_CLOSE.

v14.0.0

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

v10.0.0

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

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

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

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

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

stream.pipeline(streams, callback)

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

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

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 <Поток[]> | <Итерируемый[]> | <АсинхронноИтерируемый[]> | <Функция[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]>
  • source <Поток> | <Итерируемый> | <АсинхронноИтерируемый> | <Функция> | <ReadableStream>
    • Возвращает: <Итерируемый> | <АсинхронноИтерируемый>
  • ...transforms <Поток> | <Функция> | <TransformStream>
    • source <АсинхронноИтерируемый>
    • Возвращает: <АсинхронноИтерируемый>
  • destination <Поток> | <Функция> | <WritableStream>
    • source <АсинхронноИтерируемый>
    • Возвращает: <АсинхронноИтерируемый> | <Promise>
  • callback <Функция> Вызывается, когда конвейер полностью завершен.
    • err <Ошибка>
    • val Значение, возвращаемое Promise вызовом destination.
  • Возвращает: <Поток>

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

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

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

  • потоков, которые уже выпустили 'end' или 'close'.
  • потоков, которые уже выпустили '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)

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

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

v19.8.0, v18.16.0

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

v16.9.0

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

Устойчивость: 1 - stream.compose является экспериментальной.
  • streams <Поток[]> | <Итерируемый[]> | <АсинхронноИтерируемый[]> | <Функция[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]> | <Duplex[]> | <Функция>
  • Возвращает: <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 <Итерируемый> Объект, реализующий Symbol.asyncIterator или Symbol.iterator итерируемый протокол. Генерирует событие 'error', если передано значение null.
  • options <Объект> Опции, предоставленные 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 <Объект>
    • encoding <строка>
    • highWaterMark <число>
    • objectMode <логическое значение>
    • 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>
  • Возвращает: <логическое значение>

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

stream.isReadable(stream)

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

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

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

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

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

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

stream.Writable.toWeb(streamWritable)

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

stream.Duplex.from(src)

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

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

v16.8.0

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

  • src <Поток> | <Blob> | <ArrayBuffer> | <строка> | <Итерируемый объект> | <Асинхронно итерируемый объект> | <Функция асинхронного генератора> | <Асинхронная функция> | <Promise> | <Объект> | <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 <Объект>
    • readable <ReadableStream>
    • writable <WritableStream>
  • options <Объект>
    • allowHalfOpen <boolean>
    • decodeStrings <boolean>
    • encoding <строка>
    • highWaterMark <число>
    • 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)

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

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

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

stream.setDefaultHighWaterMark(objectMode, value)

Добавлен в: v19.9.0
  • objectMode <boolean>
  • 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 для автоматической destroy() потока при выводе 'finish' или ошибках.

v10.0.0

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

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

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

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

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

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

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

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

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

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

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

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

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

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

Метод _destroy() вызывается методом writable.destroy(). Он может быть переопределен дочерними классами, но не должен вызываться напрямую. Кроме того, 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.

Пользовательские потоки на чтение обязательно должны вызывать конструктор 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 для автоматического закрытия потока при испускании события 'end' или возникновении ошибки.

  • options <Object>
    • 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 <AbortSignal> Сигнал, представляющий возможную отмену.
const { Readable } = require('node:stream');

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

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

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

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

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

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

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

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

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

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

v8.0.0

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

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

Когда chunk — <буфер>, <массив типизированных данных>, <DataView> или <строка>, данные будут добавлены во внутреннюю очередь для использования пользователями потока. Передача 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 равно true. По умолчанию: 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

Наиболее важной частью дуплексного потока является то, что стороны чтения 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().

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

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

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

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

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

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

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

Spec-Zone.ru

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