Spec-Zone.ru › Web APIs

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

Потоки байтов для чтения — это потоки для чтения, имеющие базовый источник байтов type: "bytes", и которые поддерживают эффективную передачу данных без копирования из базового источника потребителю (минуя внутренние очереди потока). Они предназначены для случаев, когда данные могут быть предоставлены или запрошены в произвольных, потенциально очень больших, кусках, и, следовательно, избежание копирования, вероятно, повысит эффективность.

В этой статье объясняется, как потоки байтов для чтения отличаются от обычных потоков по умолчанию, и как их создавать и использовать.

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

Обзор

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

Поток байтов для чтения создается путем указания type: "bytes" в объекте underlyingSource, который может передаваться в качестве первого параметра конструктору ReadableStream(). С этим значением поток создается с ReadableByteStreamController, и именно этот объект передается базовому источнику при вызове функций обратного вызова start(controller) и pull(controller).

Основное различие между ReadableByteStreamController и контроллером по умолчанию (ReadableStreamDefaultController) заключается в том, что он имеет дополнительное свойство ReadableByteStreamController.byobRequest типа ReadableStreamBYOBRequest. Оно представляет собой ожидающий запрос на чтение потребителя, который будет выполнен как передача данных без копирования из базового источника. Свойство будет null если нет ожидающего запроса.

Запрос byobRequest доступен только тогда, когда на поток байтов для чтения сделан запрос на чтение, и нет данных во внутренних очередях потока (если данные есть, запрос выполняется из этих очередей).

Базовый источник байтов, которому необходимо передать данные, должен проверить свойство byobRequest и, если оно доступно, использовать его для передачи данных. Если свойство null, входящие данные должны быть добавлены во внутренние очереди потока с помощью ReadableByteStreamController.enqueue() (это единственный способ передачи данных при использовании потока по умолчанию).

В ReadableStreamBYOBRequest есть свойство view, которое представляет собой вид на буфер, выделенный для передачи. Данные из базового источника должны быть записаны в это свойство, а затем базовый источник должен вызвать respond(), указав количество записанных байтов. Это сигнализирует о том, что данные должны быть переданы, и ожидающий запрос на чтение потребителя разрешен. После вызова respond() в view больше нельзя записывать.

Также есть дополнительный метод ReadableStreamBYOBRequest.respondWithNewView(), которому базовый источник может передать «новый» вид, содержащий данные для передачи. Этот новый вид должен находиться в том же буфере памяти, что и оригинальный, и с того же начального смещения. Этот метод может быть использован, если базовому источнику байтов необходимо сначала передать вид в поток для обработки (например) и затем получить его обратно, прежде чем отвечать на byobRequest. В большинстве случаев этот метод не потребуется.

Потоки байтов для чтения обычно читаются с помощью ReadableStreamBYOBReader, который можно получить, вызвав ReadableStream.getReader() на потоке, указав mode: "byob" в параметре options.

Поток байтов для чтения также может быть прочитан с помощью читателя по умолчанию (ReadableStreamDefaultReader), но в этом случае объекты byobRequest создаются только при включении автоматического выделения буфера для потока (autoAllocateChunkSize был установлен для underlyingSource потока). Обратите внимание, что размер, указанный autoAllocateChunkSize, используется для размера буфера в этом случае; для байтового читателя буфер предоставляется потребителем. Если свойство не было указано, читатель по умолчанию по-прежнему «будет работать», но базовый источник никогда не получит byobRequest, и все данные будут передаваться через внутренние очереди потока.

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

Примеры

Базовый источник push с читателем байтов

Этот пример показывает, как создать читаемый байтовый поток с базовым байтовым источником push и прочитать его с помощью байтового ридера.

В отличие от байтового источника pull, данные могут поступать в любое время. Поэтому базовый источник должен использовать controller.byobRequest для передачи поступающих данных, если они есть, а в противном случае добавлять данные в внутренние очереди потока. Кроме того, поскольку данные могут поступать в любое время, поведение мониторинга настраивается в функции обратного вызова underlyingSource.start().

