Spec-Zone.ru › Node.js

API потоков веб-приложений

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

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

v18.0.0

Использование этого API больше не вызывает предупреждение во время выполнения.

v16.5.0

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

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

Реализация стандарта потоков WHATWG.

Обзор

Стандарт WHATWG Streams (или "потоки веб-приложений") определяет API для обработки потоковых данных. Он похож на API потоков Node.js Streams, но появился позже и стал "стандартным" API для потоковой передачи данных во многих средах JavaScript.

Существует три основных типа объектов:

  • ReadableStream - Представляет источник потоковых данных.
  • WritableStream - Представляет место назначения для потоковых данных.
  • TransformStream - Представляет алгоритм преобразования потоковых данных.

Пример ReadableStream

Этот пример создает простой ReadableStream, который передает текущую отметку времени performance.now() каждый второй раз, бесконечно. Для чтения данных из потока используется асинхронный итерируемый объект.

Модули MJS

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

import {
  setInterval as every,
} from 'node:timers/promises';

import {
  performance,
} from 'node:perf_hooks';

const SECOND = 1000;

const stream = new ReadableStream({
  async start(controller) {
    for await (const _ of every(SECOND))
      controller.enqueue(performance.now());
  },
});

for await (const value of stream)
  console.log(value);

Модули CJS

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

const {
  setInterval: every,
} = require('node:timers/promises');

const {
  performance,
} = require('node:perf_hooks');

const SECOND = 1000;

const stream = new ReadableStream({
  async start(controller) {
    for await (const _ of every(SECOND))
      controller.enqueue(performance.now());
  },
});

(async () => {
  for await (const value of stream)
    console.log(value);
})();

API

Класс: ReadableStream

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

new ReadableStream([underlyingSource [, strategy]])
Добавлен в: v16.5.0
  • underlyingSource <Объект>
    • start <Функция> Пользовательская функция, которая вызывается немедленно при создании ReadableStream.
      • controller <ReadableStreamDefaultController> | <ReadableByteStreamController>
      • Возвращает: undefined или промис, который разрешается со значением undefined.
    • pull <Функция> Пользовательская функция, которая вызывается повторно, когда внутренняя очередь ReadableStream не заполнена. Операция может быть синхронной или асинхронной. Если асинхронная, функция не будет вызвана снова, пока предыдущий промис не будет выполнен.
      • controller <ReadableStreamDefaultController> | <ReadableByteStreamController>
      • Возвращает: Промис, разрешаемый значением undefined.
    • cancel <Функция> Пользовательская функция, которая вызывается при отмене ReadableStream.
      • reason <любой>
      • Возвращает: Промис, разрешаемый значением undefined.
    • type <строка> Должно быть 'bytes' или undefined.
    • autoAllocateChunkSize <число> Используется только когда type равно 'bytes'. При установке ненулевого значения буфер представления автоматически выделяется для ReadableByteStreamController.byobRequest. Если значение не установлено, необходимо использовать внутренние очереди потока для передачи данных через стандартный читатель ReadableStreamDefaultReader.
  • strategy <Объект>
    • highWaterMark <число> Максимальный размер внутренней очереди, прежде чем будет применено ограничение скорости.
    • size <Функция> Пользовательская функция, используемая для определения размера каждого фрагмента данных.
      • chunk <любой>
      • Возвращает: <число>
readableStream.locked
Добавлен в: v16.5.0
  • Тип: <логическое> Устанавливается в true, если для этого <ReadableStream> существует активный читатель.

Свойство readableStream.locked по умолчанию равно false, и меняется на true при наличии активного читателя, потребляющего данные потока.

readableStream.cancel([reason])
Добавлен в: v16.5.0
  • reason <любой>
  • Возвращает: Промис, разрешаемый значением undefined после завершения отмены.
readableStream.getReader([options])
Добавлен в: v16.5.0
  • options <Объект>
    • mode <строка> 'byob' или undefined
  • Возвращает: <ReadableStreamDefaultReader> | <ReadableStreamBYOBReader>

Модули MJS

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

const stream = new ReadableStream();

const reader = stream.getReader();

console.log(await reader.read());

Модули CJS

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

const stream = new ReadableStream();

const reader = stream.getReader();

reader.read().then(console.log);

Приводит readableStream.locked к true.

