Web Streams API#

稳定性:2 - 稳定

这是 WHATWG 流标准的实现。

概览#

WHATWG 流标准(或称“web 流”)定义了用于处理流式数据的 API。它类似于 Node.js 的 Streams API,但出现得更晚,并已成为在多个 JavaScript 环境中处理流数据的“标准”API。

主要有三种类型的对象

  • ReadableStream - 代表流式数据的源。
  • WritableStream - 代表流式数据的目的地。
  • TransformStream - 代表转换流式数据的算法。

ReadableStream 示例#

此示例创建了一个简单的 ReadableStream,它每秒永久推送一次当前的 performance.now() 时间戳。使用异步可迭代对象从流中读取数据。

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);
const {
  ReadableStream,
} = require('node:stream/web');

const {
  setInterval: every,
} = require('node:timers/promises');

const {
  performance,
} = require('node:perf_hooks');

const SECOND = 1000;

const stream = new ReadableStream({
  async start(controller) {
    for await (const _ of every(SECOND))
      controller.enqueue(performance.now());
  },
});

(async () => {
  for await (const value of stream)
    console.log(value);
})();

Node.js 流的互操作性#

Node.js 流可以通过 stream.Readablestream.Writablestream.Duplex 对象上提供的 toWebfromWeb 方法与 web 流进行互转。

有关详细信息,请参阅相关文档

API#

类:ReadableStream#

new ReadableStream([underlyingSource [, strategy]])#
  • underlyingSource <Object>
    • start <Function> 创建 ReadableStream 时立即调用的用户定义函数。
    • pull <Function>ReadableStream 内部队列未满时重复调用的用户定义函数。该操作可以是同步或异步的。如果是异步的,在之前返回的 Promise 兑现之前,该函数不会再次被调用。
    • cancel <Function>ReadableStream 被取消时调用的用户定义函数。
      • reason <any>
      • 返回:以 undefined 兑现的 Promise。
    • type <string> 必须为 'bytes'undefined
    • autoAllocateChunkSize <number> 仅在 type 等于 'bytes' 时使用。当设置为非零值时,会自动分配视图缓冲区给 ReadableByteStreamController.byobRequest。如果未设置,则必须通过默认读取器 ReadableStreamDefaultReader 使用流的内部队列来传输数据。
  • strategy <Object>
  • highWaterMark <number> 应用背压之前的最大内部队列大小。
  • size <Function> 用于标识每个数据块大小的用户定义函数。
readableStream.locked#

readableStream.locked 属性默认为 false,当有活跃的读取器正在消耗流数据时,会切换为 true

readableStream.cancel([reason])#
  • reason <any>
  • 返回:取消完成后以 undefined 兑现的 Promise。