Пример сильно вдохновлен примером байтового источника push в спецификации потока. Он использует эмулированный "гипотетический сокет" источник, который предоставляет данные произвольных размеров. Чтение данных сознательно задерживается в разных точках, чтобы позволить базовому источнику использовать как передачу, так и добавление в очередь для отправки данных в поток. Поддержка обратной загрузки не демонстрируется.

Примечание: Базовый байтовый источник также может использоваться с ридером по умолчанию. Если включена автоматическая выделение буферов, контроллер предоставит буферы фиксированного размера для операций zero-copy при наличии запроса от ридера и пустых внутренних очередей потока. Если автоматическое выделение буферов не включено, все данные из байтового потока всегда добавляются в очередь. Это поведение аналогично поведению, показанному в примерах "pull: базового байтового источника".

Эмулированный базовый источник сокета

Эмулированный базовый источник имеет три важных метода:

  • select2() представляет собой ожидающий запрос в сокете. Он возвращает промис, который разрешается, когда данные становятся доступными.
  • readInto() считывает данные из сокета в предоставленный буфер и затем очищает данные.
  • close() закрывает сокет.

Реализация очень простая. Как показано ниже, select2() создаёт буфер случайного размера с случайными данными по таймауту. Созданные данные считываются в буфер и очищаются в readInto().

class MockHypotheticalSocket {
  constructor() {
    this.max_data = 800; // total amount of data to stream from "socket"
    this.max_per_read = 100; // max data per read
    this.min_per_read = 40; // min data per read
    this.data_read = 0; // total data read so far (capped is maxdata)
    this.socketData = null;
  }

  // Method returning promise when this socket is readable.
  select2() {
    // Object used to resolve promise
    const resultObj = {};
    resultObj["bytesRead"] = 0;

    return new Promise((resolve /*, reject*/) => {
      if (this.data_read >= this.max_data) {
        //out of data
        resolve(resultObj);
        return;
      }

      // Emulate slow read of data
      setTimeout(() => {
        const numberBytesReceived = this.getNumberRandomBytesSocket();
        this.data_read += numberBytesReceived;
        this.socketData = this.randomByteArray(numberBytesReceived);
        resultObj["bytesRead"] = numberBytesReceived;
        resolve(resultObj);
      }, 500);
    });
  }

  /* Read data into specified buffer offset */
  readInto(buffer, offset, length) {
    let dataLength = 0;
    if (this.socketData) {
      dataLength = this.socketData.length;
      const myView = new Uint8Array(buffer, offset, length);
      // Write the length of data specified into buffer
      // Code assumes buffer always bigger than incoming data
      for (let i = 0; i < dataLength; i++) {
        myView[i] = this.socketData[i];
      }
      this.socketData = null; // Clear "socket" data after reading
    }
    return dataLength;
  }

  // Dummy close function
  close() {
    return;
  }

  // Return random number bytes in this call of socket
  getNumberRandomBytesSocket() {
    // Capped to remaining data and the max min return-per-read range
    const remaining_data = this.max_data - this.data_read;
    const numberBytesReceived =
      remaining_data < this.min_per_read
        ? remaining_data
        : this.getRandomIntInclusive(
            this.min_per_read,
            Math.min(this.max_per_read, remaining_data),
          );
    return numberBytesReceived;
  }

  // Return random number between two values
  getRandomIntInclusive(min, max) {
    min = Math.ceil(min);
    max = Math.floor(max);
    return Math.floor(Math.random() * (max - min + 1) + min);
  }

  // Return random character string
  randomChars(length = 8) {
    let string = "";
    let choices =
      "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789!@#$%^&*()";

    for (let i = 0; i < length; i++) {
      string += choices.charAt(Math.floor(Math.random() * choices.length));
    }
    return string;
  }