readableStream.pipeThrough(transform[, options])
Добавлен в: v16.5.0
  • transform <Объект>
    • readable <ReadableStream> ReadableStream к которому transform.writable будет отправлять потенциально изменённые данные, которые получает от этого ReadableStream.
    • writable <WritableStream> WritableStream в который будут записаны данные этого ReadableStream.
  • options <Объект>
    • preventAbort <логическое> Если true, ошибки в этом ReadableStream не приведут к прерыванию transform.writable.
    • preventCancel <логическое> Если true, ошибки в целевом transform.writable не приведут к отмене этого ReadableStream.
    • preventClose <логическое> Если true, закрытие этого ReadableStream не приведёт к закрытию transform.writable.
    • signal <AbortSignal> Позволяет отменить передачу данных с помощью <AbortController>.
  • Возвращает: <ReadableStream> Из transform.readable.

Подключает этот <ReadableStream> к паре <ReadableStream> и <WritableStream>, предоставленных в аргументе transform, таким образом, что данные из этого <ReadableStream> записываются в transform.writable, возможно преобразуются, а затем отправляются в transform.readable. После настройки конвейера возвращается transform.readable.

Приводит readableStream.locked к true во время активной операции конвейера.

Модули MJS

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

const stream = new ReadableStream({
  start(controller) {
    controller.enqueue('a');
  },
});

const transform = new TransformStream({
  transform(chunk, controller) {
    controller.enqueue(chunk.toUpperCase());
  },
});

const transformedStream = stream.pipeThrough(transform);

for await (const chunk of transformedStream)
  console.log(chunk);
  // Prints: A

Модули CJS

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

const stream = new ReadableStream({
  start(controller) {
    controller.enqueue('a');
  },
});

const transform = new TransformStream({
  transform(chunk, controller) {
    controller.enqueue(chunk.toUpperCase());
  },
});

const transformedStream = stream.pipeThrough(transform);

(async () => {
  for await (const chunk of transformedStream)
    console.log(chunk);
    // Prints: A
})();
readableStream.pipeTo(destination[, options])
Добавлен в: v16.5.0
  • destination <WritableStream> <WritableStream>, в который будут записаны данные этого ReadableStream.
  • options <Объект>
    • preventAbort <логическое> Если true, ошибки в этом ReadableStream не приведут к прерыванию destination.
    • preventCancel <логическое> Если true, ошибки в destination не приведут к отмене этого ReadableStream.
    • preventClose <логическое> Если true, закрытие этого ReadableStream не приведёт к закрытию destination.
    • signal <AbortSignal> Позволяет отменить передачу данных с помощью <AbortController>.
  • Возвращает: Промис, разрешаемый значением undefined

Приводит readableStream.locked к true во время активной операции конвейера.

readableStream.tee()
История
Версия Изменения
v18.10.0, v16.18.0

Поддержка разветвления потока байтов на чтение.

v16.5.0

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

  • Возвращает: <ReadableStream[]>

Возвращает пару новых экземпляров <ReadableStream>, в которые будут перенаправлены данные этого ReadableStream. Каждый из них получит те же данные.

Приводит readableStream.locked к состоянию true.

readableStream.values([options])
Добавлен в: v16.5.0
  • options <Объект>
    • preventCancel <логическое значение> Если true, предотвращает закрытие <ReadableStream> при неожиданном завершении асинхронного итератора. По умолчанию: false.

Создаёт и возвращает асинхронный итератор, используемый для потребления данных этого ReadableStream.

Приводит readableStream.locked к состоянию true во время работы асинхронного итератора.

import { Buffer } from 'node:buffer';

const stream = new ReadableStream(getSomeSource());

for await (const chunk of stream.values({ preventCancel: true }))
  console.log(Buffer.from(chunk).toString()); copy
Асинхронное итерирование

Объект <ReadableStream> поддерживает протокол асинхронного итерирования с использованием синтаксиса for await.

import { Buffer } from 'node:buffer';

const stream = new ReadableStream(getSomeSource());

for await (const chunk of stream)
  console.log(Buffer.from(chunk).toString()); copy

Асинхронный итератор будет потреблять <ReadableStream> до его завершения.

По умолчанию, если асинхронный итератор завершится преждевременно (через break, return, или throw), <ReadableStream> будет закрыт. Чтобы предотвратить автоматическое закрытие <ReadableStream>, используйте метод readableStream.values() для получения асинхронного итератора и установите опцию preventCancel в значение true.

<ReadableStream> не должен быть заблокирован (то есть у него не должно быть активного читателя). Во время асинхронного итерирования <ReadableStream> будет заблокирован.

Передача с помощью postMessage()

Экземпляр <ReadableStream> может быть передан с помощью <MessagePort>.

