Spec-Zone.ru › Web APIs

Использование потоков чтения

В качестве разработчика JavaScript программирование чтения и обработки потоков данных, получаемых по сети по частям, очень полезно! Но как использовать функциональность потока чтения API потоков? В этой статье объясняются основы.

Примечание: Эта статья предполагает, что вы понимаете случаи использования потоков чтения и знакомы с основными концепциями. Если нет, рекомендуем сначала ознакомиться с обзором концепций и использования потоков и специальной статьей концепций API потоков, а затем вернуться.

Примечание: Если вам нужна информация о потоках записи, обратитесь к использованию потоков записи вместо этого.

Поиск примеров

В этой статье мы рассмотрим различные примеры, взятые из нашего репозитория dom-examples/streams. Там вы найдете полный исходный код, а также ссылки на примеры.

Использование fetch как потока

API Fetch позволяет получать ресурсы по сети, предоставляя современную альтернативу XHR. У него есть ряд преимуществ, и очень удобно, что браузеры недавно добавили возможность использовать ответ fetch как поток чтения.

Доступны свойства Request.body и Response.body, являющиеся геттерами, которые экспонируют содержимое тела как поток чтения.

Как показывает наш пример простого насоса потока (посмотреть в действии), его экспонирование сводится к доступу к свойству body ответа:

// Fetch the original image
fetch("./tortoise.png")
  // Retrieve its body as ReadableStream
  .then((response) => response.body);

Это предоставляет нам объект ReadableStream.

Присоединение ридера

Теперь, когда у нас есть тело потока, чтение потока требует присоединения ридера к нему. Это делается с помощью метода ReadableStream.getReader():

// Fetch the original image
fetch("./tortoise.png")
  // Retrieve its body as ReadableStream
  .then((response) => response.body)
  .then((body) => {
    const reader = body.getReader();
    // …
  });

Вызов этого метода создает ридер и блокирует его для потока — никакой другой ридер не может читать этот поток, пока этот ридер не будет освобожден, например, путем вызова ReadableStreamDefaultReader.releaseLock().

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

// Fetch the original image
fetch("./tortoise.png")
  // Retrieve its body as ReadableStream
  .then((response) => {
    const reader = response.body.getReader();
    // …
  });

Чтение потока

Теперь, когда у вас есть подключенный ридер, вы можете читать фрагменты данных из потока с помощью метода ReadableStreamDefaultReader.read(). Это считывает один фрагмент из потока, с которым вы можете сделать все, что захотите. Например, наш пример простого насоса потока продолжает добавлять каждый фрагмент в новый пользовательский ReadableStream (мы узнаем больше об этом в следующем разделе), а затем создает новый Response из него, потребляет его как Blob, создает URL-адрес объекта из этого blob с помощью URL.createObjectURL() и затем отображает его на экране в элементе <img>, эффективно создавая копию изображения, которое мы изначально получили.

// Fetch the original image
fetch("./tortoise.png")
  // Retrieve its body as ReadableStream
  .then((response) => {
    const reader = response.body.getReader();
    return new ReadableStream({
      start(controller) {
        return pump();
        function pump() {
          return reader.read().then(({ done, value }) => {
            // When no more data needs to be consumed, close the stream
            if (done) {
              controller.close();
              return;
            }
            // Enqueue the next data chunk into our target stream
            controller.enqueue(value);
            return pump();
          });
        }
      },
    });
  })
  // Create a new response out of the stream
  .then((stream) => new Response(stream))
  // Create an object URL for the response
  .then((response) => response.blob())
  .then((blob) => URL.createObjectURL(blob))
  // Update image
  .then((url) => console.log((image.src = url)))
  .catch((err) => console.error(err));

Давайте подробно рассмотрим, как используется read(). В функции pump() выше мы сначала вызываем read(), которая возвращает промис, содержащий объект результатов — в нем содержатся результаты нашего чтения в виде { done, value }:

reader.read().then(({ done, value }) => {
  /* … */
});

Результаты могут быть одного из трех типов:

  • Если фрагмент доступен для чтения, промис будет выполнен с объектом типа { value: theChunk, done: false }.
  • Если поток закрывается, промис будет выполнен с объектом типа { value: undefined, done: true }.
  • Если в потоке произошла ошибка, промис будет отклонен с соответствующей ошибкой.

Далее, мы проверяем, является ли done true. Если да, больше фрагментов для чтения нет (значение undefined ), поэтому мы возвращаемся из функции и закрываем пользовательский поток с помощью ReadableStreamDefaultController.close():

if (done) {
  controller.close();
  return;
}

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

Если done не является true, мы обрабатываем новый фрагмент, который мы прочитали (он находится в свойстве value объекта результатов), а затем снова вызываем функцию pump() для чтения следующего фрагмента.