  /* Return random Uint8Array of bytes */
  randomByteArray(bytes = 8) {
    const textEncoder = new TextEncoder();
    return textEncoder.encode(this.randomChars(bytes));
  }
}

Создание читаемого потока байтов сокета push

Следующий код показывает, как определить читаемый байтовый поток сокета "push".

Определение объекта underlyingSource передаётся в качестве первого параметра конструктору ReadableStream(). Чтобы сделать это читаемым "байтовым" потоком, мы указываем type: "bytes" как свойство объекта. Это гарантирует, что потоку передаётся ReadableByteStreamController (вместо контроллера по умолчанию (ReadableStreamDefaultController)).

Поскольку данные могут поступить в сокет до того, как потребитель будет готов их обработать, все параметры чтения базового источника настраиваются в функции обратного вызова start() (мы не ждём pull, чтобы начать обработку данных). Реализация открывает "сокет" и вызывает select2() для запроса данных. Когда возвращённый промис разрешается, код проверяет, существует ли controller.byobRequest (не null), и если да, вызывает socket.readInto() для копирования данных в запрос и передачи его. Если byobRequest не существует, нет ожидающего запроса от потребляющего потока, который может быть удовлетворён как передача zero-copy. В этом случае, controller.enqueue() используется для копирования данных во внутренние очереди потока.

Запрос select2() на получение дополнительных данных повторяется до тех пор, пока запрос не будет возвращён без данных. В этот момент контроллер используется для закрытия потока.

const stream = makeSocketStream("dummy host", "dummy port");

const DEFAULT_CHUNK_SIZE = 400;

function makeSocketStream(host, port) {
  const socket = new MockHypotheticalSocket();

  return new ReadableStream({
    type: "bytes",

    start(controller) {
      readRepeatedly().catch((e) => controller.error(e));
      function readRepeatedly() {
        return socket.select2().then(() => {
          // Since the socket can become readable even when there's
          // no pending BYOB requests, we need to handle both cases.
          let bytesRead;
          if (controller.byobRequest) {
            const v = controller.byobRequest.view;
            bytesRead = socket.readInto(v.buffer, v.byteOffset, v.byteLength);
            if (bytesRead === 0) {
              controller.close();
            }
            controller.byobRequest.respond(bytesRead);
            logSource(`byobRequest with ${bytesRead} bytes`);
          } else {
            const buffer = new ArrayBuffer(DEFAULT_CHUNK_SIZE);
            bytesRead = socket.readInto(buffer, 0, DEFAULT_CHUNK_SIZE);
            if (bytesRead === 0) {
              controller.close();
            } else {
              controller.enqueue(new Uint8Array(buffer, 0, bytesRead));
            }
            logSource(`enqueue() ${bytesRead} bytes (no byobRequest)`);
          }

          if (bytesRead === 0) {
            return;
            // no more bytes in source
          }
          return readRepeatedly();
        });
      }
    },

    cancel() {
      socket.close();
      logSource(`cancel(): socket closed`);
    },
  });
}

Обратите внимание, что readRepeatedly() возвращает промис, и мы используем его для перехвата любых ошибок при настройке или обработке операции чтения. Ошибки затем передаются контроллеру, как показано выше (см. readRepeatedly().catch((e) => controller.error(e));).

В конце предоставляется метод cancel() для закрытия базового источника; функция обратного вызова pull() не требуется и поэтому не реализована.

Потребление потока байтов push

Следующий код создаёт ReadableStreamBYOBReader для потока байтов сокета и использует его для чтения данных в буфер. Обратите внимание, что processText() вызывается рекурсивно для чтения дополнительных данных до заполнения буфера. Когда базовый источник сигнализирует о том, что больше нет данных, reader.read() будет иметь done установленным в true, что, в свою очередь, завершит операцию чтения.

Этот код почти идентичен коду для примера Базового источника pull с байтовым ридером выше. Единственное различие заключается в том, что ридер включает некоторый код для замедления чтения, поэтому вывод журнала может продемонстрировать, что данные будут помещены в очередь, если они не будут читаться достаточно быстро.

