跳转到内容
搜索文档

流(Streams)

最后更新 查看 MarkdownAgent 设置

Streams API 是一种 Web 标准 API,允许 JavaScript 以编程方式访问和处理数据流。

使用 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)

TransformStreamReadableStream.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))。此示例将子请求的响应正文泵送到最终响应正文中。不过,您可以使用更复杂的逻辑,例如向正文添加前缀或后缀,或以某种方式处理它。


常见问题


相关资源

这篇文档对您有帮助吗?