Streams API ↗ 是一种 Web 标准 API,允许 JavaScript 以编程方式访问和处理数据流。
- ReadableStream
- ReadableStream BYOBReader
- ReadableStream DefaultReader
- TransformStream
- WritableStream
- WritableStream DefaultWriter
使用 Streams API 可以避免在内存中缓冲大型请求或响应。这使您能够在 Worker 128 MB 内存限制内解析极大的请求或响应正文。这比将整个 payload 缓冲到内存中更快,因为 Worker 可以增量处理数据,并允许 Worker 在内存限制内处理数 GB 的 payload 或文件。
Workers 无需在返回 Response 之前准备完整的响应正文。您可以使用 ReadableStream 在发送响应状态行和标头后流式传输响应正文。
Worker 可以使用 ReadableStream 作为正文创建 Response 对象。通过 ReadableStream 提供的任何数据都会在可用时流式传输到客户端。
export default {
async fetch(request, env, ctx) {
// Fetch from origin server.
const response = await fetch(request);
// ... and deliver our Response while that’s running.
return new Response(response.body, response);
},
};addEventListener("fetch", (event) => {
event.respondWith(fetchAndStream(event.request));
});
async function fetchAndStream(request) {
// Fetch from origin server.
const response = await fetch(request);
// ... and deliver our Response while that’s running.
return new Response(readable.body, response);
}from workers import WorkerEntrypoint, Response, fetch
class Default(WorkerEntrypoint):
async def fetch(self, request):
# Fetch from origin server.
response = await fetch(request)
# Stream the response body to the client.
return Response(response.body, headers=response.headers)TransformStream 和 ReadableStream.pipeTo() 方法可用于在流式传输响应正文时对其进行修改:
export default {
async fetch(request, env, ctx) {
// Fetch from origin server.
const response = await fetch(request);
const { readable, writable } = new TransformStream({
transform(chunk, controller) {
controller.enqueue(modifyChunkSomehow(chunk));
},
});
// Start pumping the body. NOTE: No await!
response.body.pipeTo(writable);
// ... and deliver our Response while that’s running.
return new Response(readable, response);
},
};addEventListener("fetch", (event) => {
event.respondWith(fetchAndStream(event.request));
});
async function fetchAndStream(request) {
// Fetch from origin server.
const response = await fetch(request);
const { readable, writable } = new TransformStream({
transform(chunk, controller) {
controller.enqueue(modifyChunkSomehow(chunk));
},
});
// Start pumping the body. NOTE: No await!
response.body.pipeTo(writable);
// ... and deliver our Response while that’s running.
return new Response(readable, response);
}from workers import WorkerEntrypoint, Response
from js import ReadableStream, TextEncoder
from pyodide.ffi import create_proxy, to_js
import asyncio
class Default(WorkerEntrypoint):
async def fetch(self, request):
enc = TextEncoder.new()
async def start(controller):
for i in range(5):
controller.enqueue(enc.encode(f"chunk {i}\n"))
await asyncio.sleep(0.1)
controller.close()
stream = ReadableStream.new(
to_js({"start": create_proxy(start)})
)
return Response(stream, headers={"Content-Type": "text/plain"})此示例调用 response.body.pipeTo(writable) 但未 await 它。这样就不会阻塞 fetchAndStream() 函数其余部分的执行。它会异步继续运行,直到响应完成或客户端断开连接。
运行时在响应返回给客户端后,可以继续运行函数(response.body.pipeTo(writable))。此示例将子请求的响应正文泵送到最终响应正文中。不过,您可以使用更复杂的逻辑,例如向正文添加前缀或后缀,或以某种方式处理它。
- 流式传输大型 JSON — 解析和转换大型 JSON 请求与响应正文
- MDN 的 Streams API 文档 ↗
- Streams API 规范 ↗
- 使用 ES modules 语法 编写 Worker 代码,以获得优化体验。