// Enqueue the next data chunk into our target stream
controller.enqueue(value);
return pump();

Это стандартный шаблон, который вы увидите при использовании ридеров потоков:

  1. Вы пишете функцию, которая начинается с чтения потока.
  2. Если больше данных потока нет, вы возвращаетесь из функции.
  3. Если данных потока больше, вы обрабатываете текущий фрагмент, а затем снова вызываете функцию.
  4. Вы продолжаете цепочку вызовов функции pump() до тех пор, пока данные потока не закончатся, в этом случае выполняется шаг 2.

Убрав весь код для фактического "насоса", код можно обобщить примерно так:

fetch("http://example.com/somefile.txt")
  // Retrieve its body as ReadableStream
  .then((response) => {
    const reader = response.body.getReader();
    // read() returns a promise that resolves when a value has been received
    reader.read().then(function pump({ done, value }) {
      if (done) {
        // Do something with last chunk of data then exit reader
        return;
      }
      // Otherwise do something here to process current chunk

      // Read some more, and call this function again
      return reader.read().then(pump);
    });
  })
  .catch((err) => console.error(err));

Примечание: Функция выглядит так, как будто pump() вызывает себя, что может привести к потенциальной глубокой рекурсии. Однако, поскольку pump асинхронна, и каждый вызов pump() находится в конце обработчика промиса, это на самом деле аналогично цепочке обработчиков промисов.

Чтение потока еще проще, если он записан с использованием async/await вместо промисов:

async function readData(url) {
  const response = await fetch(url);
  const reader = response.body.getReader();
  while (true) {
    const { done, value } = await reader.read();
    if (done) {
      // Do something with last chunk of data then exit reader
      return;
    }
    // Otherwise do something here to process current chunk
  }
}

Использование fetch() с асинхронной итерацией

Существует еще более простой способ потребления fetch(), а именно итерация возвращаемого response.body с помощью синтаксиса for await...of. Это работает, потому что response.body возвращает ReadableStream, который является асинхронно итерируемым.

Используя этот подход, код примера в предыдущем разделе можно переписать так:

async function readData(url) {
  const response = await fetch(url);
  for await (const chunk of response.body) {
    // Do something with each "chunk"
  }
  // Exit when done
}

Если вы хотите остановить итерацию по потоку, вы можете отменить операцию fetch() с помощью AbortController и связанного с ним AbortSignal:

const aborter = new AbortController();
button.addEventListener("click", () => aborter.abort());
logChunks("http://example.com/somefile.txt", { signal: aborter.signal });

async function logChunks(url, { signal }) {
  const response = await fetch(url, { signal });
  for await (const chunk of response.body) {
    // Do something with the chunk
  }
}

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

const aborter = new AbortController();
button.addEventListener("click", () => aborter.abort());
logChunks("http://example.com/somefile.txt", { signal: aborter.signal });

async function logChunks(url, { signal }) {
  const response = await fetch(url);
  for await (const chunk of response.body) {
    if (signal.aborted) break; // just break out of loop
    // Do something with the chunk
  }
}

Пример асинхронного ридера

Код ниже демонстрирует более полный пример. Здесь поток fetch потребляется с помощью итератора внутри блока try/catch. На каждой итерации цикла код просто регистрирует и подсчитывает полученные байты. Если возникнет ошибка, он зарегистрирует проблему. Операция fetch() может быть отменена с помощью AbortSignal, что также будет зарегистрировано как ошибка.

let bytes = 0;

const aborter = new AbortController();
button.addEventListener("click", () => aborter.abort());
logChunks("http://example.com/somefile.txt", { signal: aborter.signal });

async function logChunks(url, { signal }) {
  try {
    const response = await fetch(url, signal);
    for await (const chunk of response.body) {
      if (signal.aborted) throw signal.reason;
      bytes += chunk.length;
      logConsumer(`Chunk: ${chunk}. Read ${bytes} characters.`);
    }
  } catch (e) {
    if (e instanceof TypeError) {
      console.log(e);
      logConsumer("TypeError: Browser may not support async iteration");
    } else {
      logConsumer(`Error in async iterator: ${e}.`);
    }
  }
}

Ниже показан пример лога выполнения кода или сообщения о том, что ваш браузер не поддерживает асинхронную итерацию ReadableStream. Правая сторона показывает полученные фрагменты; вы можете нажать кнопку "отмена", чтобы остановить fetch.

Примечание: Этот оператор fetch мокируется для демонстрации и просто возвращает ReadableStream, который генерирует случайные текстовые фрагменты. "Основной источник" слева ниже — данные, генерируемые в смокированном источнике, а столбец справа — регистрация от потребителя. (Код для смокированного источника не отображается, так как он не относится к примеру.)

