Skip to content

非公式本サイトは非公式の日本語ドキュメントであり、Cloudflare 公式サイトではありません。最新情報はdevelopers.cloudflare.comをご確認ください。

Streams

最終更新 Markdown で表示Agent セットアップ

Streams API は、JavaScript からデータのストリームにプログラムでアクセスし、処理できる Web 標準 API です。

Streams API を使うと、大きなリクエストやレスポンスをメモリにバッファせずに済みます。Worker の 128 MB メモリ制限内で、非常に大きなリクエストボディやレスポンスボディをパースできます。ペイロード全体をメモリに載せるより速く、データを段階的に処理し始められるため、メモリ制限内でマルチギガバイトのペイロードやファイルも扱えます。

Worker は 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))を実行し続けられます。この例はサブリクエストのレスポンスボディを最終的なレスポンスボディへ送ります。ただし、ボディの前後に文字列を足したり、何らかの処理をしたりする、より複雑なロジックも使えます。


よくある問題


関連リソース

役に立ちましたか?