const stream = new ReadableStream(getReadableSourceSomehow());

const { port1, port2 } = new MessageChannel();

port1.onmessage = ({ data }) => {
  data.getReader().read().then((chunk) => {
    console.log(chunk);
  });
};

port2.postMessage(stream, [stream]); copy

ReadableStream.from(iterable)

Добавлен в: v20.6.0
  • iterable <Итерируемый объект> реализующий протокол итерирования Symbol.asyncIterator или Symbol.iterator.

Утилитарный метод, создающий новый <ReadableStream> из итерируемого объекта.

МОДУЛИ MJS

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

async function* asyncIterableGenerator() {
  yield 'a';
  yield 'b';
  yield 'c';
}

const stream = ReadableStream.from(asyncIterableGenerator());

for await (const chunk of stream)
  console.log(chunk); // Prints: 'a', 'b', 'c'

МОДУЛИ CJS

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

async function* asyncIterableGenerator() {
  yield 'a';
  yield 'b';
  yield 'c';
}

(async () => {
  const stream = ReadableStream.from(asyncIterableGenerator());

  for await (const chunk of stream)
    console.log(chunk); // Prints: 'a', 'b', 'c'
})();

Класс: ReadableStreamDefaultReader

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

По умолчанию, вызов readableStream.getReader() без аргументов вернёт экземпляр ReadableStreamDefaultReader. Читатель по умолчанию обрабатывает куски данных, проходящие через поток, как непрозрачные значения, что позволяет <ReadableStream> работать со всеми типами JavaScript значений.

new ReadableStreamDefaultReader(stream)
Добавлен в: v16.5.0
  • stream <ReadableStream>

Создаёт новый <ReadableStreamDefaultReader>, заблокированный для данного <ReadableStream>.

readableStreamDefaultReader.cancel([reason])
Добавлен в: v16.5.0
  • reason <любой>
  • Возвращает: промис, выполненный с undefined.

Отменяет <ReadableStream> и возвращает промис, который выполняется, когда базовый поток отменён.

readableStreamDefaultReader.closed
Добавлен в: v16.5.0
  • Тип: <Промис> Выполняется с undefined при закрытии связанного <ReadableStream> или отклоняется, если поток ошибается или блокировка читателя освобождена до завершения закрытия потока.
readableStreamDefaultReader.read()
Добавлен в: v16.5.0
  • Возвращает: промис, выполненный с объектом:
    • value <любой>
    • done <логическое значение>

Запрашивает следующий фрагмент данных из базового <ReadableStream> и возвращает промис, который выполняется с данными, когда они станут доступны.

readableStreamDefaultReader.releaseLock()
Добавлен в: v16.5.0

Освобождает блокировку читателя для базового <ReadableStream>.

Класс: ReadableStreamBYOBReader

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

ReadableStreamBYOBReader — это альтернативный потребитель для потоков <ReadableStream> с ориентацией на байты (те, которые созданы с underlyingSource.type равным 'bytes' при создании ReadableStream).

BYOB — сокращение от "bring your own buffer". Это паттерн, позволяющий более эффективно читать данные с ориентацией на байты, избегая лишнего копирования.

import {
  open,
} from 'node:fs/promises';

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

import { Buffer } from 'node:buffer';

class Source {
  type = 'bytes';
  autoAllocateChunkSize = 1024;

  async start(controller) {
    this.file = await open(new URL(import.meta.url));
    this.controller = controller;
  }

  async pull(controller) {
    const view = controller.byobRequest?.view;
    const {
      bytesRead,
    } = await this.file.read({
      buffer: view,
      offset: view.byteOffset,
      length: view.byteLength,
    });

    if (bytesRead === 0) {
      await this.file.close();
      this.controller.close();
    }
    controller.byobRequest.respond(bytesRead);
  }
}

const stream = new ReadableStream(new Source());

async function read(stream) {
  const reader = stream.getReader({ mode: 'byob' });

  const chunks = [];
  let result;
  do {
    result = await reader.read(Buffer.alloc(100));
    if (result.value !== undefined)
      chunks.push(Buffer.from(result.value));
  } while (!result.done);

  return Buffer.concat(chunks);
}

const data = await read(stream);
console.log(Buffer.from(data).toString()); copy
new ReadableStreamBYOBReader(stream)
Добавлен в: v16.5.0
  • stream <ReadableStream>

Создаёт новый ReadableStreamBYOBReader, заблокированный для данного <ReadableStream>.

readableStreamBYOBReader.cancel([reason])
Добавлен в: v16.5.0
  • reason <любой>
  • Возвращает: промис, выполненный с undefined.

