> ## Documentation Index
> Fetch the complete documentation index at: https://bun.ll1025.cn/llms.txt
> Use this file to discover all available pages before exploring further.

# 流

> 使用 Bun 的流 API 处理二进制数据，无需一次性全部加载到内存中

流是一种重要的抽象，用于处理二进制数据而无需一次性加载到内存中。它们通常用于读取和写入文件、发送和接收网络请求，以及处理大量数据。

Bun 实现了 Web API [`ReadableStream`](https://developer.mozilla.org/en-US/docs/Web/API/ReadableStream) 和 [`WritableStream`](https://developer.mozilla.org/en-US/docs/Web/API/WritableStream)。

<Note>
  Bun 也实现了 `node:stream` 模块，包括
  [`Readable`](https://nodejs.org/api/stream.html#stream_readable_streams)、
  [`Writable`](https://nodejs.org/api/stream.html#stream_writable_streams) 和
  [`Duplex`](https://nodejs.org/api/stream.html#stream_duplex_and_transform_streams)。完整的文档请参考
  [Node.js 文档](https://nodejs.org/api/stream.html)。
</Note>

要创建一个 `ReadableStream`：

```ts theme={null}
const stream = new ReadableStream({
  start(controller) {
    controller.enqueue("hello");
    controller.enqueue("world");
    controller.close();
  },
});
```

`ReadableStream` 的内容可以使用 `for await` 语法逐块读取。

```ts theme={null}
for await (const chunk of stream) {
  console.log(chunk);
}

// hello
// world
```

***

## Direct `ReadableStream`

Bun 实现了一个优化版本的 `ReadableStream`，避免了不必要的数据复制和队列管理。

使用传统的 `ReadableStream`，数据块会被\_入队\_。每个数据块被复制到队列中，在那里等待直到流准备好发送更多数据。

```ts theme={null}
const stream = new ReadableStream({
  start(controller) {
    controller.enqueue("hello");
    controller.enqueue("world");
    controller.close();
  },
});
```

使用 direct `ReadableStream`，数据块会直接写入流。没有队列操作，也不需要将数据块复制到内存中。`controller` API 反映了这一点：你调用 `.write()` 而不是 `.enqueue()`。

```ts theme={null}
const stream = new ReadableStream({
  type: "direct", // [!code ++]
  pull(controller) {
    controller.write("hello");
    controller.write("world");
  },
});
```

使用 direct `ReadableStream` 时，目标端处理所有数据块排队。流的消费者会精确接收到传递给 `controller.write()` 的内容，不进行任何编码或修改。

### 处理背压

`controller.write()` 返回写入的字节数。当目标端的内部缓冲区已满（例如，慢速的 HTTP 客户端）时，它会返回一个**负数**。数据块仍然被接受——负返回值是一个信号，表示暂停并等待目标端排空。

要等待排空，使用 `await controller.flush(true)`：

```ts theme={null}
const stream = new ReadableStream({
  type: "direct",
  async pull(controller) {
    for (const chunk of chunks) {
      const n = controller.write(chunk);
      if (typeof n === "number" && n < 0) {
        // 目标端已积压；等待其排空后再写入更多数据
        await controller.flush(true);
      }
    }
    controller.close();
  },
});
```

对于默认（非 `direct`）的 `ReadableStream` 和异步生成器响应体，Bun 会自动应用此背压——当目标端积压时，生产方会暂停。

***

## 异步生成器流

Bun 也支持异步生成器函数作为 `Response` 和 `Request` 的数据源。使用异步生成器从异步源创建 `ReadableStream`。

```ts theme={null}
const response = new Response(
  (async function* () {
    yield "hello";
    yield "world";
  })(),
);

await response.text(); // "helloworld"
```

你也可以直接使用 `[Symbol.asyncIterator]`。

```ts theme={null}
const response = new Response({
  [Symbol.asyncIterator]: async function* () {
    yield "hello";
    yield "world";
  },
});

await response.text(); // "helloworld"
```

要对流进行更多控制，`yield` 返回 direct `ReadableStream` controller。

```ts theme={null}
const response = new Response({
  [Symbol.asyncIterator]: async function* () {
    const controller = yield "hello";
    await controller.end();
  },
});

await response.text(); // "hello"
```

***

## `Bun.ArrayBufferSink`

`Bun.ArrayBufferSink` 类是一个快速的增量写入器，用于构建未知大小的 `ArrayBuffer`。

```ts theme={null}
const sink = new Bun.ArrayBufferSink();

sink.write("h");
sink.write("e");
sink.write("l");
sink.write("l");
sink.write("o");

sink.end();
// ArrayBuffer(5) [ 104, 101, 108, 108, 111 ]
```

要改为以 `Uint8Array` 形式检索数据，在 `start` 方法中传递 `asUint8Array` 选项。

```ts theme={null}
const sink = new Bun.ArrayBufferSink();
sink.start({
  asUint8Array: true, // [!code ++]
});

sink.write("h");
sink.write("e");
sink.write("l");
sink.write("l");
sink.write("o");

sink.end();
// Uint8Array(5) [ 104, 101, 108, 108, 111 ]
```

`.write()` 方法支持字符串、类型化数组、`ArrayBuffer` 和 `SharedArrayBuffer`。

```ts theme={null}
sink.write("h");
sink.write(new Uint8Array([101, 108]));
sink.write(Buffer.from("lo").buffer);

sink.end();
```

一旦调用了 `.end()`，就不能再向 `ArrayBufferSink` 写入更多数据。但是，在缓冲流时，你可能希望持续写入数据并定期将内容 `.flush()`（例如，写入到 `WritableStream`）。为了支持这一点，在 `start` 方法中传递 `stream: true`。

```ts theme={null}
const sink = new Bun.ArrayBufferSink();
sink.start({
  stream: true, // [!code ++]
});

sink.write("h");
sink.write("e");
sink.write("l");
sink.flush();
// ArrayBuffer(3) [ 104, 101, 108 ]

sink.write("l");
sink.write("o");
sink.flush();
// ArrayBuffer(2) [ 108, 111 ]
```

`.flush()` 方法以 `ArrayBuffer`（如果 `asUint8Array: true` 则为 `Uint8Array`）的形式返回缓冲的数据，并清除内部缓冲区。

要手动设置内部缓冲区的大小（以字节为单位），为 `highWaterMark` 传递一个值：

```ts theme={null}
const sink = new Bun.ArrayBufferSink();
sink.start({
  highWaterMark: 1024 * 1024, // 1 MB  // [!code ++]
});
```

***

## 参考

```ts 查看 TypeScript 定义 expandable theme={null}
/**
 * 快速增量写入器，在 end() 后变为 ArrayBuffer。
 */
export class ArrayBufferSink {
  constructor();

  start(options?: {
    asUint8Array?: boolean;
    /**
     * 预分配此大小的内部缓冲区
     * 当数据块大小较小时，这可以显著提高性能
     */
    highWaterMark?: number;
    /**
     * 在 {@link ArrayBufferSink.flush} 时，将写入的数据作为 `Uint8Array` 返回。
     * 写入将从缓冲区开头重新开始。
     */
    stream?: boolean;
  }): void;

  write(chunk: string | ArrayBufferView | ArrayBuffer | SharedArrayBuffer): number;
  /**
   * 刷新内部缓冲区
   *
   * 如果 {@link ArrayBufferSink.start} 传递了 `stream` 选项，将返回一个 `ArrayBuffer`
   * 如果传递了 `stream` 选项和 `asUint8Array`，将返回一个 `Uint8Array`
   * 否则，将返回自上次刷新以来写入的字节数
   *
   * 此 API 将来可能会更改，以分离 Uint8ArraySink 和 ArrayBufferSink
   */
  flush(): number | Uint8Array<ArrayBuffer> | ArrayBuffer;
  end(): ArrayBuffer | Uint8Array<ArrayBuffer>;
}
```
