API веб-потоков
Реализация стандарта WHATWG Streams.
Обзор
Стандарт WHATWG Streams (или «веб-потоки») определяет API для обработки потоковых данных. Он похож на API потоков Node.js, но появился позже и стал «стандартным» API для потоковой передачи данных в различных средах JavaScript.
Существует три основных типа объектов:
-
ReadableStream— представляет источник потоковых данных. -
WritableStream— представляет назначение потоковых данных. -
TransformStream— представляет алгоритм преобразования потоковых данных.
Пример ReadableStream
В этом примере создаётся простой ReadableStream, который бесконечно отправляет текущую временную метку performance.now() раз в секунду. Для чтения данных из потока используется асинхронный итерируемый объект.
Модули JavaScript
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);CommonJS
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);
})();Взаимодействие с потоками Node.js
Потоки Node.js можно преобразовать в веб-потоки и наоборот с помощью методов toWeb и fromWeb, доступных у объектов stream.Readable, stream.Writable и stream.Duplex.
Дополнительные сведения см. в соответствующей документации:
API
Класс: ReadableStream
new ReadableStream([underlyingSource [, strategy]])
-
underlyingSource<Object>-
start<Function> Пользовательская функция, которая вызывается сразу после созданияReadableStream.-
controller<ReadableStreamDefaultController> | <ReadableByteStreamController> - Возвращает:
undefinedили промис, выполненный сundefined.
-
-
pull<Function> Пользовательская функция, которая вызывается повторно, пока внутренняя очередьReadableStreamне заполнена. Операция может быть синхронной или асинхронной. Если она асинхронная, функция не будет вызвана повторно, пока не будет выполнен ранее возвращённый промис.-
controller<ReadableStreamDefaultController> | <ReadableByteStreamController> - Возвращает: промис, выполненный с
undefined.
-
-
cancel<Function> Пользовательская функция, которая вызывается при отменеReadableStream.-
reason<any> - Возвращает: промис, выполненный с
undefined.
-
-
type<string> Должно быть'bytes'илиundefined. -
autoAllocateChunkSize<number> Используется только, еслиtypeравно'bytes'. Если задано ненулевое значение, буфер представления автоматически выделяется дляReadableByteStreamController.byobRequest. Если значение не задано, для передачи данных через считыватель по умолчаниюReadableStreamDefaultReaderнеобходимо использовать внутренние очереди потока.
-
-
strategy<Object>-
highWaterMark<number> Максимальный размер внутренней очереди, до которого не применяется обратное давление. -
size<Function> Пользовательская функция, используемая для определения размера каждого фрагмента данных.
-
readableStream.locked
- Тип: <boolean> Устанавливается в
true, если у этого <ReadableStream> есть активный считыватель.
По умолчанию свойство readableStream.locked имеет значение false; при наличии активного считывателя, потребляющего данные потока, оно переключается на true.
readableStream.cancel([reason])
-
reason<any> - Возвращает: промис, выполненный с
undefinedпосле завершения отмены.
readableStream.getReader([options])
-
options<Object>-
mode<string>'byob'илиundefined
-
- Возвращает: <ReadableStreamDefaultReader> | <ReadableStreamBYOBReader>
Модули JavaScript
import { ReadableStream } from 'node:stream/web';
const stream = new ReadableStream();
const reader = stream.getReader();
console.log(await reader.read());CommonJS
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])
-
transform<Object>-
readable<ReadableStream> ПотокReadableStream, в которыйtransform.writableбудет передавать потенциально изменённые данные, полученные из этогоReadableStream. -
writable<WritableStream> ПотокWritableStream, в который будут записываться данные этогоReadableStream.
-
-
options<Object>-
preventAbort<boolean> Еслиtrue, ошибки в этомReadableStreamне приведут к прерываниюtransform.writable. -
preventCancel<boolean> Еслиtrue, ошибки в целевомtransform.writableне приведут к отмене этогоReadableStream. -
preventClose<boolean> Еслиtrue, закрытие этогоReadableStreamне приведёт к закрытиюtransform.writable. -
signal<AbortSignal> Позволяет отменить передачу данных с помощью <AbortController>.
-
- Возвращает: <ReadableStream> Из
transform.readable.
Подключает этот <ReadableStream> к паре <ReadableStream> и <WritableStream>, переданной в аргументе transform, так что данные из этого <ReadableStream> записываются в transform.writable, могут быть преобразованы, а затем передаются в transform.readable. После настройки конвейера возвращается transform.readable.
Переводит readableStream.locked в состояние true на время выполнения операции передачи.
Модули JavaScript
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: ACommonJS
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])
-
destination<WritableStream> <WritableStream>, в который будут записываться данные этогоReadableStream. -
options<Object>-
preventAbort<boolean> Еслиtrue, ошибки в этомReadableStreamне приведут к прерываниюdestination. -
preventCancel<boolean> Еслиtrue, ошибки вdestinationне приведут к отмене этогоReadableStream. -
preventClose<boolean> Еслиtrue, закрытие этогоReadableStreamне приведёт к закрытиюdestination. -
signal<AbortSignal> Позволяет отменить передачу данных с помощью <AbortController>.
-
- Возвращает: промис, выполненный с
undefined
Переводит readableStream.locked в состояние true на время выполнения операции передачи.
readableStream.tee()
- Возвращает: <ReadableStream[]>
Возвращает пару новых экземпляров <ReadableStream>, которым будут передаваться данные этого ReadableStream. Каждый из них получит одинаковые данные.
Переводит readableStream.locked в состояние true.
readableStream.values([options])
-
options<Object>-
preventCancel<boolean> Если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)
-
iterable<Iterable> Объект, реализующий протокол итерируемых объектовSymbol.asyncIteratorилиSymbol.iterator.
Вспомогательный метод, создающий новый <ReadableStream> из итерируемого объекта.
Модули JavaScript
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'CommonJS
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'
})();Чтобы передать полученный <ReadableStream> в <WritableStream>, итерируемый объект <Iterable> должен выдавать последовательность объектов <Buffer>, <TypedArray> или <DataView>.
Модули JavaScript
import { ReadableStream } from 'node:stream/web';
import { Buffer } from 'node:buffer';
async function* asyncIterableGenerator() {
yield Buffer.from('a');
yield Buffer.from('b');
yield Buffer.from('c');
}
const stream = ReadableStream.from(asyncIterableGenerator());
await stream.pipeTo(createWritableStreamSomehow());CommonJS
const { ReadableStream } = require('node:stream/web');
const { Buffer } = require('node:buffer');
async function* asyncIterableGenerator() {
yield Buffer.from('a');
yield Buffer.from('b');
yield Buffer.from('c');
}
const stream = ReadableStream.from(asyncIterableGenerator());
(async () => {
await stream.pipeTo(createWritableStreamSomehow());
})();Класс: ReadableStreamDefaultReader
По умолчанию вызов readableStream.getReader() без аргументов возвращает экземпляр ReadableStreamDefaultReader. Считыватель по умолчанию рассматривает передаваемые через поток фрагменты данных как непрозрачные значения, благодаря чему <ReadableStream> может работать практически с любыми значениями JavaScript.
new ReadableStreamDefaultReader(stream)
-
stream<ReadableStream>
Создаёт новый <ReadableStreamDefaultReader>, блокирующий указанный <ReadableStream>.
readableStreamDefaultReader.cancel([reason])
-
reason<any> - Возвращает: промис, выполненный с
undefined.
Отменяет <ReadableStream> и возвращает промис, который выполняется после отмены базового потока.
readableStreamDefaultReader.closed
- Тип: <Promise> Выполняется с
undefined, когда связанный <ReadableStream> закрывается, или отклоняется, если в потоке возникает ошибка либо блокировка считывателя снимается до завершения закрытия потока.
readableStreamDefaultReader.read()
Запрашивает следующий фрагмент данных из базового <ReadableStream> и возвращает промис, который выполняется с этими данными, когда они становятся доступны.
readableStreamDefaultReader.releaseLock()
Снимает блокировку, установленную этим считывателем на базовый <ReadableStream>.
Класс: ReadableStreamBYOBReader
ReadableStreamBYOBReader — это альтернативный потребитель потоков <ReadableStream>, ориентированных на байты (созданных с underlyingSource.type, равным 'bytes', при создании ReadableStream).
Сокращение BYOB означает «используйте собственный буфер». Этот шаблон позволяет эффективнее считывать данные, ориентированные на байты, избегая лишнего копирования.
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)
-
stream<ReadableStream>
Создаёт новый ReadableStreamBYOBReader, блокирующий указанный <ReadableStream>.
readableStreamBYOBReader.cancel([reason])
-
reason<any> - Возвращает: промис, выполненный с
undefined.
Отменяет <ReadableStream> и возвращает промис, который выполняется после отмены базового потока.
readableStreamBYOBReader.closed
- Тип: <Promise> Выполняется с
undefined, когда связанный <ReadableStream> закрывается, или отклоняется, если в потоке возникает ошибка либо блокировка считывателя снимается до завершения закрытия потока.
readableStreamBYOBReader.read(view[, options])
-
view<Buffer> | <TypedArray> | <DataView> -
options<Object>-
min<number> Если задано, возвращённый промис будет выполнен, только когда станет доступно указанное количество элементовmin. Если значение не задано, промис выполняется, когда доступен хотя бы один элемент.
-
- Возвращает: промис, выполненный с объектом:
-
value<TypedArray> | <DataView> -
done<boolean>
-
Запрашивает следующий фрагмент данных из базового <ReadableStream> и возвращает промис, который выполняется с этими данными, когда они становятся доступны.
Не передавайте в этот метод экземпляр объекта <Buffer> из пула. Буферизованные объекты Buffer создаются с помощью Buffer.allocUnsafe() или Buffer.from() либо часто возвращаются различными обратными вызовами модуля node:fs. В таких объектах Buffer используется общий базовый объект <ArrayBuffer>, содержащий данные всех буферизованных экземпляров Buffer. Когда в readableStreamBYOBReader.read() передаётся Buffer, <TypedArray> или <DataView>, базовый буфер представления ArrayBuffer отсоединяется, что делает недействительными все существующие представления этого буфера ArrayBuffer. Это может иметь катастрофические последствия для приложения.
readableStreamBYOBReader.releaseLock()
Снимает блокировку, установленную этим считывателем на базовый <ReadableStream>.
Класс: ReadableStreamDefaultController
У каждого <ReadableStream> есть контроллер, отвечающий за внутреннее состояние и управление очередью потока. ReadableStreamDefaultController — это реализация контроллера по умолчанию для ReadableStream, не ориентированных на байты.
readableStreamDefaultController.close()
Закрывает <ReadableStream>, с которым связан этот контроллер.
readableStreamDefaultController.desiredSize
- Тип: <number>
Возвращает количество данных, необходимое для заполнения очереди <ReadableStream>.
readableStreamDefaultController.enqueue([chunk])
-
chunk<any>
Добавляет новый фрагмент данных в очередь <ReadableStream>.
readableStreamDefaultController.error([error])
-
error<any>
Сигнализирует об ошибке, из-за которой в <ReadableStream> возникает ошибка и поток закрывается.
Класс: ReadableByteStreamController
У каждого <ReadableStream> есть контроллер, отвечающий за внутреннее состояние и управление очередью потока. ReadableByteStreamController предназначен для потоков ReadableStream, ориентированных на байты.
readableByteStreamController.byobRequest
readableByteStreamController.close()
Закрывает <ReadableStream>, с которым связан этот контроллер.
readableByteStreamController.desiredSize
- Тип: <number>
Возвращает количество данных, необходимое для заполнения очереди <ReadableStream>.
readableByteStreamController.enqueue(chunk)
-
chunk<Buffer> | <TypedArray> | <DataView>
Добавляет новый фрагмент данных в очередь <ReadableStream>.
readableByteStreamController.error([error])
-
error<any>
Сигнализирует об ошибке, из-за которой в <ReadableStream> возникает ошибка и поток закрывается.
Класс: ReadableStreamBYOBRequest
При использовании ReadableByteStreamController в потоках, ориентированных на байты, и при использовании ReadableStreamBYOBReader свойство readableByteStreamController.byobRequest предоставляет доступ к экземпляру ReadableStreamBYOBRequest, представляющему текущий запрос на чтение. Этот объект используется для доступа к ArrayBuffer/TypedArray, предоставленному для заполнения запроса на чтение, и содержит методы для уведомления о предоставлении данных.
readableStreamBYOBRequest.respond(bytesWritten)
-
bytesWritten<number>
Указывает, что в readableStreamBYOBRequest.view записано bytesWritten байтов.
readableStreamBYOBRequest.respondWithNewView(view)
-
view<Buffer> | <TypedArray> | <DataView>
Указывает, что запрос выполнен: байты записаны в новое представление Buffer, TypedArray или DataView.
readableStreamBYOBRequest.view
- Тип: <Buffer> | <TypedArray> | <DataView>
Класс: WritableStream
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]])
-
underlyingSink<Object>-
start<Function> Пользовательская функция, вызываемая сразу после созданияWritableStream.-
controller<WritableStreamDefaultController> - Возвращает:
undefinedили промис, выполненный сundefined.
-
-
write<Function> Пользовательская функция, вызываемая при записи фрагмента данных вWritableStream.-
chunk<any> -
controller<WritableStreamDefaultController> - Возвращает: промис, выполненный с
undefined.
-
-
close<Function> Пользовательская функция, вызываемая при закрытииWritableStream.- Возвращает: промис, выполненный с
undefined.
- Возвращает: промис, выполненный с
-
abort<Function> Пользовательская функция, вызываемая для аварийного закрытияWritableStream.-
reason<any> - Возвращает: промис, выполненный с
undefined.
-
-
type<any> Параметрtypeзарезервирован для будущего использования и должен иметь значение undefined.
-
-
strategy<Object>-
highWaterMark<number> Максимальный размер внутренней очереди до применения обратного давления. -
size<Function> Пользовательская функция для определения размера каждого фрагмента данных.
-
writableStream.abort([reason])
-
reason<any> - Возвращает: промис, выполненный с
undefined.
Аварийно завершает WritableStream. Все записи в очереди будут отменены, а связанные с ними промисы — отклонены.
writableStream.close()
- Возвращает: промис, выполненный с
undefined.
Закрывает WritableStream, когда больше не ожидается новых записей.
writableStream.getWriter()
- Возвращает: <WritableStreamDefaultWriter>
Создает и возвращает новый экземпляр средства записи, который можно использовать для записи данных в WritableStream.
writableStream.locked
- Тип: <boolean>
По умолчанию свойство 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
new WritableStreamDefaultWriter(stream)
-
stream<WritableStream>
Создает новый WritableStreamDefaultWriter, заблокированный для указанного WritableStream.
writableStreamDefaultWriter.abort([reason])
-
reason<any> - Возвращает: промис, выполненный с
undefined.
Аварийно завершает WritableStream. Все записи в очереди будут отменены, а связанные с ними промисы — отклонены.
writableStreamDefaultWriter.close()
- Возвращает: промис, выполненный с
undefined.
Закрывает WritableStream, когда больше не ожидается новых записей.
writableStreamDefaultWriter.closed
- Тип: <Promise> Выполняется с
undefined, когда связанный <WritableStream> закрыт; отклоняется, если в потоке возникает ошибка или блокировка средства записи снята до завершения закрытия потока.
writableStreamDefaultWriter.desiredSize
- Тип: <number>
Объем данных, необходимый для заполнения очереди <WritableStream>.
writableStreamDefaultWriter.ready
- Тип: <Promise> Выполняется с
undefined, когда средство записи готово к использованию.
writableStreamDefaultWriter.releaseLock()
Снимает блокировку этого средства записи с базового <ReadableStream>.
writableStreamDefaultWriter.write([chunk])
-
chunk<any> - Возвращает: промис, выполненный с
undefined.
Добавляет новый фрагмент данных в очередь <WritableStream>.
Класс: WritableStreamDefaultController
WritableStreamDefaultController управляет внутренним состоянием <WritableStream>.
writableStreamDefaultController.error([error])
-
error<any>
Вызывается пользовательским кодом для оповещения об ошибке при обработке данных WritableStream. При вызове <WritableStream> прерывается, а ожидающие записи отменяются.
writableStreamDefaultController.signal
- Тип: <AbortSignal>
AbortSignal, который можно использовать для отмены ожидающих операций записи или закрытия при прерывании <WritableStream>.
Класс: TransformStream
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]]])
-
transformer<Object>-
start<Function> Пользовательская функция, вызываемая сразу после созданияTransformStream.-
controller<TransformStreamDefaultController> - Возвращает:
undefinedили промис, выполненный сundefined
-
-
transform<Function> Пользовательская функция, которая получает и, возможно, изменяет фрагмент данных, записанный вtransformStream.writable, прежде чем передать его вtransformStream.readable.-
chunk<any> -
controller<TransformStreamDefaultController> - Возвращает: промис, выполненный с
undefined.
-
-
flush<Function> Пользовательская функция, вызываемая непосредственно перед закрытием записываемой стороныTransformStream, чтобы обозначить конец процесса преобразования.-
controller<TransformStreamDefaultController> - Возвращает: промис, выполненный с
undefined.
-
-
readableType<any> ПараметрreadableTypeзарезервирован для будущего использования и должен иметь значениеundefined. -
writableType<any> ПараметрwritableTypeзарезервирован для будущего использования и должен иметь значениеundefined.
-
-
writableStrategy<Object>-
highWaterMark<number> Максимальный размер внутренней очереди до применения обратного давления. -
size<Function> Пользовательская функция для определения размера каждого фрагмента данных.
-
-
readableStrategy<Object>-
highWaterMark<number> Максимальный размер внутренней очереди до применения обратного давления. -
size<Function> Пользовательская функция для определения размера каждого фрагмента данных.
-
transformStream.readable
- Тип: <ReadableStream>
transformStream.writable
- Тип: <WritableStream>
Передача с помощью 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
TransformStreamDefaultController управляет внутренним состоянием TransformStream.
transformStreamDefaultController.desiredSize
- Тип: <number>
Объем данных, необходимый для заполнения очереди читаемой стороны.
transformStreamDefaultController.enqueue([chunk])
-
chunk<any>
Добавляет фрагмент данных в очередь читаемой стороны.
transformStreamDefaultController.error([reason])
-
reason<any>
Оповещает читаемую и записываемую стороны об ошибке при обработке преобразуемых данных, в результате чего обе стороны аварийно закрываются.
transformStreamDefaultController.terminate()
Закрывает читаемую сторону потока и вызывает аварийное закрытие записываемой стороны с ошибкой.
Класс: ByteLengthQueuingStrategy
byteLengthQueuingStrategy.highWaterMark
- Тип: <number>
byteLengthQueuingStrategy.size
- Тип: <Function>
Класс: CountQueuingStrategy
countQueuingStrategy.highWaterMark
- Тип: <number>
countQueuingStrategy.size
- Тип: <Function>
Класс: TextEncoderStream
new TextEncoderStream()
Создает новый экземпляр TextEncoderStream.
textEncoderStream.readable
- Тип: <ReadableStream>
textEncoderStream.writable
- Тип: <WritableStream>
Класс: TextDecoderStream
new TextDecoderStream([encoding[, options]])
-
encoding<string> Определяетencoding, поддерживаемую этим экземпляромTextDecoder. По умолчанию:'utf-8'. -
options<Object>-
fatal<boolean>true, если ошибки декодирования считаются фатальными. -
ignoreBOM<boolean> Если задано значениеtrue,TextDecoderStreamвключит метку порядка байтов в результат декодирования. Если задано значениеfalse, метка порядка байтов будет удалена из выходных данных. Этот параметр используется только еслиencodingимеет значение'utf-8','utf-16be'или'utf-16le'. По умолчанию:false.
-
Создает новый экземпляр TextDecoderStream.
textDecoderStream.fatal
- Тип: <boolean>
Значение будет true, если ошибки декодирования приводят к выбрасыванию TypeError.
textDecoderStream.ignoreBOM
- Тип: <boolean>
Значение будет true, если результат декодирования будет содержать метку порядка байтов.
textDecoderStream.readable
- Тип: <ReadableStream>
textDecoderStream.writable
- Тип: <WritableStream>
Class: CompressionStream
new CompressionStream(format)
-
format<string> Одно из значений:'deflate','deflate-raw','gzip'или'brotli'.
compressionStream.readable
- Тип: <ReadableStream>
compressionStream.writable
- Тип: <WritableStream>
Class: DecompressionStream
new DecompressionStream(format)
-
format<string> Одно из значений:'deflate','deflate-raw','gzip'или'brotli'.
decompressionStream.readable
- Тип: <ReadableStream>
decompressionStream.writable
- Тип: <WritableStream>
Вспомогательные функции для чтения потоков
Вспомогательные функции для чтения потоков предоставляют стандартные способы чтения потоков.
Они доступны с помощью:
Модули JavaScript
import {
arrayBuffer,
blob,
buffer,
json,
text,
} from 'node:stream/consumers';CommonJS
const {
arrayBuffer,
blob,
buffer,
json,
text,
} = require('node:stream/consumers');
streamConsumers.arrayBuffer(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с
ArrayBuffer, содержащим всё содержимое потока.
Модули JavaScript
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: 76CommonJS
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)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с объектом <Blob>, содержащим всё содержимое потока.
Модули JavaScript
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: 27CommonJS
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)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с объектом <Buffer>, содержащим всё содержимое потока.
Модули JavaScript
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: 27CommonJS
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.bytes(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с объектом <Uint8Array>, содержащим всё содержимое потока.
Модули JavaScript
import { bytes } 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 bytes(readable);
console.log(`from readable: ${data.length}`);
// Prints: from readable: 27CommonJS
const { bytes } = 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);
bytes(readable).then((data) => {
console.log(`from readable: ${data.length}`);
// Prints: from readable: 27
});
streamConsumers.json(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с содержимым потока, разобранным как строка в кодировке UTF-8, которая затем передаётся в
JSON.parse().
Модули JavaScript
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: 100CommonJS
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)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с содержимым потока, разобранным как строка в кодировке UTF-8.
Модули JavaScript
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: 27CommonJS
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/dist/latest-v24.x/docs/api/webstreams.html