API потоков веб-приложений
Реализация стандарта потоков WHATWG WHATWG Streams Standard.
Обзор
Стандарт потоков WHATWG (или "потоки веб-приложений") определяет 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
new ReadableStream([underlyingSource [, strategy]])
-
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<Объект>
readableStream.locked
- Тип: <логическое> Устанавливается в
true, если для этого <ReadableStream> есть активный ридер.
Свойство readableStream.locked по умолчанию false, и переключается на true, пока активен ридер, потребляющий данные потока.
readableStream.cancel([reason])
-
reason<любой> - Возвращает: промис, завершённый с
undefinedпосле завершения отмены.
readableStream.getReader([options])
-
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])
-
transform<Объект>-
readable<ReadableStream> Поток, в которыйtransform.writableбудет помещать потенциально измененные данные, полученные из этогоReadableStream. -
writable<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);
Модули 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);
})();
readableStream.pipeTo(destination[, options])
-
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()
- Возвращает: <ReadableStream[]>
Возвращает пару новых экземпляров <ReadableStream>, в которые будут переданы данные этого ReadableStream. Каждый экземпляр получит те же данные.
Приводит к тому, что readableStream.locked будет true.
readableStream.values([options])
-
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 Класс: ReadableStreamDefaultReader
По умолчанию, вызов readableStream.getReader() без аргументов вернёт экземпляр ReadableStreamDefaultReader. Читатель по умолчанию обрабатывает куски данных, проходящие через поток, как непрозрачные значения, что позволяет <ReadableStream> работать с любыми JavaScript-значениями.
new ReadableStreamDefaultReader(stream)
-
stream<ReadableStream>
Создаёт новый <ReadableStreamDefaultReader>, заблокированный для данного <ReadableStream>.
readableStreamDefaultReader.cancel([reason])
-
reason<любое> - Возвращает: промис, завершённый значением
undefined.
Отменяет <ReadableStream> и возвращает промис, который выполнится, когда базовый поток будет отменён.
readableStreamDefaultReader.closed
- Тип: <Промис> Завершается значением
undefined, когда связанный <ReadableStream> закрыт, или отклоняется, если поток ошибается или блокировка читателя освобождается до завершения закрытия потока.
readableStreamDefaultReader.read()
- Возвращает: промис, завершённый объектом:
-
value<ArrayBuffer> -
done<логическое значение>
-
Запрашивает следующий фрагмент данных из базового <ReadableStream> и возвращает промис, который выполняется с данными, как только они станут доступны.
readableStreamDefaultReader.releaseLock()
Освобождает блокировку этого читателя на базовом <ReadableStream>.
Класс: ReadableStreamBYOBReader
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)
-
stream<ReadableStream>
Создаёт новый ReadableStreamBYOBReader, заблокированный для данного <ReadableStream>.
readableStreamBYOBReader.cancel([reason])
-
reason<любое> - Возвращает: промис, завершённый значением
undefined.
Отменяет <ReadableStream> и возвращает промис, который выполнится, когда базовый поток будет отменён.
readableStreamBYOBReader.closed
- Тип: <Промис> Завершается значением
undefined, когда связанный <ReadableStream> закрыт, или отклоняется, если поток ошибается или блокировка читателя освобождается до завершения закрытия потока.
readableStreamBYOBReader.read(view)
-
view<Буфер> | <Массив типа> | <DataView> - Возвращает: промис, завершённый объектом:
-
value<ArrayBuffer> -
done<логическое значение>
-
Запрашивает следующий фрагмент данных из базового <ReadableStream> и возвращает промис, который выполняется с данными, как только они станут доступны.
Не передавайте пулы объектов <Буфер> в этот метод. Пулы объектов Buffer создаются с помощью Buffer.allocUnsafe() или Buffer.from(), или часто возвращаются различными колбэками модуля node:fs. Эти типы объектов Buffer используют общий базовый объект <ArrayBuffer>, который содержит все данные от всех экземпляров пула Buffer. Когда Buffer, <Массив типа> или <DataView> передаются в readableStreamBYOBReader.read(), базовый <ArrayBuffer> просмотра открепляется, делая недействительными все существующие представления, которые могут существовать на этом <ArrayBuffer>. Это может иметь катастрофические последствия для вашего приложения.
readableStreamBYOBReader.releaseLock()
Освобождает блокировку этого читателя на базовом <ReadableStream>.
Класс: ReadableStreamDefaultController
У каждого <ReadableStream> есть контроллер, который отвечает за внутреннее состояние и управление очередью потока. ReadableStreamDefaultController — это реализация контроллера по умолчанию для ReadableStream, которые не ориентированы на байты.
readableStreamDefaultController.close()
Закрывает <ReadableStream>, к которому привязан этот контроллер.
readableStreamDefaultController.desiredSize
- Тип: <число>
Возвращает количество данных, оставшихся для заполнения очереди <ReadableStream>.
readableStreamDefaultController.enqueue([chunk])
-
chunk<любой>
Добавляет новый фрагмент данных в очередь <ReadableStream>.
readableStreamDefaultController.error([error])
-
error<любой>
Указывает ошибку, из-за которой <ReadableStream> генерирует ошибку и закрывается.
Класс: ReadableByteStreamController
Каждый <ReadableStream> имеет контроллер, ответственный за внутреннее состояние и управление очередью потока. ReadableByteStreamController предназначен для байтовых ReadableStream.
readableByteStreamController.byobRequest
readableByteStreamController.close()
Закрывает <ReadableStream>, к которому привязан этот контроллер.
readableByteStreamController.desiredSize
- Тип: <число>
Возвращает количество данных, оставшихся для заполнения очереди <ReadableStream>.
readableByteStreamController.enqueue(chunk)
-
chunk: <Буфер> | <Массив типов> | <DataView>
Добавляет новый фрагмент данных в очередь <ReadableStream>.
readableByteStreamController.error([error])
-
error<любой>
Указывает ошибку, из-за которой <ReadableStream> генерирует ошибку и закрывается.
Класс: ReadableStreamBYOBRequest
При использовании ReadableByteStreamController в потоках байтового типа и при использовании ReadableStreamBYOBReader, свойство readableByteStreamController.byobRequest предоставляет доступ к экземпляру ReadableStreamBYOBRequest, который представляет текущий запрос чтения. Объект используется для доступа к ArrayBuffer/TypedArray, которые были предоставлены для заполнения запроса чтения, и предоставляет методы для сигнализации о том, что данные были предоставлены.
readableStreamBYOBRequest.respond(bytesWritten)
-
bytesWritten<число>
Указывает, что bytesWritten байтов были записаны в readableStreamBYOBRequest.view.
readableStreamBYOBRequest.respondWithNewView(view)
-
view<Буфер> | <Массив типов> | <DataView>
Указывает, что запрос был выполнен с байтами, записанными в новый Buffer, TypedArray или DataView.
readableStreamBYOBRequest.view
- Тип: <Буфер> | <Массив типов> | <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<Объект>-
start<Функция> Пользовательская функция, которая вызывается немедленно при созданииWritableStream.-
controller<WritableStreamDefaultController> - Возвращает:
undefinedили промис, выполненный сundefined.
-
-
write<Функция> Пользовательская функция, которая вызывается, когда фрагмент данных был записан вWritableStream.-
chunk<любой> -
controller<WritableStreamDefaultController> - Возвращает: промис, выполненный с
undefined.
-
-
close<Функция> Пользовательская функция, вызываемая при закрытииWritableStream.- Возвращает: промис, выполненный с
undefined.
- Возвращает: промис, выполненный с
-
abort<Функция> Пользовательская функция, вызываемая для прерывания закрытияWritableStream.-
reason<любой> - Возвращает: промис, выполненный с
undefined.
-
-
type<любой> Параметрtypeзарезервирован для будущего использования и должен быть undefined.
-
-
strategy<Объект>
writableStream.abort([reason])
-
reason<любой> - Возвращает: промис, выполненный с
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
- тип: промис, который выполняется с
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<Объект>-
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<Объект> -
readableStrategy<Объект>
transformStream.readable
- Тип: <Поток чтения>
transformStream.writable
- Тип: <Поток записи>
Передача с помощью 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
- Тип: <число>
Количество данных, необходимое для заполнения очереди стороны чтения.
transformStreamDefaultController.enqueue([chunk])
-
chunk<любой>
Добавляет фрагмент данных в очередь стороны чтения.
transformStreamDefaultController.error([reason])
-
reason<любой>
Сигнализирует обеим сторонам (чтения и записи) о возникновении ошибки при обработке данных преобразования, что приводит к внезапному закрытию обеих сторон.
transformStreamDefaultController.terminate()
Закрывает сторону чтения транспорта и приводит к внезапному закрытию стороны записи с ошибкой.
Класс: ByteLengthQueuingStrategy
new ByteLengthQueuingStrategy(options)
byteLengthQueuingStrategy.highWaterMark
- Тип: <число>
byteLengthQueuingStrategy.size
Класс: CountQueuingStrategy
new CountQueuingStrategy(options)
countQueuingStrategy.highWaterMark
- Тип: <число>
countQueuingStrategy.size
Класс: TextEncoderStream
new TextEncoderStream()
Создаёт новый экземпляр TextEncoderStream.
textEncoderStream.encoding
- Тип: <строка>
Кодировка, поддерживаемая экземпляром TextEncoderStream.
textEncoderStream.readable
- Тип: <ReadableStream>
textEncoderStream.writable
- Тип: <WritableStream>
Класс: TextDecoderStream
new TextDecoderStream([encoding[, options]])
-
encoding<строка> Определяетencoding, который поддерживает этот экземплярTextDecoder. По умолчанию:'utf-8'. -
options<Объект>-
fatal<логическое значение>true, если ошибки декодирования являются фатальными. -
ignoreBOM<логическое значение> Еслиtrue, тоTextDecoderStreamбудет включать маркер порядка байтов в декодированный результат. Еслиfalse, маркер порядка байтов будет удалён из вывода. Этот параметр используется только еслиencodingравен'utf-8','utf-16be'или'utf-16le'. По умолчанию:false.
-
Создаёт новый экземпляр TextDecoderStream.
textDecoderStream.encoding
- Тип: <строка>
Кодировка, поддерживаемая экземпляром TextDecoderStream.
textDecoderStream.fatal
Значение будет true, если ошибки декодирования приводят к выбрасыванию TypeError.
textDecoderStream.ignoreBOM
Значение будет true, если результат декодирования будет включать маркер порядка байтов.
textDecoderStream.readable
- Тип: <ReadableStream>
textDecoderStream.writable
- Тип: <WritableStream>
Класс: CompressionStream
new CompressionStream(format)
-
format<строка> Одно из значений:'deflate'или'gzip'.
compressionStream.readable
- Тип: <ReadableStream>
compressionStream.writable
- Тип: <WritableStream>
Класс: DecompressionStream
new DecompressionStream(format)
-
format<строка> Одно из значений:'deflate'или'gzip'.
decompressionStream.readable
- Тип: <ReadableStream>
decompressionStream.writable
- Тип: <WritableStream>
Потребители утилит
Функции потребителей утилит предоставляют общие варианты для потребления потоков.
К ним можно обратиться с помощью:
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)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <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}`);
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}`);
});
streamConsumers.blob(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <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}`);
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}`);
});
streamConsumers.buffer(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <Promise> Выполняется с <Buffer>, содержащим полное содержимое потока.
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}`);
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}`);
});
streamConsumers.json(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <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}`);
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}`);
});
streamConsumers.text(stream)
-
stream<ReadableStream> | <stream.Readable> | <AsyncIterator> - Возвращает: <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}`);
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}`);
});
© Joyent, Inc. and other Node contributors
Licensed under the MIT License.
Node.js is a trademark of Joyent, Inc. and is used with its permission.
We are not endorsed by or affiliated with Joyent.
https://nodejs.org/dist/latest-v18.x/docs/api/webstreams.html