const reader = stream.getReader({ mode: "byob" });
let buffer = new ArrayBuffer(4000);
readStream(reader);

function readStream(reader) {
  let bytesReceived = 0;
  let offset = 0;

  while (offset < buffer.byteLength) {
    // read() returns a promise that resolves when a value has been received
    reader
      .read(new Uint8Array(buffer, offset, buffer.byteLength - offset))
      .then(async function processText({ done, value }) {
        // Result objects contain two properties:
        // done  - true if the stream has already given all its data.
        // value - some data. Always undefined when done is true.

        if (done) {
          logConsumer(`readStream() complete. Total bytes: ${bytesReceived}`);
          return;
        }

        buffer = value.buffer;
        offset += value.byteLength;
        bytesReceived += value.byteLength;

        //logConsumer(`Read ${bytesReceived} bytes: ${value}`);
        logConsumer(`Read ${bytesReceived} bytes`);
        result += value;

        // Add delay to emulate when data can't be read and data is enqueued
        if (bytesReceived > 300 && bytesReceived < 600) {
          logConsumer(`Delaying read to emulate slow stream reading`);
          const delay = (ms) =>
            new Promise((resolve) => setTimeout(resolve, ms));
          await delay(1000);
        }

        // Read some more, and call this function again
        return reader
          .read(new Uint8Array(buffer, offset, buffer.byteLength - offset))
          .then(processText);
      });
  }
}

Отмена потока с помощью ридера

Мы можем использовать ReadableStreamBYOBReader.cancel() для отмены потока. В этом примере мы вызываем метод, если нажата кнопка с причиной "выбор пользователя" (другой HTML и код для кнопки не показаны). Мы также регистрируем завершение операции отмены.

button.addEventListener("click", () => {
  reader
    .cancel("user choice")
    .then(() => logConsumer("reader.cancel complete"));
});

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

Мониторинг потока на закрытие/ошибку

Свойство ReadableStreamBYOBReader.closed возвращает промис, который разрешится, когда поток будет закрыт, и отклонится, если произойдёт ошибка. Хотя в этом случае ошибок не ожидается, следующий код должен записать случай завершения.

reader.closed
  .then(() => {
    logConsumer("ReadableStreamBYOBReader.closed: resolved");
  })
  .catch(() => {
    logConsumer("ReadableStreamBYOBReader.closed: rejected:");
  });

Результат

Ниже показаны логи от базового источника push (слева) и потребителя (справа). Обратите внимание на промежуток времени, в который данные помещаются в очередь вместо передачи как операции zero-copy.

Базовый источник pull с байтовым ридером

Этот пример показывает, как данные могут быть считаны из "pull" базового байтового источника, например, файла, и переданы потоком как передача zero-copy в ReadableStreamBYOBReader.

Эмулированный базовый источник файла

Для базового источника pull мы используем следующий класс для (очень поверхностного) эмулирования nodejs FileHandle, а в частности метода read(). Класс генерирует случайные данные для представления файла. Метод read() читает "полуслучайный" фрагмент случайных данных в предоставленный буфер из указанной позиции. Метод close() ничего не делает: он предоставлен только для демонстрации места, где вы можете закрыть источник при определении конструктора потока.

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

class MockUnderlyingFileHandle {
  constructor() {
    this.maxdata = 100; // "file size"
    this.maxReadChunk = 25; // "max read chunk size"
    this.minReadChunk = 13; // "min read chunk size"
    this.filedata = this.randomByteArray(this.maxdata);
    this.position = 0;
  }