readableStream.getReader([options])#
  • 返回:<ReadableStreamDefaultReader> | <ReadableStreamBYOBReader>
  • import { ReadableStream } from 'node:stream/web';
    
    const stream = new ReadableStream();
    
    const reader = stream.getReader();
    
    console.log(await reader.read());
    const { ReadableStream } = require('node:stream/web');
    
    const stream = new ReadableStream();
    
    const reader = stream.getReader();
    
    reader.read().then(console.log);
    

    会导致 readableStream.locked 变为 true

    readableStream.pipeThrough(transform[, options])#
    • transform <Object>
      • readable <ReadableStream> transform.writable 会将从该 ReadableStream 接收并可能经过修改的数据推送到此 ReadableStream
      • writable <WritableStream>ReadableStream 的数据将写入到的 WritableStream
    • options <Object>
    • preventAbort <boolean> 当为 true 时,此 ReadableStream 中的错误不会导致 transform.writable 中止。
    • preventCancel <boolean> 当为 true 时,目标 transform.writable 中的错误不会导致此 ReadableStream 被取消。
    • preventClose <boolean> 当为 true 时,关闭此 ReadableStream 不会导致 transform.writable 关闭。
    • signal <AbortSignal> 允许使用 <AbortController> 取消数据传输。
  • 返回:<ReadableStream> 来自 transform.readable
  • 将此 <ReadableStream> 连接到 transform 参数中提供的 <ReadableStream><WritableStream> 对,使得此 <ReadableStream> 的数据被写入 transform.writable,可能经过转换,然后推送到 transform.readable。一旦管道配置完成,即返回 transform.readable

    在管道操作活跃时,会导致 readableStream.locked 变为 true

    import {
      ReadableStream,
      TransformStream,
    } from 'node:stream/web';
    
    const stream = new ReadableStream({
      start(controller) {
        controller.enqueue('a');
      },
    });
    
    const transform = new TransformStream({
      transform(chunk, controller) {
        controller.enqueue(chunk.toUpperCase());
      },
    });
    
    const transformedStream = stream.pipeThrough(transform);
    
    for await (const chunk of transformedStream)
      console.log(chunk);
      // Prints: A
    const {
      ReadableStream,
      TransformStream,
    } = require('node:stream/web');
    
    const stream = new ReadableStream({
      start(controller) {
        controller.enqueue('a');
      },
    });
    
    const transform = new TransformStream({
      transform(chunk, controller) {
        controller.enqueue(chunk.toUpperCase());
      },
    });
    
    const transformedStream = stream.pipeThrough(transform);
    
    (async () => {
      for await (const chunk of transformedStream)
        console.log(chunk);
        // Prints: A
    })();
    
    readableStream.pipeTo(destination[, options])#
    • destination <WritableStream>ReadableStream 的数据将写入到的 <WritableStream>
    • options <Object>
    • preventAbort <boolean> 当为 true 时,此 ReadableStream 中的错误不会导致 destination 中止。
    • preventCancel <boolean> 当为 true 时,destination 中的错误不会导致此 ReadableStream 被取消。
    • preventClose <boolean> 当为 true 时,关闭此 ReadableStream 不会导致 destination 关闭。
    • signal <AbortSignal> 允许使用 <AbortController> 取消数据传输。
  • 返回:以 undefined 兑现的 Promise
  • 在管道操作活跃时,会导致 readableStream.locked 变为 true

    readableStream.tee()#

    返回一对新的 <ReadableStream> 实例,此 ReadableStream 的数据将转发给它们。两者都将接收相同的数据。

    会导致 readableStream.locked 变为 true

    readableStream.values([options])#

    创建并返回一个可用于消耗此 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());
    
    异步迭代#

    <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());
    

    异步迭代器将消耗 <ReadableStream> 直至其终止。

    默认情况下,如果异步迭代器提前退出(通过 breakreturnthrow),<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]);
    

    ReadableStream.from(iterable)#

    • iterable <Iterable> 实现 Symbol.asyncIteratorSymbol.iterator 可迭代协议的对象。

    从可迭代对象创建新的 <ReadableStream> 的实用方法。

    import { ReadableStream } from 'node:stream/web';
    
    async function* asyncIterableGenerator() {
      yield 'a';
      yield 'b';
      yield 'c';
    }
    
    const stream = ReadableStream.from(asyncIterableGenerator());
    
    for await (const chunk of stream)
      console.log(chunk); // Prints: 'a', 'b', 'c'
    const { ReadableStream } = require('node:stream/web');
    
    async function* asyncIterableGenerator() {
      yield 'a';
      yield 'b';
      yield 'c';
    }
    
    (async () => {
      const stream = ReadableStream.from(asyncIterableGenerator());
    
      for await (const chunk of stream)
        console.log(chunk); // Prints: 'a', 'b', 'c'
    })();
    

    要将生成的 <ReadableStream> 管道传输到 <WritableStream> 中,<Iterable> 应产生一系列 <Buffer><TypedArray><DataView> 对象。

    import { ReadableStream } from 'node:stream/web';
    import { Buffer } from 'node:buffer';
    
    async function* asyncIterableGenerator() {
      yield Buffer.from('a');
      yield Buffer.from('b');
      yield Buffer.from('c');
    }
    
    const stream = ReadableStream.from(asyncIterableGenerator());
    
    await stream.pipeTo(createWritableStreamSomehow());
    const { ReadableStream } = require('node:stream/web');
    const { Buffer } = require('node:buffer');
    
    async function* asyncIterableGenerator() {
      yield Buffer.from('a');
      yield Buffer.from('b');
      yield Buffer.from('c');
    }
    
    const stream = ReadableStream.from(asyncIterableGenerator());
    
    (async () => {
      await stream.pipeTo(createWritableStreamSomehow());
    })();
    

    类:ReadableStreamDefaultReader#

    默认情况下,不带参数调用 readableStream.getReader() 将返回 ReadableStreamDefaultReader 的实例。默认读取器将流经流的数据块视为不透明值,这允许 <ReadableStream> 处理一般的任何 JavaScript 值。

    new ReadableStreamDefaultReader(stream)#

    创建锁定到给定 <ReadableStream> 的新 <ReadableStreamDefaultReader>

    readableStreamDefaultReader.cancel([reason])#
    • reason <any>
    • 返回:以 undefined 兑现的 Promise。

    取消 <ReadableStream> 并返回一个在底层流被取消时兑现的 Promise。

    readableStreamDefaultReader.closed#
    • 类型:<Promise> 当关联的 <ReadableStream> 关闭时以 undefined 兑现,如果流报错或读取器的锁在流完成关闭之前释放,则被拒绝。
    readableStreamDefaultReader.read()#

    从底层 <ReadableStream> 请求下一个数据块,并返回一个在数据可用时以数据兑现的 Promise。

    readableStreamDefaultReader.releaseLock()#

    释放此读取器对底层 <ReadableStream> 的锁定。

    类:ReadableStreamBYOBReader#

    ReadableStreamBYOBReader 是面向字节的 <ReadableStream>(在创建 ReadableStream 时将 underlyingSource.type 设置为 'bytes' 的流)的替代消费者。

    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());
    
    new ReadableStreamBYOBReader(stream)#

    创建锁定到给定 <ReadableStream> 的新 ReadableStreamBYOBReader

    readableStreamBYOBReader.cancel([reason])#
    • reason <any>
    • 返回:以 undefined 兑现的 Promise。

    取消 <ReadableStream> 并返回一个在底层流被取消时兑现的 Promise。

    readableStreamBYOBReader.closed#
    • 类型:<Promise> 当关联的 <ReadableStream> 关闭时以 undefined 兑现,如果流报错或读取器的锁在流完成关闭之前释放,则被拒绝。
    readableStreamBYOBReader.read(view[, options])#
  • 返回:以对象兑现的 Promise
  • 从底层 <ReadableStream> 请求下一个数据块,并返回一个在数据可用时以数据兑现的 Promise。

    请勿将池化的 <Buffer> 对象实例传递给此方法。池化的 Buffer 对象是使用 Buffer.allocUnsafe()Buffer.from() 创建的,或者通常由各种 node:fs 模块回调返回。这些类型的 Buffer 使用共享的底层 <ArrayBuffer> 对象,其中包含来自所有池化 Buffer 实例的所有数据。当 Buffer<TypedArray><DataView> 被传递给 readableStreamBYOBReader.read() 时,视图的底层 ArrayBuffer 会被分离,这会使该 ArrayBuffer 上可能存在的所有现有视图无效。这对您的应用程序可能会造成灾难性的后果。

    readableStreamBYOBReader.releaseLock()#

    释放此读取器对底层 <ReadableStream> 的锁定。

    类:ReadableStreamDefaultController#

    每个 <ReadableStream> 都有一个控制器,负责管理内部状态和流队列。ReadableStreamDefaultController 是非面向字节的 ReadableStream 的默认控制器实现。

    readableStreamDefaultController.close()#

    关闭与此控制器关联的 <ReadableStream>

    readableStreamDefaultController.desiredSize#

    返回填满 <ReadableStream> 队列所需的剩余数据量。

    readableStreamDefaultController.enqueue([chunk])#

    <ReadableStream> 的队列追加一个新数据块。

    readableStreamDefaultController.error([error])#

    发出错误信号,导致 <ReadableStream> 报错并关闭。

    类:ReadableByteStreamController#

    每个 <ReadableStream> 都有一个控制器,负责管理内部状态和流队列。ReadableByteStreamController 专用于面向字节的 ReadableStream

    readableByteStreamController.byobRequest#
    readableByteStreamController.close()#

    关闭与此控制器关联的 <ReadableStream>

    readableByteStreamController.desiredSize#

    返回填满 <ReadableStream> 队列所需的剩余数据量。

    readableByteStreamController.enqueue(chunk)#

    <ReadableStream> 的队列追加一个新数据块。

    readableByteStreamController.error([error])#

    发出错误信号,导致 <ReadableStream> 报错并关闭。

    类:ReadableStreamBYOBRequest#

    当在面向字节的流中使用 ReadableByteStreamController 且使用 ReadableStreamBYOBReader 时,readableByteStreamController.byobRequest 属性提供了对代表当前读取请求的 ReadableStreamBYOBRequest 实例的访问权限。该对象用于获取对已提供给读取请求填写的 ArrayBuffer/TypedArray 的访问权限,并提供用于发出数据已提供信号的方法。

    readableStreamBYOBRequest.respond(bytesWritten)#

    发出信号表示 bytesWritten 字节数已写入 readableStreamBYOBRequest.view

    readableStreamBYOBRequest.respondWithNewView(view)#

    发出信号表示请求已通过写入新 BufferTypedArrayDataView 的字节兑现。

    readableStreamBYOBRequest.view#

    类:WritableStream#

    WritableStream 是发送流数据的目的地。

    import {
      WritableStream,
    } from 'node:stream/web';
    
    const stream = new WritableStream({
      write(chunk) {
        console.log(chunk);
      },
    });
    
    await stream.getWriter().write('Hello World');
    
    new WritableStream([underlyingSink[, strategy]])#
    • underlyingSink <Object>
      • start <Function> 创建 WritableStream 时立即调用的用户定义函数。
      • write <Function> 当数据块写入 WritableStream 时调用的用户定义函数。
      • close <Function>WritableStream 关闭时调用的用户定义函数。
        • 返回:以 undefined 兑现的 Promise。
      • abort <Function> 用于突然关闭 WritableStream 时调用的用户定义函数。
        • reason <any>
        • 返回:以 undefined 兑现的 Promise。
      • type <any> type 选项保留供将来使用,必须为 undefined。
    • strategy <Object>
    • highWaterMark <number> 应用背压之前的最大内部队列大小。
    • size <Function> 用于标识每个数据块大小的用户定义函数。
    writableStream.abort([reason])#
    • reason <any>
    • 返回:以 undefined 兑现的 Promise。

    突然终止 WritableStream。所有排队的写入都将被取消,其关联的 Promise 将被拒绝。

    writableStream.close()#
    • 返回:以 undefined 兑现的 Promise。

    当预计没有更多写入时,关闭 WritableStream

    writableStream.getWriter()#

    创建并返回一个新的写入器实例,可用于将数据写入 WritableStream

    writableStream.locked#

    writableStream.locked 属性默认为 false,当有活跃的写入器连接到此 WritableStream 时,会切换为 true

    使用 postMessage() 进行传输#

    <WritableStream> 实例可以使用 <MessagePort> 进行传输。

    const stream = new WritableStream(getWritableSinkSomehow());
    
    const { port1, port2 } = new MessageChannel();
    
    port1.onmessage = ({ data }) => {
      data.getWriter().write('hello');
    };
    
    port2.postMessage(stream, [stream]);
    

    类:WritableStreamDefaultWriter#

    new WritableStreamDefaultWriter(stream)#

    创建锁定到给定 WritableStream 的新 WritableStreamDefaultWriter

    writableStreamDefaultWriter.abort([reason])#
    • reason <any>
    • 返回:以 undefined 兑现的 Promise。

    突然终止 WritableStream。所有排队的写入都将被取消,其关联的 Promise 将被拒绝。

    writableStreamDefaultWriter.close()#
    • 返回:以 undefined 兑现的 Promise。

    当预计没有更多写入时,关闭 WritableStream

    writableStreamDefaultWriter.closed#
    • 类型:<Promise> 当关联的 <WritableStream> 关闭时以 undefined 兑现,如果流报错或写入器的锁在流完成关闭之前释放,则被拒绝。
    writableStreamDefaultWriter.desiredSize#

    填满 <WritableStream> 队列所需的数据量。

    writableStreamDefaultWriter.ready#
    • 类型:<Promise> 当写入器准备好使用时以 undefined 兑现。
    writableStreamDefaultWriter.releaseLock()#

    释放此写入器对底层 <ReadableStream> 的锁定。

    writableStreamDefaultWriter.write([chunk])#
    • chunk <any>
    • 返回:以 undefined 兑现的 Promise。

    <WritableStream> 的队列追加一个新数据块。

    类:WritableStreamDefaultController#

    WritableStreamDefaultController 管理 <WritableStream> 的内部状态。

    writableStreamDefaultController.error([error])#

    由用户代码调用,以发出在处理 WritableStream 数据时发生错误的信号。调用时,<WritableStream> 将被中止,当前待处理的写入将被取消。

    writableStreamDefaultController.signal#

    类: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(),
    ]);
    
    new TransformStream([transformer[, writableStrategy[, readableStrategy]]])#
    • transformer <Object>
      • start <Function> 创建 TransformStream 时立即调用的用户定义函数。
      • transform <Function> 在将写入 transformStream.writable 的数据块转发给 transformStream.readable 之前,接收并可能对其进行修改的用户定义函数。
      • flush <Function>TransformStream 的可写端关闭之前立即调用的用户定义函数,标志着转换过程的结束。
      • readableType <any> readableType 选项保留供将来使用,必须undefined
      • writableType <any> writableType 选项保留供将来使用,必须undefined
    • writableStrategy <Object>
      • highWaterMark <number> 应用背压之前的最大内部队列大小。
      • size <Function> 用于标识每个数据块大小的用户定义函数。
    • readableStrategy <Object>
      • highWaterMark <number> 应用背压之前的最大内部队列大小。
      • size <Function> 用于标识每个数据块大小的用户定义函数。
    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]);
    

    类:TransformStreamDefaultController#

    TransformStreamDefaultController 管理 TransformStream 的内部状态。

    transformStreamDefaultController.desiredSize#

    填满可读端队列所需的数据量。

    transformStreamDefaultController.enqueue([chunk])#

    向可读端队列追加一个数据块。

    transformStreamDefaultController.error([reason])#

    向可读端和可写端发出信号,表示在处理转换数据时发生错误,导致两侧都突然关闭。

    transformStreamDefaultController.terminate()#

    关闭传输的可读端,并导致可写端因错误而突然关闭。

    类:ByteLengthQueuingStrategy#

    new ByteLengthQueuingStrategy(init)#
    byteLengthQueuingStrategy.highWaterMark#
    byteLengthQueuingStrategy.size#

    类:CountQueuingStrategy#

    new CountQueuingStrategy(init)#
    countQueuingStrategy.highWaterMark#
    countQueuingStrategy.size#

    类:TextEncoderStream#

    new TextEncoderStream()#

    创建新的 TextEncoderStream 实例。

    textEncoderStream.encoding#

    TextEncoderStream 实例支持的编码。

    textEncoderStream.readable#
    textEncoderStream.writable#

    类:TextDecoderStream#

    new TextDecoderStream([encoding[, options]])#
    • encoding <string> 标识此 TextDecoder 实例支持的 encoding默认值: 'utf-8'
    • options <Object>
    • fatal <boolean> 如果解码失败是致命的,则为 true
    • ignoreBOM <boolean> 当为 true 时,TextDecoderStream 将在解码结果中包含字节顺序标记。当为 false 时,字节顺序标记将从输出中删除。此选项仅在 encoding'utf-8''utf-16be''utf-16le' 时使用。默认值:false

    创建新的 TextDecoderStream 实例。

    textDecoderStream.encoding#

    TextDecoderStream 实例支持的编码。

    textDecoderStream.fatal#

    如果解码错误导致抛出 TypeError,则该值将为 true

    textDecoderStream.ignoreBOM#

    如果解码结果将包含字节顺序标记 (BOM),则该值将为 true

    textDecoderStream.readable#
    textDecoderStream.writable#

    类:CompressionStream#

    new CompressionStream(format)#
    • format <string> 'deflate''deflate-raw''gzip''brotli' 之一。
    compressionStream.readable#
    compressionStream.writable#

    类:DecompressionStream#

    new DecompressionStream(format)#
    • format <string> 'deflate''deflate-raw''gzip''brotli' 之一。
    decompressionStream.readable#
    decompressionStream.writable#

    工具类消费者#

    工具类消费者函数为消耗流提供了常用选项。

    可通过以下方式访问

    import {
      arrayBuffer,
      blob,
      buffer,
      json,
      text,
    } from 'node:stream/consumers';
    const {
      arrayBuffer,
      blob,
      buffer,
      json,
      text,
    } = require('node:stream/consumers');
    
    streamConsumers.arrayBuffer(stream)#
    import { arrayBuffer } from 'node:stream/consumers';
    import { Readable } from 'node:stream';
    import { TextEncoder } from 'node:util';
    
    const encoder = new TextEncoder();
    const dataArray = encoder.encode('hello world from consumers!');
    
    const readable = Readable.from(dataArray);
    const data = await arrayBuffer(readable);
    console.log(`from readable: ${data.byteLength}`);
    // Prints: from readable: 76
    const { arrayBuffer } = require('node:stream/consumers');
    const { Readable } = require('node:stream');
    const { TextEncoder } = require('node:util');
    
    const encoder = new TextEncoder();
    const dataArray = encoder.encode('hello world from consumers!');
    const readable = Readable.from(dataArray);
    arrayBuffer(readable).then((data) => {
      console.log(`from readable: ${data.byteLength}`);
      // Prints: from readable: 76
    });
    
    streamConsumers.blob(stream)#
    import { blob } from 'node:stream/consumers';
    
    const dataBlob = new Blob(['hello world from consumers!']);
    
    const readable = dataBlob.stream();
    const data = await blob(readable);
    console.log(`from readable: ${data.size}`);
    // Prints: from readable: 27
    const { blob } = require('node:stream/consumers');
    
    const dataBlob = new Blob(['hello world from consumers!']);
    
    const readable = dataBlob.stream();
    blob(readable).then((data) => {
      console.log(`from readable: ${data.size}`);
      // Prints: from readable: 27
    });
    
    streamConsumers.buffer(stream)#
    import { buffer } from 'node:stream/consumers';
    import { Readable } from 'node:stream';
    import { Buffer } from 'node:buffer';
    
    const dataBuffer = Buffer.from('hello world from consumers!');
    
    const readable = Readable.from(dataBuffer);
    const data = await buffer(readable);
    console.log(`from readable: ${data.length}`);
    // Prints: from readable: 27
    const { buffer } = require('node:stream/consumers');
    const { Readable } = require('node:stream');
    const { Buffer } = require('node:buffer');
    
    const dataBuffer = Buffer.from('hello world from consumers!');
    
    const readable = Readable.from(dataBuffer);
    buffer(readable).then((data) => {
      console.log(`from readable: ${data.length}`);
      // Prints: from readable: 27
    });
    
    streamConsumers.bytes(stream)#
    import { bytes } from 'node:stream/consumers';
    import { Readable } from 'node:stream';
    import { Buffer } from 'node:buffer';
    
    const dataBuffer = Buffer.from('hello world from consumers!');
    
    const readable = Readable.from(dataBuffer);
    const data = await bytes(readable);
    console.log(`from readable: ${data.length}`);
    // Prints: from readable: 27
    const { bytes } = require('node:stream/consumers');
    const { Readable } = require('node:stream');
    const { Buffer } = require('node:buffer');
    
    const dataBuffer = Buffer.from('hello world from consumers!');
    
    const readable = Readable.from(dataBuffer);
    bytes(readable).then((data) => {
      console.log(`from readable: ${data.length}`);
      // Prints: from readable: 27
    });
    
    streamConsumers.json(stream)#
    import { json } from 'node:stream/consumers';
    import { Readable } from 'node:stream';
    
    const items = Array.from(
      {
        length: 100,
      },
      () => ({
        message: 'hello world from consumers!',
      }),
    );
    
    const readable = Readable.from(JSON.stringify(items));
    const data = await json(readable);
    console.log(`from readable: ${data.length}`);
    // Prints: from readable: 100
    const { json } = require('node:stream/consumers');
    const { Readable } = require('node:stream');
    
    const items = Array.from(
      {
        length: 100,
      },
      () => ({
        message: 'hello world from consumers!',
      }),
    );
    
    const readable = Readable.from(JSON.stringify(items));
    json(readable).then((data) => {
      console.log(`from readable: ${data.length}`);
      // Prints: from readable: 100
    });
    
    streamConsumers.text(stream)#
    import { text } from 'node:stream/consumers';
    import { Readable } from 'node:stream';
    
    const readable = Readable.from('Hello world from consumers!');
    const data = await text(readable);
    console.log(`from readable: ${data.length}`);
    // Prints: from readable: 27
    const { text } = require('node:stream/consumers');
    const { Readable } = require('node:stream');
    
    const readable = Readable.from('Hello world from consumers!');
    text(readable).then((data) => {
      console.log(`from readable: ${data.length}`);
      // Prints: from readable: 27
    });