Создание собственного пользовательского потока чтения

В примере простого насоса потока, который мы изучали в этой статье, есть вторая часть — как только мы прочитали изображение из тела fetch по частям, мы добавили их в другой, пользовательский поток, созданный нами самими. Как это сделать? Конструктор ReadableStream().

Конструктор ReadableStream()

Легко читать из потока, когда браузер предоставляет его, как в случае с Fetch, но иногда вам нужно создать пользовательский поток и заполнить его своими частями. Конструктор ReadableStream() позволяет сделать это с помощью синтаксиса, который на первый взгляд кажется сложным, но на самом деле не так уж и плох.

Общая структура синтаксиса выглядит так:

const stream = new ReadableStream(
  {
    start(controller) {},
    pull(controller) {},
    cancel() {},
    type,
    autoAllocateChunkSize,
  },
  {
    highWaterMark: 3,
    size: () => 1,
  },
);

Конструктор принимает два объекта в качестве параметров. Первый объект обязателен и создаёт модель в JavaScript базового источника, из которого читаются данные. Второй объект необязателен и позволяет указать пользовательскую стратегию очереди для вашего потока. Вам редко придётся это делать, поэтому пока сосредоточимся на первом.

Первый объект может содержать до пяти элементов, из которых только первый является обязательным:

  1. start(controller) — Метод, вызываемый один раз сразу после того, как конструктор ReadableStream был создан. В этом методе необходимо включить код, который настраивает функциональность потока, например, начинает генерировать данные или иначе получает доступ к источнику.
  2. pull(controller) — Метод, который, если включён, вызывается повторно до тех пор, пока внутренняя очередь потока не заполнится. Это можно использовать для управления потоком по мере добавления новых частей в очередь.
  3. cancel() — Метод, который, если включён, будет вызван, если приложение сигнализирует о необходимости отмены потока (например, если вызван ReadableStream.cancel()). Содержимое должно выполнить все необходимые действия для освобождения доступа к источнику потока.
  4. type и autoAllocateChunkSize — Эти элементы (если включены) указывают, что поток должен быть байтовым потоком. Байтовые потоки рассматриваются отдельно в Использовании потоков чтения байтовых данных, так как их назначение и использование несколько отличаются от обычных (по умолчанию) потоков.

Если снова взглянуть на наш простой пример кода, можно увидеть, что наш конструктор ReadableStream() включает только один метод — start(), который служит для считывания всех данных из нашего потока fetch.

// Fetch the original image
fetch("./tortoise.png")
  // Retrieve its body as ReadableStream
  .then((response) => {
    const reader = response.body.getReader();
    return new ReadableStream({
      start(controller) {
        return pump();
        function pump() {
          return reader.read().then(({ done, value }) => {
            // When no more data needs to be consumed, close the stream
            if (done) {
              controller.close();
              return;
            }
            // Enqueue the next data chunk into our target stream
            controller.enqueue(value);
            return pump();
          });
        }
      },
    });
  });

Контроллеры ReadableStream

Вы заметите, что методы start() и pull(), переданные в конструктор ReadableStream(), получают параметры controller — это экземпляры класса ReadableStreamDefaultController, которые можно использовать для управления потоком.

В нашем примере мы используем метод контроллера enqueue() для помещения значения в пользовательский поток после его чтения из тела fetch.

Кроме того, когда мы закончим чтение тела fetch, мы используем метод контроллера close() для закрытия пользовательского потока — все ранее помещённые в очередь части всё ещё можно прочитать из него, но больше ничего нельзя поместить в очередь, и поток закрывается после завершения чтения.

Чтение из пользовательских потоков

В нашем примере с простым насосом для потоков мы потребляем пользовательский поток чтения, передавая его в вызов конструктора Response, после чего потребляем его как blob().

readableStream
  .then((stream) => new Response(stream))
  .then((response) => response.blob())
  .then((blob) => URL.createObjectURL(blob))
  .then((url) => console.log((image.src = url)))
  .catch((err) => console.error(err));

Но пользовательский поток всё ещё является экземпляром ReadableStream , что означает, что вы можете прикрепить к нему читателя. В качестве примера обратите внимание на нашу демонстрацию простого случайного потока (посмотреть в живом режиме), которая создаёт пользовательский поток, помещает в него несколько случайных строк и затем считывает данные из потока, как только нажимается кнопка Остановить генерацию строк.

Примечание: Для потребления потока с использованием FetchEvent.respondWith() содержимое, помещённое в очередь потока, должно быть типа Uint8Array; например, закодированное с использованием TextEncoder.