  // Read data from "file" at position/length into specified buffer offset
  read(buffer, offset, length, position) {
    // Object used to resolve promise
    const resultObj = {};
    resultObj["buffer"] = buffer;
    resultObj["bytesRead"] = 0;

    return new Promise((resolve /*, reject*/) => {
      if (position >= this.maxdata) {
        //out of data
        resolve(resultObj);
        return;
      }

      // Simulate a file read that returns random numbers of bytes
      // Read minimum of bytes requested and random bytes that can be returned
      let readLength =
        Math.floor(
          Math.random() * (this.maxReadChunk - this.minReadChunk + 1),
        ) + this.minReadChunk;
      readLength = length > readLength ? readLength : length;

      // Read random data into supplied buffer
      const myView = new Uint8Array(buffer, offset, readLength);
      // Write the length of data specified
      for (let i = 0; i < readLength; i++) {
        myView[i] = this.filedata[position + i];
        resultObj["bytesRead"] = i + 1;
        if (position + i + 1 >= this.maxdata) {
          break;
        }
      }
      // Emulate slow read of data
      setTimeout(() => {
        resolve(resultObj);
      }, 1000);
    });
  }

  // Dummy close function
  close() {
    return;
  }

  // Return random character string
  randomChars(length = 8) {
    let string = "";
    let choices =
      "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789!@#$%^&*()";

    for (let i = 0; i < length; i++) {
      string += choices.charAt(Math.floor(Math.random() * choices.length));
    }
    return string;
  }

  // Return random Uint8Array of bytes
  randomByteArray(bytes = 8) {
    const textEncoder = new TextEncoder();
    return textEncoder.encode(this.randomChars(bytes));
  }
}

Создание читаемого байтового потока файла

Следующий код показывает, как определить читаемый байтовый поток файла.

Так же, как и в предыдущем примере, определение объекта underlyingSource передаётся в качестве первого параметра конструктору ReadableStream(). Чтобы сделать это читаемым "байтовым" потоком, мы указываем type: "bytes" как свойство объекта. Это гарантирует, что потоку передаётся ReadableByteStreamController.

Функция start() просто открывает дескриптор файла, который затем закрывается в функции обратного вызова cancel(). cancel() предоставляет возможность очистить ресурсы, если вызывается ReadableStream.cancel() или ReadableStreamDefaultController.close().

Большая часть интересного кода находится в функции обратного вызова pull(). Она копирует данные из файла в ожидающий запрос на чтение (ReadableByteStreamController.byobRequest), а затем вызывает respond(), чтобы указать, сколько данных находится в буфере, и передать их. Если из файла было передано 0 байт, мы знаем, что всё было скопировано, и вызываем close() на контроллере, что, в свою очередь, вызовет cancel() на базовом источнике.

const stream = makeReadableByteFileStream("dummy file.txt");

function makeReadableByteFileStream(filename) {
  let fileHandle;
  let position = 0;
  return new ReadableStream({
    type: "bytes", // An underlying byte stream!
    start(controller) {
      // Called to initialise the underlying source.
      // For a file source open a file handle (here we just create the mocked object).
      fileHandle = new MockUnderlyingFileHandle();
      logSource(
        `start(): ${controller.constructor.name}.byobRequest = ${controller.byobRequest}`,
      );
    },
    async pull(controller) {
      // Called when there is a pull request for data
      const theView = controller.byobRequest.view;
      const { bytesRead, buffer } = await fileHandle.read(
        theView.buffer,
        theView.byteOffset,
        theView.byteLength,
        position,
      );
      if (bytesRead === 0) {
        await fileHandle.close();
        controller.close();
        controller.byobRequest.respond(0);
        logSource(
          `pull() with byobRequest. Close controller (read bytes: ${bytesRead})`,
        );
      } else {
        position += bytesRead;
        controller.byobRequest.respond(bytesRead);
        logSource(`pull() with byobRequest. Transfer ${bytesRead} bytes`);
      }
    },
    cancel(reason) {
      // This is called if the stream is cancelled (via reader or controller).
      // Clean up any resources
      fileHandle.close();
      logSource(`cancel() with reason: ${reason}`);
    },
  });
}

Потребление байтового потока

