Streams API ↗ は、JavaScript からデータのストリームにプログラムでアクセスし、処理できる Web 標準 API です。
- ReadableStream
- ReadableStream BYOBReader
- ReadableStream DefaultReader
- TransformStream
- WritableStream
- WritableStream DefaultWriter
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)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 仕様 ↗
- Worker のコードは ES modules 構文 で書くと、より最適化された体験になります。