Отменяет <ReadableStream> и возвращает промис, который выполняется, когда базовый поток отменён.

readableStreamBYOBReader.closed
Добавлен в: v16.5.0
  • Тип: <Промис> Выполняется с undefined при закрытии связанного <ReadableStream> или отклоняется, если поток ошибается или блокировка читателя освобождена до завершения закрытия потока.
readableStreamBYOBReader.read(view[, options])
История
Версия Изменения
v21.7.0

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

v16.5.0

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

  • view <Буфер> | <Массив типов> | <DataView>
  • options <Объект>
    • min <число> Если установлено, возвращаемый промис будет выполнен только тогда, когда станет доступно min элементов. Если не задано, промис выполняется, когда доступен хотя бы один элемент.
  • Возвращает: промис, выполненный с объектом:
    • value <Массив типов> | <DataView>
    • done <логическое значение>

Запрашивает следующий фрагмент данных из базового <ReadableStream> и возвращает промис, который выполняется с данными, когда они станут доступны.

Не передавайте экземпляр объекта <Buffer> из пула в этот метод. Объекты из пула Buffer создаются с помощью Buffer.allocUnsafe(), или Buffer.from(), или часто возвращаются различными node:fs модульными обработчиками. Эти типы Buffer используют общий базовый объект <ArrayBuffer>, который содержит все данные от всех экземпляров Buffer из пула. Когда Buffer, <TypedArray> или <DataView> передаются в readableStreamBYOBReader.read(), базовый объект ArrayBuffer представления отсоединяется, что делает недействительными все существующие представления, которые могут существовать в этом ArrayBuffer. Это может иметь катастрофические последствия для вашего приложения.

readableStreamBYOBReader.releaseLock()
Добавлен в: v16.5.0

Освобождает блокировку потока чтения (reader) на базовом <ReadableStream>.

Класс: ReadableStreamDefaultController

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

У каждого <ReadableStream> есть контроллер, который отвечает за внутреннее состояние и управление очередью потока. ReadableStreamDefaultController — это реализация контроллера по умолчанию для ReadableStream которые не ориентированы на байты.

readableStreamDefaultController.close()
Добавлен в: v16.5.0

Закрывает <ReadableStream>, к которому связан этот контроллер.

readableStreamDefaultController.desiredSize
Добавлен в: v16.5.0
  • Тип: <число>

Возвращает количество данных, оставшихся для заполнения очереди <ReadableStream>.

readableStreamDefaultController.enqueue([chunk])
Добавлен в: v16.5.0
  • chunk <любой тип>

Добавляет новый фрагмент данных в очередь <ReadableStream>.

readableStreamDefaultController.error([error])
Добавлен в: v16.5.0
  • error <любой тип>

Сигнализирует об ошибке, из-за которой <ReadableStream> получает ошибку и закрывается.

Класс: ReadableByteStreamController

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

Поддержка обработки запроса BYOB от выпущенного читателя.

v16.5.0

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

У каждого <ReadableStream> есть контроллер, который отвечает за внутреннее состояние и управление очередью потока. ReadableByteStreamController предназначен для байтовых ReadableStream.

readableByteStreamController.byobRequest
Добавлен в: v16.5.0
  • Тип: <ReadableStreamBYOBRequest>
readableByteStreamController.close()
Добавлен в: v16.5.0

Закрывает <ReadableStream>, к которому связан этот контроллер.

readableByteStreamController.desiredSize
Добавлен в: v16.5.0
  • Тип: <число>

Возвращает количество данных, оставшихся для заполнения очереди <ReadableStream>.

readableByteStreamController.enqueue(chunk)
Добавлен в: v16.5.0
  • chunk: <Buffer> | <TypedArray> | <DataView>

Добавляет новый фрагмент данных в очередь <ReadableStream>.

readableByteStreamController.error([error])
Добавлен в: v16.5.0
  • error <любой тип>

Сигнализирует об ошибке, из-за которой <ReadableStream> получает ошибку и закрывается.

Класс: ReadableStreamBYOBRequest

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

При использовании ReadableByteStreamController в байтовых потоках и при использовании ReadableStreamBYOBReader, свойство readableByteStreamController.byobRequest предоставляет доступ к экземпляру ReadableStreamBYOBRequest, который представляет текущий запрос чтения. Объект используется для получения доступа к ArrayBuffer/TypedArray, предоставленному для заполнения запроса чтения, и предоставляет методы для сигнализации о предоставлении данных.