Следующий код создаёт ReadableStreamBYOBReader для потока байтов файла и использует его для чтения данных в буфер. Обратите внимание, что processText() вызывается рекурсивно для чтения дополнительных данных до заполнения буфера. Когда базовый источник сигнализирует о том, что больше нет данных, reader.read() будет иметь done установленным в true, что, в свою очередь, завершит операцию чтения.

const reader = stream.getReader({ mode: "byob" });
let buffer = new ArrayBuffer(200);
readStream(reader);

function readStream(reader) {
  let bytesReceived = 0;
  let offset = 0;

  // read() returns a promise that resolves when a value has been received
  reader
    .read(new Uint8Array(buffer, offset, buffer.byteLength - offset))
    .then(function processText({ done, value }) {
      // Result objects contain two properties:
      // done  - true if the stream has already given all its data.
      // value - some data. Always undefined when done is true.

      if (done) {
        logConsumer(`readStream() complete. Total bytes: ${bytesReceived}`);
        return;
      }

      buffer = value.buffer;
      offset += value.byteLength;
      bytesReceived += value.byteLength;

      logConsumer(
        `Read ${value.byteLength} (${bytesReceived}) bytes: ${value}`,
      );
      result += value;

      // Read some more, and call this function again
      return reader
        .read(new Uint8Array(buffer, offset, buffer.byteLength - offset))
        .then(processText);
    });
}

Наконец, мы добавляем обработчик, который отменит поток, если будет нажата кнопка (другой HTML и код для кнопки не показаны).

button.addEventListener("click", () => {
  reader.cancel("user choice").then(() => {
    logConsumer(`reader.cancel complete`);
  });
});

Результат

Ниже показаны логи от базового источника pull (слева) и потребителя (справа). Обратите особое внимание на:

  • start() функция получает ReadableByteStreamController
  • буфер, переданный ридеру, достаточно велик, чтобы охватить весь "файл". Базовый источник данных предоставляет данные случайными фрагментами.

Базовый источник pull с ридером по умолчанию

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

Создание потока байтов из файла для чтения с автоматическим выделением буфера

Единственное отличие в нашем исходном источнике заключается в том, что мы должны указать autoAllocateChunkSize, и этот размер будет использован в качестве размера буфера представления для controller.byobRequest, а не задаваться потребителем.

const DEFAULT_CHUNK_SIZE = 20;
const stream = makeReadableByteFileStream("dummy file.txt");

function makeReadableByteFileStream(filename) {
  let fileHandle;
  let position = 0;
  return new ReadableStream({
    type: "bytes", // An underlying byte stream!
    start(controller) {
      // Called to initialise the underlying source.
      // For a file source open a file handle (here we just create the mocked object).
      fileHandle = new MockUnderlyingFileHandle();
      logSource(
        `start(): ${controller.constructor.name}.byobRequest = ${controller.byobRequest}`,
      );
    },
    async pull(controller) {
      // Called when there is a pull request for data
      const theView = controller.byobRequest.view;
      const { bytesRead, buffer } = await fileHandle.read(
        theView.buffer,
        theView.byteOffset,
        theView.byteLength,
        position,
      );
      if (bytesRead === 0) {
        await fileHandle.close();
        controller.close();
        controller.byobRequest.respond(0);
        logSource(
          `pull() with byobRequest. Close controller (read bytes: ${bytesRead})`,
        );
      } else {
        position += bytesRead;
        controller.byobRequest.respond(bytesRead);
        logSource(`pull() with byobRequest. Transfer ${bytesRead} bytes`);
      }
    },
    cancel(reason) {
      // This is called if the stream is cancelled (via reader or controller).
      // Clean up any resources
      fileHandle.close();
      logSource(`cancel() with reason: ${reason}`);
    },
    autoAllocateChunkSize: DEFAULT_CHUNK_SIZE, // Only relevant if using a default reader
  });
}

Обработка потока байтов с помощью стандартного читателя

