INTEGRITY Dokumentace

Streams

Streams API je standardní webové API, které umožňuje JavaScriptu programově přistupovat ke streamům dat a zpracovávat je.

Pomocí Streams API se vyhnete ukládání velkých požadavků nebo odpovědí do paměti. Díky tomu můžete zpracovat i velmi rozsáhlá těla požadavků nebo odpovědí v rámci limitu 128 MB paměti Workeru. Je to rychlejší než načítání celého obsahu do paměti najednou, protože Worker může data zpracovávat postupně, a zvládne tak i data nebo soubory o velikosti několika gigabajtů v rámci svých paměťových limitů.

Workers nemusí připravit celé tělo odpovědi, než vrátí Response. Můžete použít ReadableStream pro streamování těla odpovědi po odeslání stavového řádku a hlaviček odpovědi.

Worker může vytvořit Response objekt pomocí ReadableStream jako tělo. Veškerá data poskytnutá prostřednictvím ReadableStream se bude klientovi streamovat postupně, jak bude k dispozici.

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 a ReadableStream.pipeTo() metodu lze použít k úpravě těla odpovědi během jeho streamování:

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"})

Tento příklad volá response.body.pipeTo(writable) ale ne await to. Je to proto, aby to neblokovalo další postup zbytku fetchAndStream() funkci. Pokračuje v asynchronním běhu, dokud není odpověď dokončena nebo dokud se klient neodpojí.

Runtime může pokračovat ve spouštění funkce (response.body.pipeTo(writable)) after a response is returned to the client. This example pumps the subrequest response body to the final response body. However, you can use more complicated logic, such as adding a prefix or a suffix to the body or to process it somehow.


Časté problémy