Spec-Zone.ru › Node.js

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

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

v15.0.0

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

v14.0.0

Метод callback вызывается, если эмитируются 'finish' или 'error'.

v10.0.0

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

v8.0.0

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

v0.9.4

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

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

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

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

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

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

Истина после вызова 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
  • <логическое>

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

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

Получатель свойства objectMode данного потока Writable.

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

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

v8.0.0

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

v6.0.0

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

v0.9.4

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

  • chunk <строка> | <Буфер> | <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 при получении события 'end' в источнике Readable потоке, чтобы поток назначения перестал быть записываемым. Чтобы отключить это поведение по умолчанию, можно передать опцию 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])
История
Версия Изменения
v22.0.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, v18.18.0
Уровень стабильности: 1 - Экспериментально

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

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

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

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

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

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

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

readable.iterator([options])
Добавлена в: v16.3.0
Уровень стабильности: 1 - Экспериментально
  • options <Объект>
    • destroyOnReturn <логическое значение> Если установлено в false, вызов return для асинхронного итератора или выход из цикла for await...of итерирования с помощью break, return, или throw не уничтожит поток. По умолчанию: true.
  • Возвращает: <AsyncIterator> для потребления потока.

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

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

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

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

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

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

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

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

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

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

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

v17.4.0, v16.14.0

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

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

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

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

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

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

v17.4.0, v16.14.0

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

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

Этот метод позволяет перебирать поток. Для каждого фрагмента в потоке будет вызвана функция 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 <СигналПрерывания> позволяет отменить операцию 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 <Функция> | <AsyncФункция> функция для вызова для каждого фрагмента потока.
    • data <любой> фрагмент данных из потока.
    • options <Объект>
      • signal <СигналПрерывания> прерывается, если поток уничтожен, что позволяет прервать вызов fn раньше.
  • options <Объект>
    • concurrency <число> максимальное одновременное вызов fn для вызова потока за раз. По умолчанию: 1.
    • signal <СигналПрерывания> позволяет уничтожить поток, если сигнал прерван.
  • Возвращает: <Промис> промис, равный 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])
Added in: v17.5.0, v16.17.0
Устойчивость: 1 - Экспериментальная
  • fn <Функция> | <AsyncFunction> функция для вызова для каждого фрагмента потока.
    • data <любой тип> фрагмент данных из потока.
    • options <Объект>
      • signal <AbortSignal> прерывается, если поток уничтожен, позволяя прервать вызов fn раньше.
  • options <Объект>
    • concurrency <число> максимальное одновременное обращение к fn для вызова по потоку одновременно. По умолчанию: 1.
    • signal <AbortSignal> позволяет уничтожить поток, если сигнал прерван.
  • Возвращает: <Обещание> обещание, оценивающее первый фрагмент, для которого fn вычислилось с истинным значением или undefined , если элемент не был найден.

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

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

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

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

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

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

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

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

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

import { Readable } from 'node:stream';

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

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

import { Readable } from 'node:stream';

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

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

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

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

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

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

console.log(folderSize); copy

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

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

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

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

console.log(folderSize); copy

Потоки duplex и transform

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

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

v0.9.4

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

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

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

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

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

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

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

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

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

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

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

v8.0.0

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

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

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

После вызова destroy() все последующие вызовы будут no-op и не будут генерировать дополнительных ошибок, кроме тех, что могут быть сгенерированы _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 <Поток> | <ReadableStream> | <WritableStream> Поток чтения и/или записи/веб-поток.
  • 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', 'end', 'finish' и 'close') после вызова callback. Причина в том, что непредвиденные события 'error' (из-за неправильной реализации потоков) не вызывают непредвиденных сбоев. Если это поведение нежелательно, функция очистки, возвращаемая функцией, должна вызываться в обратном вызове:

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

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

stream.pipeline(streams, callback)

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

Добавлена поддержка 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 <Stream[]> | <Iterable[]> | <AsyncIterable[]> | <Function[]> | <ReadableStream[]> | <WritableStream[]> | <TransformStream[]>
  • source <Stream> | <Iterable> | <AsyncIterable> | <Function> | <ReadableStream>
    • Возвращает: <Iterable> | <AsyncIterable>
  • ...transforms <Stream> | <Function> | <TransformStream>
    • source <AsyncIterable>
    • Возвращает: <AsyncIterable>
  • destination <Stream> | <Function> | <WritableStream>
    • source <AsyncIterable>
    • Возвращает: <AsyncIterable> | <Promise>
  • callback <Function> Вызывается, когда обработка данных завершена.
    • err <Error>
    • val Значение, возвращенное методом Promise вызываемым destination.
  • Возвращает: <Stream>

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

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

// Use the pipeline API to easily pipe a series of streams
// together and get notified when the pipeline is fully done.

// A pipeline to gzip a potentially huge tar file efficiently:

pipeline(
  fs.createReadStream('archive.tar'),
  zlib.createGzip(),
  fs.createWriteStream('archive.tar.gz'),
  (err) => {
    if (err) {
      console.error('Pipeline failed.', err);
    } else {
      console.log('Pipeline succeeded.');
    }
  },
); copy

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

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

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

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

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

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

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

stream.compose(...streams)

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

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

v19.8.0, v18.16.0

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

v16.9.0

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

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

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

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

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

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

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

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

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

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

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

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

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

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

let res = '';

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

stream.Readable.isDisturbed(stream)

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

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