readableStreamBYOBRequest.respond(bytesWritten)
Добавлен в: v16.5.0
  • bytesWritten <число>

Сигнализирует о том, что bytesWritten байтов были записаны в readableStreamBYOBRequest.view.

readableStreamBYOBRequest.respondWithNewView(view)
Добавлен в: v16.5.0
  • view <Buffer> | <TypedArray> | <DataView>

Сигнализирует, что запрос был выполнен с байтами, записанными в новое Buffer, TypedArray, или DataView.

readableStreamBYOBRequest.view
Добавлен в: v16.5.0
  • Тип: <Buffer> | <TypedArray> | <DataView>

Класс: WritableStream

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

WritableStream — это пункт назначения, в который отправляются данные потока.

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

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

await stream.getWriter().write('Hello World'); copy
new WritableStream([underlyingSink[, strategy]])
Добавлен в: v16.5.0
  • underlyingSink <Объект>
    • start <Функция> Пользовательская функция, которая вызывается немедленно при создании WritableStream.
      • controller <WritableStreamDefaultController>
      • Возвращает: undefined или промис, выполненный с undefined.
    • write <Функция> Пользовательская функция, которая вызывается, когда фрагмент данных был записан в WritableStream.
      • chunk <любой>
      • controller <WritableStreamDefaultController>
      • Возвращает: промис, выполненный с undefined.
    • close <Функция> Пользовательская функция, которая вызывается при закрытии WritableStream.
      • Возвращает: промис, выполненный с undefined.
    • abort <Функция> Пользовательская функция, вызываемая для прерывания WritableStream.
      • reason <любой>
      • Возвращает: промис, выполненный с undefined.
    • type <любой> Параметр type зарезервирован для будущего использования и должен быть неопределённым.
  • strategy <Объект>
    • highWaterMark <число> Максимальный размер внутренней очереди перед применением обратной связи.
    • size <Функция> Пользовательская функция для определения размера каждого фрагмента данных.
      • chunk <любой>
      • Возвращает: <число>
writableStream.abort([reason])
Добавлен в: v16.5.0
  • reason <любой>
  • Возвращает: промис, выполненный с undefined.

Прерывает WritableStream. Все очереди записей будут отменены, а связанные промисы отклонены.

writableStream.close()
Добавлен в: v16.5.0
  • Возвращает: промис, выполненный с undefined.

Закрывает WritableStream, когда больше нет ожидаемых записей.

writableStream.getWriter()
Добавлен в: v16.5.0
  • Возвращает: <WritableStreamDefaultWriter>

Создаёт и возвращает новый экземпляр писателя, который может быть использован для записи данных в WritableStream.

writableStream.locked
Добавлен в: v16.5.0
  • Тип: <логическое>

Свойство writableStream.locked по умолчанию false, и переключается на true, пока активен присоединённый писатель к этому WritableStream.

Передача с помощью postMessage()

Экземпляр <WritableStream> может быть передан с помощью <MessagePort>.

const stream = new WritableStream(getWritableSinkSomehow());

const { port1, port2 } = new MessageChannel();

port1.onmessage = ({ data }) => {
  data.getWriter().write('hello');
};

port2.postMessage(stream, [stream]); copy

Класс: WritableStreamDefaultWriter

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

new WritableStreamDefaultWriter(stream)
Добавлен в: v16.5.0
  • stream <WritableStream>

Создаёт новый WritableStreamDefaultWriter, который заблокирован для данного WritableStream.

writableStreamDefaultWriter.abort([reason])
Добавлен в: v16.5.0
  • reason <любой>
  • Возвращает: промис, выполненный с undefined.

Прерывает WritableStream. Все очереди записей будут отменены, а связанные промисы отклонены.

writableStreamDefaultWriter.close()
Добавлен в: v16.5.0
  • Возвращает: промис, выполненный с undefined.

Закрывает WritableStream, когда больше нет ожидаемых записей.

writableStreamDefaultWriter.closed
Добавлен в: v16.5.0
  • Тип: <Промис> Выполняется с undefined, когда связанный <WritableStream> закрывается или отклоняется, если поток ошибается или блокировка писателя освобождается до завершения закрытия потока.
writableStreamDefaultWriter.desiredSize
Добавлен в: v16.5.0
  • Тип: <число>

Количество данных, необходимое для заполнения очереди <WritableStream>.

writableStreamDefaultWriter.ready
Добавлен в: v16.5.0
  • Тип: <Промис> Выполняется с undefined, когда писатель готов к использованию.
