INTEGRITY Документация

Streams

Streams API представляет собой веб-стандартный API, который позволяет JavaScript программно обращаться к потокам данных и обрабатывать их.

Используйте Streams API, чтобы не буферизировать в памяти большие запросы и ответы. Это позволяет обрабатывать очень большие тела запросов и ответов в пределах лимита памяти Worker в 128 МБ. Такой подход быстрее буферизации всей полезной нагрузки в памяти, поскольку Worker начинает обрабатывать данные постепенно, и позволяет ему работать с файлами и полезной нагрузкой размером в несколько гигабайт, не выходя за пределы лимита памяти.

Workers не обязаны подготавливать все тело ответа перед тем, как вернуть Response. Можно использовать ReadableStream чтобы передавать тело ответа потоком после отправки строки статуса и заголовков ответа.

Worker может создать Response объект с помощью ReadableStream в качестве тела. Любые данные, переданные через 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)

A 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)) после того как ответ возвращён клиенту. В этом примере тело ответа подзапроса перекачивается в тело итогового ответа. Однако можно использовать и более сложную логику, например добавлять префикс или суффикс к телу или как-то иначе его обрабатывать.


Частые проблемы