Следующий код создаёт ReadableStreamDefaultReader для потока байтов из файла, вызывая stream.getReader(); без указания режима и использует его для чтения данных в буфер. Действие кода такое же, как в предыдущем примере, за исключением того, что буфер предоставляется потоком, а не потребителем.

const reader = stream.getReader();
readStream(reader);

function readStream(reader) {
  let bytesReceived = 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) {
      logConsumer(`readStream() complete. Total bytes: ${bytesReceived}`);
      return;
    }

    bytesReceived += value.length;
    logConsumer(
      `Read ${value.length} (${bytesReceived}). Current bytes = ${value}`,
    );
    result += value;

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

Наконец, мы добавляем обработчик, который отменит поток при нажатии на кнопку (другой HTML и код для кнопки не показаны).

button.addEventListener("click", () => {
  reader.cancel("user choice").then(() => {
    logConsumer(`reader.cancel complete`);
  });
});

Результат

Ниже показаны журналы из исходного потока байтов (слева) и потребителя (справа).

Обратите внимание, что фрагменты теперь имеют максимальную ширину 20 байт, поскольку это размер автоматически выделенного буфера, указанного в исходном источнике байтов (autoAllocateChunkSize). Эти фрагменты создаются как передачи без копирования.

Исходный поток байтов с использованием стандартного читателя и без выделения памяти

Для полноты картины, мы также можем использовать стандартный читатель с источником байтов, который не поддерживает автоматическое выделение буфера.

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

const stream = makeReadableByteFileStream("dummy file.txt");
const DEFAULT_CHUNK_SIZE = 40;

function makeReadableByteFileStream(filename) {
  let fileHandle;
  let position = 0;
  return new ReadableStream({
    type: "bytes", // An underlying byte stream!
    start(controller) {
      // Called to initialise the underlying source.
      // For a file source open a file handle (here we just create the mocked object).
      fileHandle = new MockUnderlyingFileHandle();
      logSource(
        `start(): ${controller.constructor.name}.byobRequest = ${controller.byobRequest}`,
      );
    },
    async pull(controller) {
      // Called when there is a pull request for data
      if (controller.byobRequest) {
        const theView = controller.byobRequest.view;
        const { bytesRead, buffer } = await fileHandle.read(
          theView.buffer,
          theView.byteOffset,
          theView.byteLength,
          position,
        );
        if (bytesRead === 0) {
          await fileHandle.close();
          controller.close();
          controller.byobRequest.respond(0);
          logSource(
            `pull() with byobRequest. Close controller (read bytes: ${bytesRead})`,
          );
        } else {
          position += bytesRead;
          controller.byobRequest.respond(bytesRead);
          logSource(`pull() with byobRequest. Transfer ${bytesRead} bytes`);
        }
      } else {
        // No BYOBRequest so enqueue data to stream
        // NOTE, this branch would only execute for a default reader if autoAllocateChunkSize is not defined.
        const myNewBuffer = new Uint8Array(DEFAULT_CHUNK_SIZE);
        const { bytesRead, buffer } = await fileHandle.read(
          myNewBuffer.buffer,
          myNewBuffer.byteOffset,
          myNewBuffer.byteLength,
          position,
        );
        if (bytesRead === 0) {
          await fileHandle.close();
          controller.close();
          controller.enqueue(myNewBuffer);
          logSource(
            `pull() with no byobRequest. Close controller (read bytes: ${bytesRead})`,
          );
        } else {
          position += bytesRead;
          controller.enqueue(myNewBuffer);
          logSource(`pull() with no byobRequest. enqueue() ${bytesRead} bytes`);
        }
      }
    },
    cancel(reason) {
      // This is called if the stream is cancelled (via reader or controller).
      // Clean up any resources
      fileHandle.close();
      logSource(`cancel() with reason: ${reason}`);
    },
  });
}

Результат

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

См. также

  • Концепции API потоков
  • Обзор концепций и использования потоков
  • Использование потоков для чтения

© 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_byte_streams

Spec-Zone.ru

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