stream.isErrored(stream)

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

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

stream.isReadable(stream)

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

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

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

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

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

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

stream.Writable.toWeb(streamWritable)

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

stream.Duplex.from(src)

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

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

v16.8.0

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

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

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

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

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

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

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

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

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

Модули MJS

import { Duplex } from 'node:stream';
import {
  ReadableStream,
  WritableStream,
} from 'node:stream/web';

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

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

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

duplex.write('hello');

for await (const chunk of duplex) {
  console.log('readable', chunk);
}

Модули CJS

const { Duplex } = require('node:stream');
const {
  ReadableStream,
  WritableStream,
} = require('node:stream/web');

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

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

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

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

stream.Duplex.toWeb(streamDuplex)

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

Модули MJS

import { Duplex } from 'node:stream';

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

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

const { value } = await readable.getReader().read();
console.log('readable', value);

Модули CJS

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

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

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

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

stream.addAbortSignal(signal, stream)

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

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

v15.4.0

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

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

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

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

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

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

Или с использованием AbortSignal с потоком чтения как асинхронной итерируемой последовательности:

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

Или с использованием AbortSignal с ReadableStream:

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

addAbortSignal(controller.signal, rs);

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

const reader = rs.getReader();

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

stream.getDefaultHighWaterMark(objectMode)

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

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

stream.setDefaultHighWaterMark(objectMode, value)

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

v15.5.0

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

v14.0.0

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

v11.2.0, v10.16.0

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

v10.0.0

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

  • options <Объект>
    • highWaterMark <число> Уровень буфера, когда stream.write() начинает возвращать false. По умолчанию: 65536 (64 КБ), или 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.
    • write <Функция> Реализация метода stream._write().
    • writev <Функция> Реализация метода stream._writev().
    • destroy <Функция> Реализация метода stream._destroy().
    • final <Функция> Реализация метода stream._final().
    • construct <Функция> Реализация метода stream._construct().
    • autoDestroy <логическое значение> Должен ли поток автоматически вызывать .destroy() после завершения. По умолчанию: true.
    • signal <AbortSignal> Сигнал, представляющий возможность отмены.
const { Writable } = require('node:stream');

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

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

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

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

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

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

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

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

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

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

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

Эта необязательная функция будет вызвана через tick после возвращения конструктора потока, откладывая любые вызовы _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 для записи, преобразованные из 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 <Буфер> | <строка> Экземпляр буфера или строка, содержащая данные для записи. Значение будет строкой, если 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. Хотя этот конкретный экземпляр потока на запись Writable не обладает практической ценностью, пример иллюстрирует каждый из необходимых элементов пользовательского экземпляра потока Writable:

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

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

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

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

class StringWritable extends Writable {
  constructor(options) {
    super(options);
    this._decoder = new StringDecoder(options && options.defaultEncoding);
    this.data = '';
  }
  _write(chunk, encoding, callback) {
    if (encoding === 'buffer') {
      chunk = this._decoder.write(chunk);
    }
    this.data += chunk;
    callback();
  }
  _final(callback) {
    this.data += this._decoder.end();
    callback();
  }
}

const euro = [[0xE2, 0x82], [0xAC]].map(Buffer.from);
const w = new StringWritable();

w.write('currency: ');
w.write(euro[0]);
w.end(euro[1]);

console.log(w.data); // currency: € copy

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

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

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

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

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

v15.5.0

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

v14.0.0

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

v11.2.0, v10.16.0

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

  • options <Объект>
    • highWaterMark <число> Максимальное количество байтов для хранения во внутрене буфере перед прекращением чтения из базового ресурса. По умолчанию: 65536 (64 КБ) или 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()) на создаваемом чтении.

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])
История
Версия Изменения
v22.0.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 указано. По умолчанию: false.
    • readableHighWaterMark <число> Устанавливает highWaterMark для стороны чтения потока. Не имеет эффекта, если указано highWaterMark.
    • writableHighWaterMark <число> Устанавливает highWaterMark для стороны записи потока. Не имеет эффекта, если указано highWaterMark.
const { Duplex } = require('node:stream');

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

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

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

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

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

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

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

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

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

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

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

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

class MyDuplex extends Duplex {
  constructor(source, options) {
    super(options);
    this[kSource] = source;
  }

  _write(chunk, encoding, callback) {
    // The underlying source only deals with strings.
    if (Buffer.isBuffer(chunk))
      chunk = chunk.toString();
    this[kSource].writeSomeData(chunk);
    callback();
  }

  _read(size) {
    this[kSource].fetchSomeData(size, (data, encoding) => {
      this.push(Buffer.from(data, encoding));
    });
  }
} copy

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

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

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

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

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

// All Transform streams are also Duplex Streams.
const myTransform = new Transform({
  writableObjectMode: true,

  transform(chunk, encoding, callback) {
    // Coerce the chunk to a number if necessary.
    chunk |= 0;

    // Transform the chunk into something else.
    const data = chunk.toString(16);

    // Push the data onto the readable queue.
    callback(null, '0'.repeat(data.length % 2) + data);
  },
});

myTransform.setEncoding('ascii');
myTransform.on('data', (chunk) => console.log(chunk));

myTransform.write(1);
// Prints: 01
myTransform.write(10);
// Prints: 0a
myTransform.write(100);
// Prints: 64 copy

Реализация потока преобразования

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

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

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

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

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

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

Spec-Zone.ru

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