Workers Streams 串流處理:邊收邊送與 SSE 實戰 | Cloudflare 完整教學

2026/08/07
Workers Streams 串流處理:邊收邊送與 SSE 實戰 | Cloudflare 完整教學

Cloudflare Workers 中,當你要處理「很大」或「即時」的資料——代理一個上百 MB 的檔案、把 AI 模型逐字吐出的 token 即時推給前端——你不可能把整份資料塞進只有 128 MB 的記憶體。答案是 Streams(串流):讓資料像水一樣一段段流過,邊接收邊回傳,記憶體裡任何時刻只留很小的一塊。這一篇我們深入 ReadableStream / WritableStream / TransformStream 三大元件、pipeThroughpipeTo 的接管、背壓(backpressure) 的自動流量控制,以及怎麼用 SSELLM 串流輸出

前言

上一篇《Cache API》我們讓重複的請求「變快」——先查快取、命中就秒回。但你可能注意到,那一篇我一直很小心地處理 body:response.clone()、「body 只能讀一次」、不要一次讀進記憶體。這些小心翼翼的背後,其實藏著 Workers 的另一個核心機制在支撐——串流(Streams)。這一篇,我們就把它正面攤開。

先給一句話定義:Streams 是一套符合 Web 標準的 API,讓你把資料當成「一連串陸續抵達的區塊(chunk)」來處理,而不是等它整份到齊。 它有三個角色:ReadableStream(可讀串流,資料的來源端)、WritableStream(可寫串流,資料的目標端)、TransformStream(轉換串流,夾在中間邊流邊改)。Workers 的 Response body 本身就是一個 ReadableStreamfetch 回來的 body 也是——這代表串流早就無所不在,只是你以前用 await response.text() 把它「攤平」成一整塊而已。

打個比方:串流就像自來水管,一次讀進記憶體則像先把整個游泳池的水搬進你家浴缸。 如果你要把水從 A 送到 B,聰明的做法是接一條水管(ReadableStreamWritableStream),打開水龍頭讓水邊流邊送,你家浴缸(記憶體)任何時刻只裝著管子裡那一小段水;笨的做法是先把整池水全搬進浴缸、再一桶桶往 B 倒——浴缸(128 MB 記憶體)根本裝不下一個游泳池,直接溢出來(崩潰)。而如果你想在送水途中「加點東西」(過濾、加料、逐字包裝),就在水管中間接一個處理器(TransformStream),這就是 pipeThrough。水管還有個聰明的地方:如果 B 端喝水的速度變慢,水管會自動關小水龍頭,不會硬灌到爆——這就是背壓(backpressure)

讀完本篇,你會掌握:

  • 三種 Stream 與資料流動——ReadableStream(來源)、WritableStream(目標)、TransformStream(中間轉換),以及它們怎麼串成一條管線
  • 串流回應(邊收邊送)——不緩衝整份 body,直接把 ReadableStream 包進 Response 回傳,處理遠大於記憶體的資料
  • pipeThroughpipeTo、背壓——兩種接管方式的差異,以及為什麼絕不能 await pipeTo
  • SSE 與 LLM 串流輸出——用 text/event-streamTransformStream 把 token 逐字推給前端

核心概念

三種 Stream 與資料流動示意

一條完整的串流管線,可以用一句話描述:資料從 ReadableStream 流出,(可選地)經過一或多個 TransformStream 加工,最後流進 WritableStream

ReadableStream ──▶ [TransformStream] ──▶ [TransformStream] ──▶ WritableStream
   (來源)              (轉換 1)              (轉換 2)              (目標)

三個角色各自的職責:

元件角色你拿它做什麼
ReadableStream資料來源端response.bodyfetch(...).bodyrequest.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 TransformStreamnew 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 即時改寫串流內容

現在我們要在轉發途中「加工」——把上游回傳的文字內容逐塊轉成大寫,示範 pipeThroughTransformStream 的組合。真實場景可能是:即時翻譯、遮蔽敏感字、注入內容。

// 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 的串流回應,用 TransformStreamwritable 端手動逐筆寫入 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,在 finallyawait 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 TransformStreamwriter、寫完 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 記憶體下代理。
  • pipeThroughpipeTo、背壓——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 的 WebSocketPairaccept / send / close,以及 101 升級回應怎麼建立一條全雙工連線,並看它如何搭配 Durable Objects 管理長連線。串流讓你「一直推」,WebSocket 則讓雙方「隨時互相說話」。

想深入官方 Streams API 細節,可以隨時參考 Cloudflare 官方 Streams API 文件。準備好讓你的 Worker 邊收邊送、逐字串流了嗎?我們下一篇《WebSockets》見。

BenZ Software Developer

熱愛技術的軟體開發者,在這裡分享程式開發經驗與學習筆記。

本週主打

AI 自動化入門包

你每天手動在做的那些煩事,其實 AI 可以自己跑。這份給你 10 個照著做就會的自動化工作流 + 50 個複製即用的提示詞,不用會寫程式。

看看這個產品 →