writableStreamDefaultWriter.releaseLock()
Добавлен в: v16.5.0

Освобождает блокировку этого писателя на базовом <ReadableStream>.

writableStreamDefaultWriter.write([chunk])
Добавлен в: v16.5.0
  • chunk: <любой>
  • Возвращает: промис, выполненный с undefined.

Добавляет новый фрагмент данных в очередь <WritableStream>.

Класс: WritableStreamDefaultController

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

WritableStreamDefaultController управляет внутренним состоянием <WritableStream>.

writableStreamDefaultController.error([error])
Добавлен в: v16.5.0
  • error <любой>

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

writableStreamDefaultController.signal
  • Тип: <AbortSignal> AbortSignal, который можно использовать для отмены ожидающих операций записи или закрытия, когда <WritableStream> прерывается.

Класс: TransformStream

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

TransformStream состоит из <ReadableStream> и <WritableStream>, которые соединены таким образом, что данные, записанные в WritableStream, принимаются и, возможно, преобразуются, прежде чем быть помещены в очередь ReadableStream.

import {
  TransformStream,
} from 'node:stream/web';

const transform = new TransformStream({
  transform(chunk, controller) {
    controller.enqueue(chunk.toUpperCase());
  },
});

await Promise.all([
  transform.writable.getWriter().write('A'),
  transform.readable.getReader().read(),
]); copy
new TransformStream([transformer[, writableStrategy[, readableStrategy]]])
Добавлен в: v16.5.0
  • transformer <Объект>
    • start <Функция> Пользовательская функция, которая вызывается немедленно при создании TransformStream.
      • controller <TransformStreamDefaultController>
      • Возвращает: undefined или промис, выполненный с undefined
    • transform <Функция> Пользовательская функция, которая получает и потенциально изменяет фрагмент данных, записанный в transformStream.writable, перед передачей его в transformStream.readable.
      • chunk <любой>
      • controller <TransformStreamDefaultController>
      • Возвращает: Промис, выполненный с undefined.
    • flush <Функция> Пользовательская функция, которая вызывается непосредственно перед закрытием стороны записи TransformStream, сигнализируя об окончании процесса преобразования.
      • controller <TransformStreamDefaultController>
      • Возвращает: Промис, выполненный с undefined.
    • readableType <любой> параметр readableType зарезервирован для будущего использования и обязательно должен быть undefined.
    • writableType <любой> параметр writableType зарезервирован для будущего использования и обязательно должен быть undefined.
  • writableStrategy <Объект>
    • highWaterMark <число> Максимальный размер внутренней очереди перед применением обратной связи по ограничению пропускной способности.
    • size <Функция> Пользовательская функция, используемая для определения размера каждого фрагмента данных.
      • chunk <любой>
      • Возвращает: <число>
  • readableStrategy <Объект>
    • highWaterMark <число> Максимальный размер внутренней очереди перед применением обратной связи по ограничению пропускной способности.
    • size <Функция> Пользовательская функция, используемая для определения размера каждого фрагмента данных.
      • chunk <любой>
      • Возвращает: <число>
transformStream.readable
Добавлена в: v16.5.0
  • Тип: <Поток чтения>
transformStream.writable
Добавлена в: v16.5.0
  • Тип: <Поток записи>
Передача с помощью postMessage()

Экземпляр <TransformStream> может быть передан с помощью <MessagePort>.

const stream = new TransformStream();

const { port1, port2 } = new MessageChannel();

port1.onmessage = ({ data }) => {
  const { writable, readable } = data;
  // ...
};

port2.postMessage(stream, [stream]); copy

Класс: TransformStreamDefaultController

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

Класс TransformStreamDefaultController управляет внутренним состоянием TransformStream.

transformStreamDefaultController.desiredSize
Добавлена в: v16.5.0
  • Тип: <число>

Количество данных, необходимое для заполнения очереди со стороны чтения.

transformStreamDefaultController.enqueue([chunk])
Добавлена в: v16.5.0
  • chunk <любой>

Добавляет фрагмент данных в очередь со стороны чтения.

transformStreamDefaultController.error([reason])
Добавлена в: v16.5.0
  • reason <любой>

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

transformStreamDefaultController.terminate()
Добавлена в: v16.5.0

Закрывает сторону чтения потока и вызывает внезапное закрытие стороны записи с ошибкой.

Класс: ByteLengthQueuingStrategy

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

new ByteLengthQueuingStrategy(init)
Добавлена в: v16.5.0
  • init <Объект>
    • highWaterMark <число>
