Узнайте, как использовать потоки данных для чтения, записи и преобразования с помощью Streams API.
API Streams позволяет программно получать доступ к потокам данных, поступающим по сети или создаваемым любыми способами локально, и обрабатывать их с помощью JavaScript. Потоковая обработка предполагает разбиение ресурса, который вы хотите получить, отправить или преобразовать, на небольшие фрагменты, а затем обработку этих фрагментов побитово. Хотя потоковая обработка и так используется браузерами при получении таких ресурсов, как HTML или видео для отображения на веб-страницах, эта возможность никогда не была доступна JavaScript до появления функции fetch с потоками в 2015 году.
Обратите внимание, что потоковая передача данных технически возможна с помощью XMLHttpRequest , но это не самое оптимальное решение. Вот пример использования XMLHttpRequest на GitHub .
Раньше, если вам нужно было обработать какой-либо ресурс (будь то видео, текстовый файл и т. д.), вам приходилось загружать весь файл, ждать его десериализации в подходящий формат, а затем обрабатывать его. С появлением потоковой обработки в JavaScript все меняется. Теперь вы можете обрабатывать необработанные данные с помощью JavaScript постепенно, как только они становятся доступны на стороне клиента, без необходимости создавать буфер, строку или двоичные данные. Это открывает множество вариантов использования, некоторые из которых я перечислю ниже:
- Видеоэффекты: передача читаемого видеопотока через поток преобразований, который применяет эффекты в реальном времени.
- Сжатие (де)компрессии данных: передача файлового потока через поток преобразования, который выборочно (де)компрессирует его.
- Декодирование изображений: передача потока HTTP-ответа через поток преобразования, который декодирует байты в данные растрового изображения, а затем через другой поток преобразования, который преобразует растровые изображения в PNG. Если это установлено внутри обработчика
fetchсервисного работника, это позволяет прозрачно добавлять полифилы к новым форматам изображений, таким как AVIF.
Поддержка браузеров
ReadableStream и WritableStream
ТрансформСтрим
Основные концепции
Прежде чем подробно рассказать о различных типах потоков, позвольте мне представить несколько основных понятий.
Куски
Блок данных — это отдельный фрагмент данных , который записывается в поток или считывается из него. Он может быть любого типа; потоки могут даже содержать блоки данных разных типов. В большинстве случаев блок данных не является наиболее атомарной единицей данных для данного потока. Например, поток байтов может содержать блоки, состоящие из 16 единиц Uint8Array размером 16 КиБ, вместо отдельных байтов.
Читаемые потоки
Читаемый поток представляет собой источник данных, из которого можно считывать информацию. Другими словами, данные поступают из читаемого потока. Конкретнее, читаемый поток — это экземпляр класса ReadableStream .
Записываемые потоки
Поток, доступный для записи, представляет собой место назначения для данных, в которое можно записывать информацию. Другими словами, данные поступают в поток, доступный для записи. Конкретно, поток, доступный для записи, является экземпляром класса WritableStream .
Преобразовать потоки
Поток преобразования состоит из пары потоков : потока для записи, известного как записываемая сторона, и потока для чтения, известного как читаемая сторона. В качестве метафоры из реальной жизни можно привести синхронный переводчик , который переводит с одного языка на другой на лету. В случае потока преобразования запись в записываемую сторону приводит к тому, что новые данные становятся доступными для чтения из читаемой стороны. Конкретно, любой объект со свойством writable и свойством readable может служить потоком преобразования. Однако стандартный класс TransformStream упрощает создание такой пары, которая должным образом запутана.
Трубные цепи
Потоки в основном используются путем их соединения друг с другом по каналам. Читаемый поток может быть напрямую соединен с записываемым потоком с помощью метода pipeTo() читаемого потока, или же он может быть сначала соединен через один или несколько потоков преобразования с помощью метода pipeThrough() читаемого потока. Набор потоков, соединенных таким образом, называется цепочкой каналов.
Обратное давление
После построения трубопроводной цепи она будет передавать сигналы, указывающие на скорость потока частиц. Если какой-либо этап цепи еще не может принимать частицы, он передает сигнал в обратном направлении по трубопроводной цепи, пока в конечном итоге исходному источнику не будет дано указание прекратить производство частиц с такой скоростью. Этот процесс нормализации потока называется противодавлением.
Тиинг
Читаемый поток можно заблокировать (название происходит от формы заглавной буквы «Т») с помощью метода tee() . Это заблокирует поток, то есть сделает его недоступным для непосредственного использования; однако при этом будут созданы два новых потока , называемых ветвями, которые можно использовать независимо друг от друга. Блокировка также важна, потому что потоки нельзя перемотать назад или перезапустить, об этом подробнее позже.
Механика читаемого потока
Читаемый поток — это источник данных, представленный в JavaScript объектом ReadableStream , который поступает из базового источника. Конструктор ReadableStream() создает и возвращает объект читаемого потока из заданных обработчиков. Существует два типа базовых источников:
- Источники push-уведомлений постоянно передают вам данные после того, как вы к ним обращаетесь, и только от вас зависит, начнете ли вы, приостановите или отмените доступ к потоку. Примерами могут служить прямые видеотрансляции, события, отправляемые сервером, или WebSocket.
- Для работы с источниками данных, требующими получения информации, необходимо явно запрашивать у них данные после установления соединения. Примерами могут служить HTTP-операции с помощью вызовов
fetch()илиXMLHttpRequest.
Потоковые данные считываются последовательно небольшими фрагментами, называемыми блоками . Блоки, помещенные в поток, называются блоками, поставленными в очередь . Это означает, что они находятся в очереди, ожидая чтения. Внутренняя очередь отслеживает блоки, которые еще не были прочитаны.
Стратегия организации очереди — это объект, определяющий, как поток должен сигнализировать об обратном давлении в зависимости от состояния его внутренней очереди. Стратегия организации очереди присваивает размер каждому фрагменту и сравнивает общий размер всех фрагментов в очереди с заданным числом, известным как « верхняя граница» .
Фрагменты данных в потоке считываются устройством чтения . Это устройство извлекает данные по одному фрагменту за раз, позволяя выполнять любые необходимые операции. Устройство чтения и связанный с ним код обработки называются потребителем .
Следующая конструкция в этом контексте называется контроллером . Каждому читаемому потоку соответствует контроллер, который, как следует из названия, позволяет управлять потоком.
Одновременно читать поток может только один пользователь; когда создается пользователь и начинает читать поток (то есть становится активным пользователем ), он блокируется для этого потока. Если вы хотите, чтобы другой пользователь взял на себя чтение вашего потока, обычно необходимо освободить первого пользователя, прежде чем делать что-либо еще (хотя можно также блокировать потоки).
Создание читаемого потока
Для создания читаемого потока необходимо вызвать его конструктор ReadableStream() . Конструктор имеет необязательный аргумент underlyingSource , представляющий собой объект с методами и свойствами, определяющими поведение созданного экземпляра потока.
underlyingSource
Для этого можно использовать следующие необязательные методы, определяемые разработчиком:
-
start(controller): Вызывается немедленно при создании объекта. Метод может получить доступ к источнику потока и выполнить любые другие действия, необходимые для настройки функциональности потока. Если этот процесс должен выполняться асинхронно, метод может возвращать промис для сигнализации успеха или неудачи. Параметрcontroller, передаваемый этому методу, представляет собойReadableStreamDefaultController. -
pull(controller): Может использоваться для управления потоком по мере получения новых фрагментов. Вызывается многократно, пока внутренняя очередь фрагментов потока не заполнена, до тех пор, пока очередь не достигнет своего максимального значения. Если результатом вызоваpull()является промис,pull()не будет вызываться снова, пока этот промис не будет выполнен. Если промис отклоняется, поток получит ошибку. -
cancel(reason): Вызывается, когда потребитель потока отменяет поток.
const readableStream = new ReadableStream({
start(controller) {
/* … */
},
pull(controller) {
/* … */
},
cancel(reason) {
/* … */
},
});
Контроллер ReadableStreamDefaultController поддерживает следующие методы:
-
ReadableStreamDefaultController.close()закрывает связанный поток. -
ReadableStreamDefaultController.enqueue()добавляет заданный фрагмент в соответствующий поток. - Вызов
ReadableStreamDefaultController.error()приводит к ошибке при любом последующем взаимодействии с соответствующим потоком.
/* … */
start(controller) {
controller.enqueue('The first chunk!');
},
/* … */
queuingStrategy
Второй, также необязательный, аргумент конструктора ReadableStream() — это queuingStrategy . Это объект, который опционально определяет стратегию постановки в очередь для потока и принимает два параметра:
-
highWaterMark: Неотрицательное число, указывающее на максимальный уровень воды в ручье при использовании данной стратегии организации очередей. -
size(chunk): Функция, которая вычисляет и возвращает конечный неотрицательный размер заданного фрагмента. Результат используется для определения обратного давления, которое проявляется через соответствующее свойствоReadableStreamDefaultController.desiredSize. Он также определяет, когда вызывается методpull()базового источника.
const readableStream = new ReadableStream({
/* … */
},
{
highWaterMark: 10,
size(chunk) {
return chunk.length;
},
},
);
Методы getReader() и read()
Для чтения из читаемого потока вам потребуется объект Reader, который будет представлять собой ReadableStreamDefaultReader . Метод getReader() интерфейса ReadableStream создает объект Reader и блокирует поток для него. Пока поток заблокирован, получить доступ к другому объекту Reader невозможно, пока этот объект не будет освобожден.
Метод read() ` интерфейса ReadableStreamDefaultReader возвращает промис, предоставляющий доступ к следующему фрагменту во внутренней очереди потока. Он выполняет или отклоняет промис с результатом в зависимости от состояния потока. Возможны следующие варианты:
- Если фрагмент кода доступен, обещание будет выполнено с помощью объекта следующего вида:
{ value: chunk, done: false }. - Если поток прервётся, обещание будет выполнено с помощью объекта следующего вида:
{ value: undefined, done: true }. - Если в потоке возникнет ошибка, обещание будет отклонено с указанием соответствующей ошибки.
const reader = readableStream.getReader();
while (true) {
const { done, value } = await reader.read();
if (done) {
console.log('The stream is done.');
break;
}
console.log('Just read a chunk:', value);
}
locked помещение
Проверить, заблокирован ли поток для чтения, можно, обратившись к его свойству ReadableStream.locked .
const locked = readableStream.locked;
console.log(`The stream is ${locked ? 'indeed' : 'not'} locked.`);
Примеры читаемого потокового кода
Следующий пример кода демонстрирует все шаги в действии. Во-первых, создайте объект ReadableStream , который в своем аргументе underlyingSource (то есть, классе TimestampSource ) определяет метод start() . Этот метод указывает controller потока добавлять метку времени enqueue() каждую секунду в течение десяти секунд. Наконец, он указывает контроллеру close() . Вы используете этот поток, создавая объект Reader с помощью метода getReader() и вызывая read() до тех пор, пока поток не будет done .
class TimestampSource {
#interval
start(controller) {
this.#interval = setInterval(() => {
const string = new Date().toLocaleTimeString();
// Add the string to the stream.
controller.enqueue(string);
console.log(`Enqueued ${string}`);
}, 1_000);
setTimeout(() => {
clearInterval(this.#interval);
// Close the stream after 10s.
controller.close();
}, 10_000);
}
cancel() {
// This is called if the reader cancels.
clearInterval(this.#interval);
}
}
const stream = new ReadableStream(new TimestampSource());
async function concatStringStream(stream) {
let result = '';
const reader = stream.getReader();
while (true) {
// The `read()` method returns a promise that
// resolves when a value has been received.
const { done, value } = await reader.read();
// Result objects contain two properties:
// `done` - `true` if the stream has already given you all its data.
// `value` - Some data. Always `undefined` when `done` is `true`.
if (done) return result;
result += value;
console.log(`Read ${result.length} characters so far`);
console.log(`Most recently read chunk: ${value}`);
}
}
concatStringStream(stream).then((result) => console.log('Stream complete', result));
Асинхронная итерация
Проверка done потока на каждой итерации цикла read() может быть не самым удобным API. К счастью, скоро появится лучший способ сделать это: асинхронная итерация.
for await (const chunk of stream) {
console.log(chunk);
}
Сегодня обходным путем для использования асинхронной итерации является реализация этого поведения с помощью полифилла.
if (!ReadableStream.prototype[Symbol.asyncIterator]) {
ReadableStream.prototype[Symbol.asyncIterator] = async function* () {
const reader = this.getReader();
try {
while (true) {
const {done, value} = await reader.read();
if (done) {
return;
}
yield value;
}
}
finally {
reader.releaseLock();
}
}
}
Создание читаемого потока
Метод ` tee() ` интерфейса ReadableStream перехватывает текущий читаемый поток, возвращая массив из двух элементов, содержащий две результирующие ветви в виде новых экземпляров ReadableStream . Это позволяет двум читателям одновременно читать поток. Например, это можно сделать в сервис-воркере, если вы хотите получить ответ от сервера и передать его в браузер, а также передать в кэш сервис-воркера. Поскольку тело ответа нельзя обработать более одного раза, для этого требуется две копии. Чтобы отменить поток, необходимо отменить обе результирующие ветви. Перехват потока обычно блокирует его на время его обработки, предотвращая блокировку другими читателями.
const readableStream = new ReadableStream({
start(controller) {
// Called by constructor.
console.log('[start]');
controller.enqueue('a');
controller.enqueue('b');
controller.enqueue('c');
},
pull(controller) {
// Called `read()` when the controller's queue is empty.
console.log('[pull]');
controller.enqueue('d');
controller.close();
},
cancel(reason) {
// Called when the stream is canceled.
console.log('[cancel]', reason);
},
});
// Create two `ReadableStream`s.
const [streamA, streamB] = readableStream.tee();
// Read streamA iteratively one by one. Typically, you
// would not do it this way, but you certainly can.
const readerA = streamA.getReader();
console.log('[A]', await readerA.read()); //=> {value: "a", done: false}
console.log('[A]', await readerA.read()); //=> {value: "b", done: false}
console.log('[A]', await readerA.read()); //=> {value: "c", done: false}
console.log('[A]', await readerA.read()); //=> {value: "d", done: false}
console.log('[A]', await readerA.read()); //=> {value: undefined, done: true}
// Read streamB in a loop. This is the more common way
// to read data from the stream.
const readerB = streamB.getReader();
while (true) {
const result = await readerB.read();
if (result.done) break;
console.log('[B]', result);
}
Читаемые потоки байтов
Для потоков, представляющих байты, предоставляется расширенная версия читаемого потока для эффективной обработки байтов, в частности, за счет минимизации копирований. Потоки байтов позволяют использовать считыватели с собственным буфером (BYOB). Реализация по умолчанию может предоставлять различные выходные данные, такие как строки или буферы массивов в случае WebSocket, тогда как потоки байтов гарантируют вывод байтов. Кроме того, считыватели BYOB обладают преимуществами в плане стабильности. Это связано с тем, что если буфер отключается, это гарантирует, что запись в один и тот же буфер не будет произведена дважды, тем самым избегая состояний гонки. Считыватели BYOB могут сократить количество запусков сборки мусора браузером, поскольку они могут повторно использовать буферы.
Создание читаемого потока байтов
Вы можете создать читаемый поток байтов, передав дополнительный параметр type в конструктор ReadableStream() .
new ReadableStream({ type: 'bytes' });
underlyingSource
Исходному потоку читаемых байтов передается объект ReadableByteStreamController для управления. Метод ReadableByteStreamController.enqueue() принимает аргумент chunk , значением которого является ArrayBufferView . Свойство ReadableByteStreamController.byobRequest возвращает текущий запрос на слияние BYOB или null, если такового нет. Наконец, свойство ReadableByteStreamController.desiredSize возвращает желаемый размер для заполнения внутренней очереди управляемого потока.
queuingStrategy
Второй, также необязательный, аргумент конструктора ReadableStream() — это queuingStrategy . Это объект, который опционально определяет стратегию постановки в очередь для потока и принимает один параметр:
-
highWaterMark: Неотрицательное число байтов, указывающее на верхнюю границу потока, использующего данную стратегию организации очереди. Это используется для определения обратного давления, проявляющегося через соответствующее свойствоReadableByteStreamController.desiredSize. Оно также определяет, когда вызывается методpull()базового источника.
Методы getReader() и read()
Затем вы можете получить доступ к ReadableStreamBYOBReader установив соответствующий параметр mode : ReadableStream.getReader({ mode: "byob" }) . Это позволяет более точно контролировать выделение буфера, чтобы избежать копирования. Для чтения из байтового потока необходимо вызвать ReadableStreamBYOBReader.read(view) , где view — это ArrayBufferView .
Пример читаемого кода потока байтов
const reader = readableStream.getReader({ mode: "byob" });
let startingAB = new ArrayBuffer(1_024);
const buffer = await readInto(startingAB);
console.log("The first 1024 bytes, or less:", buffer);
async function readInto(buffer) {
let offset = 0;
while (offset < buffer.byteLength) {
const { value: view, done } =
await reader.read(new Uint8Array(buffer, offset, buffer.byteLength - offset));
buffer = view.buffer;
if (done) {
break;
}
offset += view.byteLength;
}
return buffer;
}
Следующая функция возвращает читаемые потоки байтов, позволяющие эффективно считывать случайно сгенерированный массив без копирования. Вместо использования заранее заданного размера блока в 1024 байта, она пытается заполнить предоставленный разработчиком буфер, обеспечивая полный контроль.
const DEFAULT_CHUNK_SIZE = 1_024;
function makeReadableByteStream() {
return new ReadableStream({
type: 'bytes',
pull(controller) {
// Even when the consumer is using the default reader,
// the auto-allocation feature allocates a buffer and
// passes it to us via `byobRequest`.
const view = controller.byobRequest.view;
view = crypto.getRandomValues(view);
controller.byobRequest.respond(view.byteLength);
},
autoAllocateChunkSize: DEFAULT_CHUNK_SIZE,
});
}
Механика потока записи
Записываемый поток — это место назначения, куда можно записывать данные, представленное в JavaScript объектом WritableStream . Он служит абстракцией над нижележащим приемником — низкоуровневым приемником ввода-вывода, в который записываются необработанные данные.
Данные записываются в поток через записывающее устройство (Writer) , по одному фрагменту за раз. Фрагмент может принимать множество форм, как и фрагменты в считывающем устройстве (Reader). Для создания фрагментов, готовых к записи, можно использовать любой код; записывающее устройство и связанный с ним код называются производителем (Producer ).
Когда создается объект записи и он начинает запись в поток ( активный объект записи ), говорят, что он заблокирован для него. Только один объект записи может записывать в записываемый поток одновременно. Если вы хотите, чтобы другой объект записи начал запись в ваш поток, вам обычно нужно освободить его, прежде чем вы сможете подключить к нему другой объект записи.
Внутренняя очередь отслеживает фрагменты данных, записанные в поток, но еще не обработанные нижележащим приемником.
Стратегия организации очереди — это объект, определяющий, как поток должен сигнализировать об обратном давлении в зависимости от состояния его внутренней очереди. Стратегия организации очереди присваивает размер каждому фрагменту и сравнивает общий размер всех фрагментов в очереди с заданным числом, известным как « верхняя граница» .
Итоговая конструкция называется контроллером . Каждому записываемому потоку соответствует контроллер, позволяющий управлять потоком (например, прерывать его).
Создание потока, доступного для записи.
Интерфейс WritableStream API Streams предоставляет стандартную абстракцию для записи потоковых данных в место назначения, известное как приемник. Этот объект имеет встроенные механизмы обратного давления и организации очередей. Вы создаете записываемый поток, вызывая его конструктор WritableStream() . Он имеет необязательный параметр underlyingSink , который представляет собой объект с методами и свойствами, определяющими поведение созданного экземпляра потока.
underlyingSink
В underlyingSink могут входить следующие необязательные методы, определяемые разработчиком. Параметром controller , передаваемым некоторым из этих методов, является объект WritableStreamDefaultController .
-
start(controller): Этот метод вызывается сразу после создания объекта. Содержимое этого метода должно быть направлено на получение доступа к базовому приемнику. Если этот процесс должен выполняться асинхронно, он может возвращать промис для сигнализации успеха или неудачи. -
write(chunk, controller): Этот метод будет вызван, когда новый фрагмент данных (указанный в параметреchunk) будет готов к записи в базовый приемник. Он может возвращать промис, сигнализирующий об успехе или неудаче операции записи. Этот метод будет вызван только после успешного завершения предыдущих операций записи и никогда после закрытия или прерывания потока. -
close(controller): Этот метод будет вызван, если приложение сообщит о завершении записи фрагментов в поток. Содержимое должно выполнить все необходимые действия для завершения записи в базовый приемник и освобождения доступа к нему. Если этот процесс асинхронный, он может вернуть промис для сигнализации успеха или неудачи. Этот метод будет вызван только после того, как все записи в очередь будут успешно завершены. -
abort(reason): Этот метод будет вызван, если приложение сообщит о желании внезапно закрыть поток и перевести его в состояние ошибки. Он может очистить все удерживаемые ресурсы, подобноclose(), ноabort()будет вызван даже если операции записи поставлены в очередь. Эти фрагменты будут отброшены. Если этот процесс асинхронный, он может вернуть промис для сигнализации успеха или неудачи. ПараметрreasonсодержитDOMStringописывающий причину прерывания потока.
const writableStream = new WritableStream({
start(controller) {
/* … */
},
write(chunk, controller) {
/* … */
},
close(controller) {
/* … */
},
abort(reason) {
/* … */
},
});
Интерфейс WritableStreamDefaultController API Streams представляет собой контроллер, позволяющий управлять состоянием WritableStream во время его инициализации, по мере отправки новых фрагментов на запись или по завершении записи. При создании WritableStream базовому приемнику передается соответствующий экземпляр WritableStreamDefaultController для управления. WritableStreamDefaultController имеет только один метод: WritableStreamDefaultController.error() , который приводит к ошибке при любых последующих взаимодействиях с соответствующим потоком. WritableStreamDefaultController также поддерживает свойство signal , которое возвращает экземпляр AbortSignal , позволяющий остановить операцию WritableStream при необходимости.
/* … */
write(chunk, controller) {
try {
// Try to do something dangerous with `chunk`.
} catch (error) {
controller.error(error.message);
}
},
/* … */
queuingStrategy
Второй, также необязательный, аргумент конструктора WritableStream() — это queuingStrategy . Это объект, который опционально определяет стратегию постановки в очередь для потока и принимает два параметра:
-
highWaterMark: Неотрицательное число, указывающее на максимальный уровень воды в ручье при использовании данной стратегии организации очередей. -
size(chunk): Функция, которая вычисляет и возвращает конечный неотрицательный размер заданного фрагмента. Результат используется для определения обратного давления, которое проявляется через соответствующее свойствоWritableStreamDefaultWriter.desiredSize.
Методы getWriter() и write()
Для записи в записываемый поток необходим объект типа WritableStreamDefaultWriter . Метод getWriter() интерфейса WritableStream возвращает новый экземпляр WritableStreamDefaultWriter и блокирует поток для этого экземпляра. Пока поток заблокирован, получить доступ к другому объекту записи невозможно, пока текущий не будет освобожден.
Метод write() интерфейса WritableStreamDefaultWriter записывает переданный фрагмент данных в WritableStream и его базовый приемник, а затем возвращает промис, который разрешается, указывая на успех или неудачу операции записи. Следует отметить, что значение «успеха» зависит от базового приемника; это может означать, что фрагмент данных принят, а не обязательно, что он безопасно сохранен в конечном пункте назначения.
const writer = writableStream.getWriter();
const resultPromise = writer.write('The first chunk!');
locked помещение
Проверить, заблокирован ли записываемый поток, можно, обратившись к его свойству WritableStream.locked .
const locked = writableStream.locked;
console.log(`The stream is ${locked ? 'indeed' : 'not'} locked.`);
Пример кода для записи потока
Приведённый ниже пример кода демонстрирует все этапы в действии.
const writableStream = new WritableStream({
start(controller) {
console.log('[start]');
},
async write(chunk, controller) {
console.log('[write]', chunk);
// Wait for next write.
await new Promise((resolve) => setTimeout(() => {
document.body.textContent += chunk;
resolve();
}, 1_000));
},
close(controller) {
console.log('[close]');
},
abort(reason) {
console.log('[abort]', reason);
},
});
const writer = writableStream.getWriter();
const start = Date.now();
for (const char of 'abcdefghijklmnopqrstuvwxyz') {
// Wait to add to the write queue.
await writer.ready;
console.log('[ready]', Date.now() - start, 'ms');
// The Promise is resolved after the write finishes.
writer.write(char);
}
await writer.close();
Передача потока данных, доступных для чтения, в поток данных, доступных для записи.
Поток, доступный для чтения, можно перенаправить в поток, доступный для записи, с помощью метода pipeTo() потока, доступного для чтения. Метод ReadableStream.pipeTo() перенаправляет текущий поток ReadableStream в заданный WritableStream и возвращает промис, который выполняется при успешном завершении процесса перенаправления или отклоняется в случае возникновения ошибок.
const readableStream = new ReadableStream({
start(controller) {
// Called by constructor.
console.log('[start readable]');
controller.enqueue('a');
controller.enqueue('b');
controller.enqueue('c');
},
pull(controller) {
// Called when controller's queue is empty.
console.log('[pull]');
controller.enqueue('d');
controller.close();
},
cancel(reason) {
// Called when the stream is canceled.
console.log('[cancel]', reason);
},
});
const writableStream = new WritableStream({
start(controller) {
// Called by constructor
console.log('[start writable]');
},
async write(chunk, controller) {
// Called upon writer.write()
console.log('[write]', chunk);
// Wait for next write.
await new Promise((resolve) => setTimeout(() => {
document.body.textContent += chunk;
resolve();
}, 1_000));
},
close(controller) {
console.log('[close]');
},
abort(reason) {
console.log('[abort]', reason);
},
});
await readableStream.pipeTo(writableStream);
console.log('[finished]');
Создание потока преобразования
Интерфейс TransformStream API Streams представляет собой набор преобразуемых данных. Вы создаете поток преобразования, вызывая его конструктор TransformStream() , который создает и возвращает объект потока преобразования из заданных обработчиков. Конструктор TransformStream() принимает в качестве первого аргумента необязательный объект JavaScript, представляющий transformer . Такие объекты могут содержать любой из следующих методов:
transformer
-
start(controller): Этот метод вызывается сразу после создания объекта. Обычно он используется для добавления префиксных фрагментов в очередь с помощьюcontroller.enqueue(). Эти фрагменты будут считываться со стороны чтения, но не зависят от каких-либо записей на стороне записи. Если этот начальный процесс асинхронный, например, потому что требуется некоторое усилие для получения префиксных фрагментов, функция может вернуть промис для сигнализации успеха или неудачи; отклоненный промис вызовет ошибку в потоке. Любые выброшенные исключения будут повторно выброшены конструкторомTransformStream(). -
transform(chunk, controller): Этот метод вызывается, когда новый фрагмент, первоначально записанный на записываемую сторону, готов к преобразованию. Реализация потока гарантирует, что эта функция будет вызвана только после успешного завершения предыдущих преобразований и никогда до завершенияstart()или после вызоваflush(). Эта функция выполняет фактическую работу по преобразованию потока. Она может добавить результаты в очередь с помощьюcontroller.enqueue(). Это позволяет одному фрагменту, записанному на записываемую сторону, привести к появлению нуля или нескольких фрагментов на читаемой стороне, в зависимости от того, сколько раз вызываетсяcontroller.enqueue(). Если процесс преобразования асинхронный, эта функция может возвращать промис для сигнализации об успехе или неудаче преобразования. Отклоненный промис вызовет ошибку как на читаемой, так и на записываемой сторонах потока преобразования. Если методtransform()не указан, используется тождественное преобразование, которое добавляет в очередь неизмененные фрагменты с записываемой стороны на читаемую сторону. -
flush(controller): Этот метод вызывается после того, как все фрагменты, записанные на записываемую сторону, были успешно преобразованы с помощьюtransform(), и записываемая сторона вот-вот будет закрыта. Обычно это используется для добавления фрагментов с суффиксом на читаемую сторону до того, как она также будет закрыта. Если процесс сброса является асинхронным, функция может возвращать промис для сигнализации успеха или неудачи; результат будет передан вызывающей сторонеstream.writable.write(). Кроме того, отклоненный промис вызовет ошибку как на читаемой, так и на записываемой сторонах потока. Выброс исключения обрабатывается так же, как и возврат отклоненного промиса.
const transformStream = new TransformStream({
start(controller) {
/* … */
},
transform(chunk, controller) {
/* … */
},
flush(controller) {
/* … */
},
});
Стратегии организации очередей writableStrategy и readableStrategy
Вторым и третьим необязательными параметрами конструктора TransformStream() являются необязательные стратегии организации очередей writableStrategy и readableStrategy . Они определены в соответствии с разделами, посвященными потокам чтения и записи соответственно.
Пример кода преобразования потока
Приведённый ниже пример кода демонстрирует работу потока преобразования.
// Note that `TextEncoderStream` and `TextDecoderStream` exist now.
// This example shows how you would have done it before.
const textEncoderStream = new TransformStream({
transform(chunk, controller) {
console.log('[transform]', chunk);
controller.enqueue(new TextEncoder().encode(chunk));
},
flush(controller) {
console.log('[flush]');
controller.terminate();
},
});
(async () => {
const readStream = textEncoderStream.readable;
const writeStream = textEncoderStream.writable;
const writer = writeStream.getWriter();
for (const char of 'abc') {
writer.write(char);
}
writer.close();
const reader = readStream.getReader();
for (let result = await reader.read(); !result.done; result = await reader.read()) {
console.log('[value]', result.value);
}
})();
Пропуск читаемого потока через поток преобразования
Метод pipeThrough() интерфейса ReadableStream предоставляет возможность последовательной передачи текущего потока через поток преобразования или любую другую пару «запись/чтение». Передача потока через канал обычно блокирует его на время передачи, предотвращая блокировку другими читателями.
const transformStream = new TransformStream({
transform(chunk, controller) {
console.log('[transform]', chunk);
controller.enqueue(new TextEncoder().encode(chunk));
},
flush(controller) {
console.log('[flush]');
controller.terminate();
},
});
const readableStream = new ReadableStream({
start(controller) {
// called by constructor
console.log('[start]');
controller.enqueue('a');
controller.enqueue('b');
controller.enqueue('c');
},
pull(controller) {
// called read when controller's queue is empty
console.log('[pull]');
controller.enqueue('d');
controller.close(); // or controller.error();
},
cancel(reason) {
// called when rs.cancel(reason)
console.log('[cancel]', reason);
},
});
(async () => {
const reader = readableStream.pipeThrough(transformStream).getReader();
for (let result = await reader.read(); !result.done; result = await reader.read()) {
console.log('[value]', result.value);
}
})();
Следующий пример кода (немного надуманный) показывает, как можно реализовать «кричащую» версию функции fetch() , которая преобразует весь текст в верхний регистр, обрабатывая возвращаемый промис ответа как поток и преобразуя его по частям. Преимущество такого подхода заключается в том, что вам не нужно ждать загрузки всего документа, что может иметь огромное значение при работе с большими файлами.
function upperCaseStream() {
return new TransformStream({
transform(chunk, controller) {
controller.enqueue(chunk.toUpperCase());
},
});
}
function appendToDOMStream(el) {
return new WritableStream({
write(chunk) {
el.append(chunk);
}
});
}
fetch('./lorem-ipsum.txt').then((response) =>
response.body
.pipeThrough(new TextDecoderStream())
.pipeThrough(upperCaseStream())
.pipeTo(appendToDOMStream(document.body))
);
Демо
Приведённая ниже демонстрация показывает работу потоков, доступных для чтения, записи и преобразования. Она также включает примеры цепочек вызовов pipeThrough() и pipeTo() , а также демонстрирует работу функции tee() . При желании вы можете запустить демонстрацию в отдельном окне или просмотреть исходный код .
Полезные потоки доступны в браузере.
В браузере встроен ряд полезных потоков. Вы можете легко создать ReadableStream из объекта Blob. Метод stream() интерфейса Blob возвращает объект ReadableStream , который при чтении возвращает данные, содержащиеся в объекте Blob. Также помните, что объект File — это особый тип объекта Blob , и его можно использовать в любом контексте, в котором используется объект Blob.
const readableStream = new Blob(['hello world'], { type: 'text/plain' }).stream();
Потоковые варианты функций TextDecoder.decode() и TextEncoder.encode() называются соответственно TextDecoderStream и TextEncoderStream .
const response = await fetch('https://streams.spec.whatwg.org/');
const decodedStream = response.body.pipeThrough(new TextDecoderStream());
Сжатие или распаковка файла легко осуществляется с помощью потоков преобразования CompressionStream и DecompressionStream соответственно. Приведенный ниже пример кода показывает, как можно загрузить спецификацию Streams, сжать файл (gzip) прямо в браузере и записать сжатый файл непосредственно на диск.
const response = await fetch('https://streams.spec.whatwg.org/');
const readableStream = response.body;
const compressedStream = readableStream.pipeThrough(new CompressionStream('gzip'));
const fileHandle = await showSaveFilePicker();
const writableStream = await fileHandle.createWritable();
compressedStream.pipeTo(writableStream);
Примерами потоков, доступных для записи, являются FileSystemWritableFileStream из File System Access API и экспериментальные потоки запросов fetch() .
API последовательного порта активно использует как потоки для чтения, так и для записи.
// Prompt user to select any serial port.
const port = await navigator.serial.requestPort();
// Wait for the serial port to open.
await port.open({ baudRate: 9_600 });
const reader = port.readable.getReader();
// Listen to data coming from the serial device.
while (true) {
const { value, done } = await reader.read();
if (done) {
// Allow the serial port to be closed later.
reader.releaseLock();
break;
}
// value is a Uint8Array.
console.log(value);
}
// Write to the serial port.
const writer = port.writable.getWriter();
const data = new Uint8Array([104, 101, 108, 108, 111]); // hello
await writer.write(data);
// Allow the serial port to be closed later.
writer.releaseLock();
Наконец, API WebSocketStream интегрирует потоки данных с API WebSocket.
const wss = new WebSocketStream(WSS_URL);
const { readable, writable } = await wss.connection;
const reader = readable.getReader();
const writer = writable.getWriter();
while (true) {
const { value, done } = await reader.read();
if (done) {
break;
}
const result = await process(value);
await writer.write(result);
}
Полезные ресурсы
- Спецификация потоков
- Сопутствующие демонстрации
- Потоки полифила
- 2016 год — год веб-трансляций
- Асинхронные итераторы и генераторы
- Визуализатор потока
Благодарности
Данная статья была рецензирована Джейком Арчибальдом , Франсуа Бофором , Сэмом Даттоном , Маттиасом Бюленсом , Сурмой , Джо Медли и Адамом Райсом . Посты Джейка Арчибальда в блоге очень помогли мне в понимании потоков данных. Некоторые примеры кода вдохновлены исследованиями пользователя GitHub @bellbind , а часть текста в значительной степени основана на веб-документации MDN по потокам данных . Авторы стандарта Streams проделали огромную работу по написанию этой спецификации.