批次處理:MapReduce 的精神、join 的兩條路,與不可變輸入的美德

· tech

#distributed-systems#book-notes#data-engineering

📑 目錄

Part II 在單一系統內把一致性守住了;Part III 的主題換成資料在系統之間流動——而最古老、也最可靠的流動方式,是批次(batch)。DDIA 講 MapReduce 的切入點很別緻:先講 Unix 哲學——小工具、各做好一件事、用統一的介面(檔案與串流)用 pipe 接起來。然後一句話點題:MapReduce 就是跨一千台機器的 Unix pipe——輸入不可變、輸出寫成新檔案、工具之間靠統一格式銜接。這個精神,比 MapReduce 本身活得久得多。

MapReduce 解剖:map、shuffle、reduce——貴的是中間那步

以「算每個網址的點擊數」為例,三步走完:

輸入(不可變) log 分片 1 log 分片 2 log 分片 3 ① map(就地、平行) (/a,1)(/b,1)… (/a,1)(/c,1)… (/b,1)(/a,1)… 逐筆抽 (key, value),不搬資料 ② shuffle(按 key 重分發) 同 key 跨網路聚到同一台+排序 唯一大搬家的一步=最貴 ③ reduce(整組聚合) /a:(1,1,1)→ /a=3 /b:(1,1)→ /b=2 … 輸出:寫成「新」檔案 map 就地跑(把運算搬到資料旁)· shuffle 是唯一的大搬家,也是一切成本所在 groupBy、join、去重……凡是「同 key 要相聚」的操作,背後都是一次 shuffle
map 在資料所在的機器就地逐筆抽出 (key, value)——運算搬到資料旁,不搬資料;shuffle同一個 key 的資料跨網路聚到同一台並排序——整個作業唯一的大搬家,也是最貴的一步;reduce 對聚齊的每組 key 做聚合,輸出寫成新檔案(輸入永遠不動)。Spark 的 stage 邊界、效能調校的一半功夫,全在 shuffle 這一步——凡是「同 key 要相聚」的操作(groupBy、join、去重),背後都是一次 shuffle

批次 join 的兩條路:要嘛都搬,要嘛帶小抄

批次世界最常見的重活是 join(點擊 log join 使用者表)。單機的 join 演算法我在 SQL 系列講過;分散式版本的核心問題變成:兩份資料散在不同機器上,同 key 的列要怎麼相遇? 兩條路:

reduce-side(sort-merge):都搬 點擊 log(大) 使用者表(大) 兩邊都按 user_id shuffle 同 key 在同一台 reducer 相遇 ✓ 通用:不需要任何前提 ✗ 兩邊都大搬家,最重 map-side(broadcast):帶小抄 點擊 log(大) 使用者表(小) 整份複製到每台 mapper 1就地查小抄 join mapper 2就地查小抄 join ✓ 完全不 shuffle,快非常多 ✗ 前提:小表裝得進記憶體 Spark 的 broadcast join 就是右邊這條——判準只有一句:「小表,裝得進記憶體嗎?」
Reduce-side join:兩邊都按 join key 做 shuffle,同 key 的列在同一台 reducer 相遇、sort-merge 合併——通用(不需要任何前提),但兩邊都大搬家、最重。Map-side broadcast join:一邊夠小,就把它整份複製到每台 mapper 當「小抄」,大表就地查表完成 join——完全跳過 shuffle,快非常多,但小表必須裝得進記憶體。Spark 的 broadcast join 正是這條路的直系後代;選哪條,判準就一句:小的那邊,裝得進記憶體嗎?

MapReduce 本身有個大缺點:每個作業的輸出都要完整落地到 HDFS,下個作業再讀回來——一條十步的管線就要寫十次、讀十次磁碟。後來的資料流引擎(Spark、Flink)就是針對這點進化:把整條管線畫成 operator 的 DAG、中間結果盡量留在記憶體、失敗靠 lineage 重算局部。MapReduce 這個「產品」被取代了,但它的精神——分區平行、把運算搬到資料旁、shuffle 集中成本——原封不動活在每個現代引擎裡。

批次最被低估的美德:輸入不可變,錯了可以重來

DDIA 這章有個容易被跳過、我卻認為最重要的觀察:批次處理繼承了 Unix 最好的品格——輸入唯讀、輸出寫到新地方。這帶來一種被作者稱為「人為容錯(human fault tolerance)」的能力:程式寫錯了、邏輯有 bug、昨天的報表算壞了——修好程式,對著原封不動的輸入重跑一次就好。錯誤不會累積、不會污染源頭,最壞情況就是浪費一輪運算。對照之下,直接 UPDATE 資料庫的管線,一個 bug 就把唯一的真相改壞了,還原得靠備份和眼淚。可重跑,是批次給資料工程最大的禮物——Medallion 保留不可變的 Bronze 層Airflow 的冪等與 backfill,全是這個美德的現代化身。

反思

「人為容錯」是我最想幫批次講的一句話

大家嫌批次舊、嫌它慢,但做了幾年資料我最感激它的,恰恰是這個不性感的美德:人一定會犯錯,而批次讓錯誤變得可逆。 邏輯寫錯?修好重跑。昨天的維度表壞了?從 Bronze 重建。這種「錯了可以回到出發點」的安全感,streaming 世界要花十倍力氣才換得到(事件流過就過了)。所以我給團隊的紀律一直是:原始層唯讀、永遠保留、所有下游都當成「可以隨時燒掉重蓋」的衍生品——這正是 Medallion 的靈魂,而它的理論根據就在這章。系統容錯靠副本,人為容錯靠不可變;後者被討論得太少,出的事卻更多。

Shuffle 是分散式運算的房租——看懂它,調校就有了地圖

MapReduce 的三步裡,map 和 reduce 都「就地」,只有 shuffle 在搬資料——而分散式運算的成本,幾乎全部住在這一步。想通之後,一堆散落的實戰知識瞬間歸位:Spark 的 stage 為什麼在 shuffle 邊界切開、data skew 為什麼要命(某個 key 的資料全擠到一台)、broadcast join 為什麼快(整個跳過 shuffle)、為什麼 groupBy 前先 filter 能省大錢(搬之前先減量)。我現在看任何慢掉的批次作業,第一個問題永遠是:它 shuffle 了多少資料?能不能少搬、早減量、或乾脆帶小抄不搬? 這一問,就是調校的一半。

工具死了,哲學還活著——學東西要學到那一層

MapReduce 作為產品已經退役,但這章讀來毫不過時,因為它教的是三個還活著的思想:把運算搬到資料旁(而不是反過來)、用不可變的檔案當工具之間的統一介面(Unix pipe 的放大版)、把跨機器的複雜度集中到 shuffle 一個地方。Spark 是它們的新殼,dbt 的 model 鏈、Medallion 的分層,骨子裡也是「每步讀不可變輸入、寫新輸出」的同一套。這再次驗證我學技術的順位:API 一兩年就換,架構三五年一改,但這種層級的設計思想,一用二十年。 讀 DDIA 這種書的意義就在這——它讓你在下一個新工具發布時,一眼認出「喔,又是那個思想換了件衣服」。下一篇講它的另一半:當資料不再是「一批一批」,而是「一直來」——串流。