byteLengthQueuingStrategy.highWaterMark
Добавлена в: v16.5.0
  • Тип: <число>
byteLengthQueuingStrategy.size
Добавлена в: v16.5.0
  • Тип: <Функция>
    • chunk <любой>
    • Возвращает: <число>

Класс: CountQueuingStrategy

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

Этот класс теперь доступен в глобальном объекте.

v16.5.0

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

new CountQueuingStrategy(init)
Добавлена в: v16.5.0
  • init <Объект>
    • highWaterMark <число>
countQueuingStrategy.highWaterMark
Добавлена в: v16.5.0
  • Тип: <число>
countQueuingStrategy.size
Добавлена в: v16.5.0
  • Тип: <Функция>
    • chunk <любой>
    • Возвращает: <число>

Класс: TextEncoderStream

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

Этот класс теперь доступен в глобальном объекте.

v16.6.0

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

new TextEncoderStream()
Добавлена в: v16.6.0

Создаёт новый экземпляр TextEncoderStream.

textEncoderStream.encoding
Добавлена в: v16.6.0
  • Тип: <строка>

Кодировка, поддерживаемая экземпляром TextEncoderStream.

textEncoderStream.readable
Добавлен в: v16.6.0
  • Тип: <Поток чтения>
textEncoderStream.writable
Добавлен в: v16.6.0
  • Тип: <Поток записи>

Класс: TextDecoderStream

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

Этот класс теперь доступен в глобальном объекте.

v16.6.0

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

new TextDecoderStream([encoding[, options]])
Добавлен в: v16.6.0
  • encoding <строка> Определяет кодировку, которую поддерживает этот экземпляр TextDecoder. По умолчанию: 'utf-8'.
  • options <объект>
    • fatal <логическое значение> true если ошибки декодирования приводят к ошибке.
    • ignoreBOM <логическое значение> Если true, экземпляр TextDecoderStream будет включать маркер порядка байтов в декодированный результат. Если false, маркер порядка байтов будет удалён из вывода. Этот параметр используется только когда encoding равен 'utf-8', 'utf-16be', или 'utf-16le'. По умолчанию: false.

Создаёт новый экземпляр TextDecoderStream.

textDecoderStream.encoding
Добавлен в: v16.6.0
  • Тип: <строка>

Кодировка, поддерживаемая экземпляром TextDecoderStream.

textDecoderStream.fatal
Добавлен в: v16.6.0
  • Тип: <логическое значение>

Значение будет true если ошибки декодирования приводят к выбрасыванию TypeError.

textDecoderStream.ignoreBOM
Добавлен в: v16.6.0
  • Тип: <логическое значение>

Значение будет true если результат декодирования будет включать маркер порядка байтов.

textDecoderStream.readable
Добавлен в: v16.6.0
  • Тип: <Поток чтения>
textDecoderStream.writable
Добавлен в: v16.6.0
  • Тип: <Поток записи>

Класс: CompressionStream

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

Этот класс теперь доступен в глобальном объекте.

v17.0.0

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

new CompressionStream(format)
История
Версия Изменения
v21.2.0, v20.12.0

format теперь принимает значение deflate-raw.

v17.0.0

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

  • format <строка> Одно из значений 'deflate', 'deflate-raw', или 'gzip'.
compressionStream.readable
Добавлен в: v17.0.0
  • Тип: <Поток чтения>
compressionStream.writable
Добавлен в: v17.0.0
  • Тип: <Поток записи>

Класс: DecompressionStream

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

Этот класс теперь доступен в глобальном объекте.

v17.0.0

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

new DecompressionStream(format)
История
Версия Изменения
v21.2.0, v20.12.0

format теперь принимает значение deflate-raw.

v17.0.0

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

  • format <строка> Одно из значений 'deflate', 'deflate-raw', или 'gzip'.
decompressionStream.readable
Добавлен в: v17.0.0
  • Тип: <Поток чтения>
decompressionStream.writable
Добавлен в: v17.0.0
  • Тип: <Поток записи>

Потребители утилит

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

Функции потребителей утилит предоставляют общие параметры для потребления потоков.

К ним можно обратиться с помощью:

Модули MJS

import {
  arrayBuffer,
  blob,
  buffer,
  json,
  text,
} from 'node:stream/consumers';

Модули CJS

const {
  arrayBuffer,
  blob,
  buffer,
  json,
  text,
} = require('node:stream/consumers');
streamConsumers.arrayBuffer(stream)
Добавлен в: v16.7.0
  • stream <Поток чтения> | <stream.Readable> | <Асинхронный итератор>
  • Возвращает: <Promise> Выполняется с ArrayBuffer содержащим полное содержимое потока.

