Workers Streams 串流處理:邊收邊送與 SSE 實戰 | Cloudflare 完整教學
在 Cloudflare Workers 中,當你要處理「很大」或「即時」的資料——代理一個上百 MB 的檔案、把 AI 模型逐字吐出的 token 即時推給前端——你不可能把整份資料塞進只有 128 MB 的記憶體。答案是 Streams(串流):讓資料像水一樣一段段流過,邊接收邊回傳,記憶體裡任何時刻只留很小的一塊。這一篇我們深入
ReadableStream/WritableStream/TransformStream三大元件、pipeThrough與pipeTo的接管、背壓(backpressure) 的自動流量控制,以及怎麼用 SSE 做 LLM 串流輸出。
前言
上一篇《Cache API》我們讓重複的請求「變快」——先查快取、命中就秒回。但你可能注意到,那一篇我一直很小心地處理 body:response.clone()、「body 只能讀一次」、不要一次讀進記憶體。這些小心翼翼的背後,其實藏著 Workers 的另一個核心機制在支撐——串流(Streams)。這一篇,我們就把它正面攤開。
先給一句話定義:Streams 是一套符合 Web 標準的 API,讓你把資料當成「一連串陸續抵達的區塊(chunk)」來處理,而不是等它整份到齊。 它有三個角色:ReadableStream(可讀串流,資料的來源端)、WritableStream(可寫串流,資料的目標端)、TransformStream(轉換串流,夾在中間邊流邊改)。Workers 的 Response body 本身就是一個 ReadableStream,fetch 回來的 body 也是——這代表串流早就無所不在,只是你以前用 await response.text() 把它「攤平」成一整塊而已。
打個比方:串流就像自來水管,一次讀進記憶體則像先把整個游泳池的水搬進你家浴缸。 如果你要把水從 A 送到 B,聰明的做法是接一條水管(ReadableStream → WritableStream),打開水龍頭讓水邊流邊送,你家浴缸(記憶體)任何時刻只裝著管子裡那一小段水;笨的做法是先把整池水全搬進浴缸、再一桶桶往 B 倒——浴缸(128 MB 記憶體)根本裝不下一個游泳池,直接溢出來(崩潰)。而如果你想在送水途中「加點東西」(過濾、加料、逐字包裝),就在水管中間接一個處理器(TransformStream),這就是 pipeThrough。水管還有個聰明的地方:如果 B 端喝水的速度變慢,水管會自動關小水龍頭,不會硬灌到爆——這就是背壓(backpressure)。
讀完本篇,你會掌握:
- 三種 Stream 與資料流動——
ReadableStream(來源)、WritableStream(目標)、TransformStream(中間轉換),以及它們怎麼串成一條管線 - 串流回應(邊收邊送)——不緩衝整份 body,直接把
ReadableStream包進Response回傳,處理遠大於記憶體的資料 pipeThrough與pipeTo、背壓——兩種接管方式的差異,以及為什麼絕不能await pipeTo- SSE 與 LLM 串流輸出——用
text/event-stream與TransformStream把 token 逐字推給前端
核心概念
三種 Stream 與資料流動示意
一條完整的串流管線,可以用一句話描述:資料從 ReadableStream 流出,(可選地)經過一或多個 TransformStream 加工,最後流進 WritableStream。
ReadableStream ──▶ [TransformStream] ──▶ [TransformStream] ──▶ WritableStream
(來源) (轉換 1) (轉換 2) (目標)
三個角色各自的職責:
| 元件 | 角色 | 你拿它做什麼 |
|---|---|---|
ReadableStream | 資料來源端 | response.body、fetch(...).body、request.body 都是它;提供 chunk 讓人讀 |
WritableStream | 資料目標端 | 接收 chunk 並寫出去;你用 writer.write(chunk) 灌資料 |
TransformStream | 中間轉換器 | 同時有一個 writable(進料口)和一個 readable(出料口);進去什麼、出來什麼由你決定 |
TransformStream 是最關鍵、也最容易搞混的角色。它其實是「一對相連的串流」:你把資料寫進它的 writable 端,經過你定義的 transform 函式加工後,從它的 readable 端流出來。所以它的典型用法是:把上游的 readable 接到它的 writable,再把它的 readable 接給下游。
// TransformStream 就是「一進一出」的一對串流
const transform = new TransformStream({
transform(chunk, controller) {
// 進來一個 chunk,你決定要輸出什麼(可改寫、可過濾、可放大)
controller.enqueue(chunk); // enqueue 就是「往 readable 端塞一塊」
},
});
// transform.writable 是進料口,transform.readable 是出料口
運作原理:pipeThrough、pipeTo 與背壓
把管線接起來,有兩個核心方法:
pipeThrough(transform):把一個ReadableStream接上一個TransformStream,回傳轉換後的新ReadableStream。因為回傳的還是ReadableStream,所以可以鏈式串接多個轉換:source.pipeThrough(a).pipeThrough(b)。pipeTo(writable):把一個ReadableStream的所有資料灌進一個WritableStream,回傳一個 Promise(在整條串流流完、或出錯時 settle)。它是管線的終點。
// pipeThrough 回傳 ReadableStream,可以繼續鏈式接
const transformed = source.pipeThrough(upperCaseTransform).pipeThrough(gzipTransform);
// pipeTo 是終點,回傳 Promise(注意:在 Workers 裡通常不要 await 它)
transformed.pipeTo(destinationWritable);
接著是整個串流機制裡最重要、也最反直覺的概念——背壓(backpressure)。
背壓是串流的「自動流量控制」機制:當下游消費資料的速度跟不上上游生產的速度時,串流會自動叫上游「慢一點」,避免資料在記憶體裡越積越多而爆掉。 想像水管接到一個喝水很慢的人,水管會自動把上游水龍頭關小,等他喝完再放。串流內部用一個「緩衝佇列 + 高水位線(high water mark)」來實現:佇列滿了就對上游施加背壓,佇列空了才放行。
這帶出一個 Workers 裡最經典的坑:絕對不要 await 一個代理場景的 pipeTo。 看這段:
// ❌ 錯誤:await pipeTo 會造成背壓死鎖
async fetch(request: Request): Promise<Response> {
const { readable, writable } = new TransformStream();
const upstream = await fetch("https://origin.example.com/big-file");
await upstream.body!.pipeTo(writable); // ← 卡死在這裡
return new Response(readable); // ← 永遠跑不到
}
為什麼死鎖?pipeTo 要「流完」才 resolve,但 readable 這端的消費者(瀏覽器)根本還沒開始讀——因為你連 Response 都還沒 return。於是:寫入端等消費端讀、消費端等你 return、你又在等 pipeTo 完成,三方互鎖。正確版本是讓 pipeTo 在背景跑,馬上 return:
// ✅ 正確:不 await pipeTo,讓它在背景流,立刻 return Response
async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
const { readable, writable } = new TransformStream();
const upstream = await fetch("https://origin.example.com/big-file");
// 不 await!讓串流在背景流動,消費者一開始讀,背壓就疏通了
ctx.waitUntil(upstream.body!.pipeTo(writable));
return new Response(readable, upstream); // 馬上把 readable 包進 Response
}
關鍵術語:identity TransformStream 與 Workers 特有串流
- identity TransformStream:
new TransformStream()不帶任何參數時,就是一個「原樣轉發」的恆等串流——進去什麼、原封不動出來什麼。它極其有用:當你想「先拿到一個writable讓自己手動灌資料、同時拿到一個readable包進Response」時,identity transform 正是那座橋。上面的代理範例就是靠它把writable(供pipeTo灌)與readable(供Response讀)配成一對。 - Workers 特有的串流:Cloudflare 在標準之外加了兩個實用工具——
FixedLengthStream(能精確設定Content-Length,適合你事先知道總長度、要讓瀏覽器顯示下載進度條的場景)與DigestStream(一邊串流一邊算 SHA-256 等雜湊,不必把整份資料讀進來就能驗證完整性)。
實作範例
概念講完,我們寫三個可執行的 Worker:串流回應大檔、TransformStream 即時改寫內容、SSE 逐字推送。
範例一:串流回應——邊收邊送、不撐爆記憶體
最基本的串流:代理一個上游的大檔案,一邊接收、一邊轉發給使用者,記憶體裡永遠只有一小段。
// src/index.ts —— 串流代理大檔,不把整份 body 讀進記憶體
export default {
async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
// 向上游要一個大檔(可能上百 MB)
const upstream = await fetch("https://origin.example.com/huge-video.mp4");
if (!upstream.body) {
return new Response("上游沒有 body", { status: 502 });
}
// 關鍵:直接把上游的 ReadableStream(upstream.body)包進 Response 回傳。
// Workers 會邊從上游收 chunk、邊往使用者送,記憶體只留很小一段。
// 第二個參數 upstream 帶入原 status / headers。
return new Response(upstream.body, upstream);
},
} satisfies ExportedHandler<Env>;
這裡的精髓是你完全沒碰 body 的內容——沒有 await upstream.text()、沒有 arrayBuffer()。upstream.body 是個 ReadableStream,把它交給 Response,Workers 自動處理「邊收邊送」與背壓。如果你手癢寫成 const buf = await upstream.arrayBuffer(),一個 200 MB 的檔案就會把 128 MB 記憶體撐爆而崩潰。
範例二:TransformStream 即時改寫串流內容
現在我們要在轉發途中「加工」——把上游回傳的文字內容逐塊轉成大寫,示範 pipeThrough 與 TransformStream 的組合。真實場景可能是:即時翻譯、遮蔽敏感字、注入內容。
// src/index.ts —— 用 TransformStream 邊流邊改寫(此例:轉大寫)
export default {
async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
const upstream = await fetch("https://origin.example.com/article.txt");
if (!upstream.body) return new Response("no body", { status: 502 });
// TransformStream 的 transform 對「每個 chunk」呼叫一次。
// chunk 是 Uint8Array(位元組),要先 decode 成文字、改寫、再 encode 回位元組。
const decoder = new TextDecoder();
const encoder = new TextEncoder();
const upperCaseStream = new TransformStream<Uint8Array, Uint8Array>({
transform(chunk, controller) {
const text = decoder.decode(chunk, { stream: true }); // stream:true 處理跨 chunk 的多位元組字元
controller.enqueue(encoder.encode(text.toUpperCase()));
},
});
// pipeThrough 把上游串流接上轉換器,回傳「轉換後的 ReadableStream」
const transformed = upstream.body.pipeThrough(upperCaseStream);
// 直接把轉換後的串流包進 Response——全程沒有把整篇文章讀進記憶體
return new Response(transformed, {
headers: { "Content-Type": "text/plain; charset=utf-8" },
});
},
} satisfies ExportedHandler<Env>;
幾個關鍵細節:
- chunk 是位元組(
Uint8Array)不是字串:所以要用TextDecoder/TextEncoder在文字與位元組間轉換。decoder.decode(chunk, { stream: true })的stream: true很重要——它讓 decoder 記住「跨 chunk 被切斷的多位元組字元(如中文、emoji)」,等下一塊補齊,避免亂碼。 controller.enqueue(...):把加工後的 chunk 「塞進」轉換串流的readable端往下游流。pipeThrough回傳的還是ReadableStream:所以你可以繼續.pipeThrough(anotherTransform)串更多層。
範例三:Server-Sent Events(SSE)——LLM 逐字串流輸出
這是串流最迷人的應用:把 AI 模型逐字吐出的 token 即時推給前端,讓使用者看到「打字機」效果,而不是苦等整段回答生成完。做法是回傳一個 text/event-stream 的串流回應,用 TransformStream 的 writable 端手動逐筆寫入 SSE 格式的事件。
SSE 的資料格式極簡單:每一筆事件是 data: <內容>\n\n(兩個換行結尾代表一筆事件結束)。
// src/index.ts —— 用 SSE 逐字串流輸出(模擬 LLM token 流)
export default {
async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
// identity TransformStream:拿到一對相連的 writable / readable
const { readable, writable } = new TransformStream();
const writer = writable.getWriter();
const encoder = new TextEncoder();
// 把「逐字產生 + 寫入」放進背景執行,不阻擋 Response 回傳
ctx.waitUntil(
(async () => {
const tokens = ["Cloudflare", " Workers", " 讓", " 串流", " 變得", " 超簡單", "。"];
try {
for (const token of tokens) {
// SSE 格式:data: <內容> 後接兩個換行
const event = `data: ${JSON.stringify({ token })}\n\n`;
await writer.write(encoder.encode(event)); // await 這裡的 write 沒問題(背壓自然節流)
await sleep(120); // 模擬 token 之間的生成間隔
}
// 送一個結束訊號給前端
await writer.write(encoder.encode(`data: [DONE]\n\n`));
} finally {
await writer.close(); // 務必關閉 writer,否則連線永遠不結束
}
})(),
);
// SSE 必備的三個 header
return new Response(readable, {
headers: {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
"Connection": "keep-alive",
},
});
},
} satisfies ExportedHandler<Env>;
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
前端接收極為簡單,用瀏覽器內建的 EventSource:
// 前端:用 EventSource 接收 SSE 串流
const source = new EventSource("/stream");
source.onmessage = (e) => {
if (e.data === "[DONE]") {
source.close(); // 收到結束訊號就關閉
return;
}
const { token } = JSON.parse(e.data);
document.getElementById("output").textContent += token; // 逐字追加,打字機效果
};
如果你真的接 Workers AI 的 LLM,env.AI.run(model, { ..., stream: true }) 會直接回傳一個 SSE 格式的 ReadableStream,你甚至可以直接 return new Response(aiStream, { headers: { "Content-Type": "text/event-stream" } }),或用 pipeThrough 在中間插入自己的轉換(例如過濾敏感詞)。原理和上面完全一致。
常見錯誤與最佳實踐
坑一:一次把整份 body 讀進記憶體(await response.text() / arrayBuffer()),大檔直接爆記憶體。
這是串流場景最致命的錯誤。Workers 記憶體上限 128 MB,一旦你對一個可能很大的 body 呼叫 await response.text()、await response.arrayBuffer()、await response.json(),就是要求把整份資料緩衝進記憶體——200 MB 的檔案當場 OOM 崩潰;就算沒爆,使用者也得等你「全部讀完、全部處理完」才收到第一個位元組(TTFB 極差)。正確做法:能串流就串流。代理直接 new Response(upstream.body, upstream);要改寫就 body.pipeThrough(transform)。只有在 body 確定很小、且你必須拿到完整內容才能運算(如驗整份簽章)時,才緩衝。
坑二:await pipeTo(...) 造成背壓死鎖。
如前面所述,代理場景裡 readable 的消費者是使用者的瀏覽器,而它要等你 return Response 後才開始讀。你若 await pipeTo,就會卡在「等串流流完」——但串流因為沒人讀而永遠流不完,死鎖。正確做法:不要 await 代理用的 pipeTo。讓它在背景跑,馬上 return new Response(readable);若需要在管線結束後做收尾,用 ctx.waitUntil(source.pipeTo(writable)),而不是 await。
坑三:手動用 writer 卻忘了 writer.close(),連線永遠掛著。
當你用 writable.getWriter() 手動逐筆寫入(像 SSE 範例那樣),寫完後一定要 await writer.close()。忘了關,串流的 readable 端就永遠等不到「結束」訊號,使用者的連線會一直懸著(瀏覽器一直轉圈圈、EventSource 不觸發結束、連線數被佔滿)。正確做法:把寫入邏輯包在 try / finally,在 finally 裡 await writer.close(),確保就算中途出錯也會正確關閉。出錯時也可以用 writer.abort(err) 帶著錯誤關閉。
坑四:串流回應忘了設正確的 headers。
串流本身流得再好,headers 不對前端也吃不到。最常見的是 SSE:必須設 Content-Type: text/event-stream(少了它瀏覽器不會當成事件流)、Cache-Control: no-cache(避免中間層快取事件流)。若你事先知道總長度、想讓瀏覽器顯示下載進度,用 Workers 的 FixedLengthStream 精確帶上 Content-Length。一般串流則不要自己亂設 Content-Length(你邊送邊算根本不知道總長),交給 Workers 用 chunked transfer 處理即可。
最佳實踐小結: 記住一個心法——「不要把水倒進浴缸,接一條水管就好」。預設就用串流思維:body 直接 pipeThrough / pipeTo,需要手動控制時用 identity TransformStream 拿 writer、寫完 close()、背景任務用 ctx.waitUntil、永遠不 await 代理的 pipeTo。這幾條守住,你的 Worker 就能在 128 MB 的小小記憶體裡,優雅地處理遠比它大的資料、以及即時逐字的 AI 輸出。
小結
上一篇《Cache API》我們讓「重複的請求」變快——先查快取、命中秒回。這一篇,我們讓「大」與「即時」變得可能,深入了 Workers 的串流機制:
- 三種 Stream——
ReadableStream(來源,如response.body)、WritableStream(目標,用writer.write)、TransformStream(中間邊流邊改,一進一出)。 - 串流回應——直接把
ReadableStream包進Response,邊收邊送、不緩衝整份 body,200 MB 的檔案也能在 128 MB 記憶體下代理。 pipeThrough與pipeTo、背壓——pipeThrough接轉換器並回傳新ReadableStream(可鏈式);pipeTo是終點回傳 Promise;代理場景絕不await pipeTo,否則背壓死鎖。- identity TransformStream——
new TransformStream()原樣轉發,是「拿一對writable/readable」的橋,手動灌資料 + 包進 Response 的關鍵。 - SSE 與 LLM 串流——
text/event-stream+data: ...\n\n格式,用writer逐筆推、記得close(),做出打字機般的即時 AI 輸出。
串流讓資料「流」了起來,但它是單向的——伺服器往下推。當你需要客戶端也能即時往回送訊息,做真正的雙向即時互動(聊天室、多人協作、線上遊戲),就得換一種武器。下一篇《WebSockets》,我們會深入 Workers 的 WebSocketPair、accept / send / close,以及 101 升級回應怎麼建立一條全雙工連線,並看它如何搭配 Durable Objects 管理長連線。串流讓你「一直推」,WebSocket 則讓雙方「隨時互相說話」。
想深入官方 Streams API 細節,可以隨時參考 Cloudflare 官方 Streams API 文件。準備好讓你的 Worker 邊收邊送、逐字串流了嗎?我們下一篇《WebSockets》見。