後端面試:如何用 Node.js stream.compose 設計可取消的資料流水線?
題幹與適用場景
一個 Node.js 服務接收大檔案上傳,需要依序完成解壓縮、病毒掃描、格式驗證和物件儲存寫入。舊實作手動監聽 data 事件,偶爾記憶體上升、用戶端斷線後工作仍繼續,以及中間步驟失敗時暫存檔沒有清理。
請使用 stream.compose() 或等價的 pipeline 組合方式,說明可讀、可寫和 Transform 階段如何連接,如何利用背壓、AbortSignal、錯誤傳播和最終清理控制一次工作。
面試官考察點
- 能否區分
pipe、pipeline和compose的邊界與生命週期。 - 能否解釋背壓如何限制生產速度,而不是簡單提高佇列上限。
- 能否讓取消、例外和用戶端斷線沿整條流水線傳播。
- 能否處理非同步生成器、資源釋放、冪等寫入和可觀測性。
回答前需要釐清的問題
- 輸入來自 HTTP request、檔案還是物件儲存 SDK?輸出是否支援分片和重試?
- 每個階段是 Node stream、Web Streams、AsyncIterable 還是普通函式?
- 病毒掃描和驗證是否會產生外部程序、暫存檔或資料庫狀態?
- 取消後物件儲存是否允許中止 multipart upload,工作是否需要恢復?
30 秒回答
我會把每個階段定義成可讀、可寫或 Transform/AsyncIterable,並用 stream.compose() 組合成一個 Duplex,再由 pipeline 驅動到最終目的地。生產者只在下游可接受時繼續生成,讓背壓限制記憶體。把請求斷開、逾時和業務取消合併為一個 AbortSignal,傳給支援 signal 的階段;任何錯誤都讓整條鏈失敗。暫存檔、子程序和 multipart upload 在 finally 或 abort handler 中清理,寫入使用冪等鍵。指標記錄吞吐量、佇列長度、峰值 RSS、取消原因和清理結果。
分步驟深入解答
先定義每個階段的契約
明確輸入輸出型別、chunk 大小、是否允許 null、是否會阻塞,以及階段完成時誰擁有資源。非同步生成器必須消費 source 並按需 yield;普通函式不能偷偷把整個檔案讀進記憶體。
用 compose 組合可重用階段
compose 把多個 stream、AsyncIterable 或函式連接成新的 Duplex;它內部按 pipeline 語義處理相鄰階段。示例:
import { compose } from 'node:stream';
async function* validate(source) {
for await (const chunk of source) {
checkChunk(chunk);
yield chunk;
}
}
const processing = compose(decompress(), validate, scan());實際程式還要把 processing 接到目標 writable,並統一監聽完成與錯誤,而不是每段各自吞掉例外。
讓背壓成為預設控制面
下游 writable 尚未準備好時,Readable/Transform 應暫停繼續生產。不要在 data 回呼中無界 push,也不要只靠增大 high water mark 掩蓋消費者變慢。壓測時記錄每階段佇列、吞吐量和 RSS,確認最慢階段決定整體速度。
傳播取消與用戶端斷線
把請求斷開、deadline、人工取消合併到一個 AbortController。將 signal 傳給支援它的 compose 階段和外部 SDK;abort 後停止讀取、銷毀下游並等待關閉完成。取消是正常結束原因,日誌需區分 abort、業務失敗和網路錯誤。
統一錯誤與資源清理
讓一個擁有者呼叫 pipeline/compose,統一接收第一個錯誤並觸發銷毀。暫存檔、掃描程序、socket 和 multipart upload 必須在成功、失敗和取消三條路徑都釋放;清理失敗要告警並帶工作 ID,不能覆蓋原始錯誤。
設計冪等與觀測指標
用上傳 ID 和階段版本生成冪等鍵;物件儲存寫入完成後再提交業務狀態。記錄輸入位元組、輸出位元組、處理時長、峰值 RSS、取消原因、錯誤階段和清理耗時。重試只從可恢復邊界重新開始,避免重複寫入。
高品質示範回答
我會先把上傳、解壓縮、掃描、驗證和寫入定義為有明確輸入輸出與資源所有權的階段,再用 compose 組合可讀、可寫或非同步迭代器,最後交給統一的 pipeline 驅動。背壓由下游消費速度控制,禁止在 data 事件中無界快取。請求斷開、逾時和人工取消共享一個 AbortSignal,並傳遞給支援 signal 的階段和儲存 SDK。統一擁有者處理第一個錯誤,銷毀整條鏈;暫存檔、掃描程序和 multipart upload 在成功、失敗、取消路徑都清理。寫入用冪等鍵,指標覆蓋吞吐量、佇列、RSS、取消原因和清理結果,重試只從安全邊界恢復。
常見錯誤
- 把整條流收集成 Buffer,再聲稱使用了 compose。
- 只監聽最終 writable 的
error,遺漏非同步生成器或外部程序錯誤。 - 讓用戶端斷線後來源流繼續讀取並寫入物件儲存。
- 用無限增大 high water mark 解決背壓問題。
- 只在成功路徑刪除暫存檔,忽略 abort 和例外。
- 重試沒有冪等鍵,導致物件重複或業務狀態重複提交。
追問及應對
compose 和 pipeline 如何分工?
compose 負責把階段組合成可重用的 Duplex;pipeline 負責驅動端到端連接、傳播錯誤並等待關閉。兩者可以組合使用,但所有權和錯誤處理應集中。
非同步生成器拋錯會發生什麼?
錯誤應讓組合流失敗並銷毀相鄰階段。呼叫方仍需等待 pipeline 完成清理,並記錄實際失敗階段,不能只依賴程序未退出。
如何驗證背壓真的生效?
用比消費者更快的輸入和受控慢寫入壓測,觀察佇列長度與 RSS 是否有界,並確認吞吐量接近最慢階段而非無限快取。
取消後能否直接重試?
先確認所有資源已關閉、暫存狀態可識別,再從冪等且可恢復的邊界重試。不可中斷的外部寫入應先完成補償或標記未知狀態。