Модули MJS

import { arrayBuffer } from 'node:stream/consumers';
import { Readable } from 'node:stream';
import { TextEncoder } from 'node:util';

const encoder = new TextEncoder();
const dataArray = encoder.encode('hello world from consumers!');

const readable = Readable.from(dataArray);
const data = await arrayBuffer(readable);
console.log(`from readable: ${data.byteLength}`);
// Prints: from readable: 76

Модули CJS

const { arrayBuffer } = require('node:stream/consumers');
const { Readable } = require('node:stream');
const { TextEncoder } = require('node:util');

const encoder = new TextEncoder();
const dataArray = encoder.encode('hello world from consumers!');
const readable = Readable.from(dataArray);
arrayBuffer(readable).then((data) => {
  console.log(`from readable: ${data.byteLength}`);
  // Prints: from readable: 76
});
streamConsumers.blob(stream)
Добавлен в: v16.7.0
  • stream <Поток чтения> | <stream.Readable> | <Асинхронный итератор>
  • Возвращает: <Promise> Выполняется с <Blob> содержащим полное содержимое потока.

Модули MJS

import { blob } from 'node:stream/consumers';

const dataBlob = new Blob(['hello world from consumers!']);

const readable = dataBlob.stream();
const data = await blob(readable);
console.log(`from readable: ${data.size}`);
// Prints: from readable: 27

Модули CJS

const { blob } = require('node:stream/consumers');

const dataBlob = new Blob(['hello world from consumers!']);

const readable = dataBlob.stream();
blob(readable).then((data) => {
  console.log(`from readable: ${data.size}`);
  // Prints: from readable: 27
});
streamConsumers.buffer(stream)
Добавлен в: v16.7.0
  • stream <Поток чтения> | <stream.Readable> | <Асинхронный итератор>
  • Возвращает: <Promise> Выполняется с <Буфер> содержащим полное содержимое потока.

Модули MJS

import { buffer } from 'node:stream/consumers';
import { Readable } from 'node:stream';
import { Buffer } from 'node:buffer';

const dataBuffer = Buffer.from('hello world from consumers!');

const readable = Readable.from(dataBuffer);
const data = await buffer(readable);
console.log(`from readable: ${data.length}`);
// Prints: from readable: 27

Модули CJS

const { buffer } = require('node:stream/consumers');
const { Readable } = require('node:stream');
const { Buffer } = require('node:buffer');

const dataBuffer = Buffer.from('hello world from consumers!');

const readable = Readable.from(dataBuffer);
buffer(readable).then((data) => {
  console.log(`from readable: ${data.length}`);
  // Prints: from readable: 27
});
streamConsumers.json(stream)
Добавлен в: v16.7.0
  • stream <Поток чтения> | <stream.Readable> | <Асинхронный итератор>
  • Возвращает: <Promise> Выполняется с содержимым потока, интерпретированным как строка UTF-8, которая затем передаётся через JSON.parse().

Модули MJS

import { json } from 'node:stream/consumers';
import { Readable } from 'node:stream';

const items = Array.from(
  {
    length: 100,
  },
  () => ({
    message: 'hello world from consumers!',
  }),
);

const readable = Readable.from(JSON.stringify(items));
const data = await json(readable);
console.log(`from readable: ${data.length}`);
// Prints: from readable: 100

Модули CJS

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

const items = Array.from(
  {
    length: 100,
  },
  () => ({
    message: 'hello world from consumers!',
  }),
);

const readable = Readable.from(JSON.stringify(items));
json(readable).then((data) => {
  console.log(`from readable: ${data.length}`);
  // Prints: from readable: 100
});
streamConsumers.text(stream)
Добавлен в: v16.7.0
  • stream <Поток чтения> | <stream.Readable> | <Асинхронный итератор>
  • Возвращает: <Promise> Выполняется со содержимым потока, интерпретированным как строка UTF-8.

Модули MJS

import { text } from 'node:stream/consumers';
import { Readable } from 'node:stream';

const readable = Readable.from('Hello world from consumers!');
const data = await text(readable);
console.log(`from readable: ${data.length}`);
// Prints: from readable: 27

Модули CJS

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

const readable = Readable.from('Hello world from consumers!');
text(readable).then((data) => {
  console.log(`from readable: ${data.length}`);
  // Prints: from readable: 27
});

© 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/webstreams.html

Spec-Zone.ru

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