Конструктор пользовательского потока имеет метод start(), который использует вызов setInterval() для генерации случайной строки каждую секунду. Затем используется ReadableStreamDefaultController.enqueue() для помещения её в поток. Когда нажимается кнопка, интервал отменяется, и вызывается функция readStream(), чтобы снова прочитать данные из потока. Мы также закрываем поток, так как перестали помещать в него части.

let interval;
const stream = new ReadableStream({
  start(controller) {
    interval = setInterval(() => {
      const string = randomChars();
      // Add the string to the stream
      controller.enqueue(string);
      // show it on the screen
      const listItem = document.createElement("li");
      listItem.textContent = string;
      list1.appendChild(listItem);
    }, 1000);
    button.addEventListener("click", () => {
      clearInterval(interval);
      readStream();
      controller.close();
    });
  },
  pull(controller) {
    // We don't really need a pull in this example
  },
  cancel() {
    // This is called if the reader cancels,
    // so we should stop generating strings
    clearInterval(interval);
  },
});

В самой функции readStream() мы фиксируем читателя к потоку с помощью ReadableStream.getReader(), а затем следуем той же схеме, что и раньше — читаем каждую часть с помощью read(), проверяем, является ли done true , и завершаем процесс, если это так, и читаем следующую часть и обрабатываем её, если нет, прежде чем повторно вызывать метод read().

function readStream() {
  const reader = stream.getReader();
  let charsReceived = 0;
  let result = "";

  // read() returns a promise that resolves
  // when a value has been received
  reader.read().then(function processText({ done, value }) {
    // 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) {
      console.log("Stream complete");
      para.textContent = result;
      return;
    }

    charsReceived += value.length;
    const chunk = value;
    const listItem = document.createElement("li");
    listItem.textContent = `Read ${charsReceived} characters so far. Current chunk = ${chunk}`;
    list2.appendChild(listItem);

    result += chunk;

    // Read some more, and call this function again
    return reader.read().then(processText);
  });
}

Закрытие и отмена потоков

Мы уже показали примеры использования ReadableStreamDefaultController.close() для закрытия читателя. Как мы уже говорили, все ранее помещённые в очередь части всё ещё можно прочитать, но больше ничего нельзя поместить в очередь, потому что поток закрыт.

Если вы хотите полностью избавиться от потока и отбросить все помещённые в очередь части, используйте ReadableStream.cancel() или ReadableStreamDefaultReader.cancel().

Разветвление потока

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

Например, вы можете сделать это в ServiceWorker, если хотите получить ответ от сервера и передать его в браузер, а также в кэш Service Worker. Поскольку тело ответа нельзя потреблять более одного раза, а поток нельзя читать более чем одному читателю одновременно, вам потребуются две копии для этого.

Мы предоставляем пример этого в нашей демонстрации разветвления потока (посмотреть в живом режиме). Этот пример работает примерно так же, как наш пример с простым случайным потоком, за исключением того, что когда нажимается кнопка для остановки генерации случайных строк, пользовательский поток берётся и разветвляется, а затем оба полученных потока читаются:

function teeStream() {
  const teedOff = stream.tee();
  readStream(teedOff[0], list2);
  readStream(teedOff[1], list3);
}

Цепочки соединений

Ещё одной особенностью потоков является возможность соединять потоки друг с другом (называемые цепочками соединений). Это включает в себя два метода — ReadableStream.pipeThrough(), который пропускает поток чтения через пару писатель/читатель для преобразования одного формата данных в другой, и ReadableStream.pipeTo(), который пропускает поток чтения к писателю, выступающему в качестве конечной точки цепочки соединений.

У нас есть пример под названием Распаковка блоков PNG (посмотреть в живом режиме), который получает изображение как поток, затем пропускает его через пользовательский поток преобразования PNG, который извлекает блоки PNG из потока бинарных данных.

// Fetch the original image
fetch("png-logo.png")
  // Retrieve its body as ReadableStream
  .then((response) => response.body)
  // Create a gray-scaled PNG stream out of the original
  .then((rs) => logReadableStream("Fetch Response Stream", rs))
  .then((body) => body.pipeThrough(new PNGTransformStream()))
  .then((rs) => logReadableStream("PNG Chunk Stream", rs));

У нас пока нет примера, использующего TransformStream.

Сводка

Это объясняет основы "стандартных" потоков чтения.

См. Использование потоков чтения байтовых данных для получения информации о том, как использовать потоки чтения байтовых данных: потоки с базовым байтовым источником, которые могут выполнять эффективные передачу без копирования к потребителю, минуя внутренние очереди потока.

© 2005–2024 MDN contributors.
Licensed under the Creative Commons Attribution-ShareAlike License v2.5 or later.
https://developer.mozilla.org/en-US/docs/Web/API/Streams_API/Using_readable_streams

Spec-Zone.ru

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