跳至主要內容
論文精煉 · Big data processing

Spark: Cluster Computing with Working Sets

Resilient distributed dataset:可快取、能靠 lineage 重建的分散式集合,讓叢集把工作集留在記憶體中反覆重用。

作者Matei Zaharia、Mosharaf Chowdhury、Michael J. Franklin、Scott Shenker 等人,University of California, Berkeley 發表於HotCloud 2010(第 2 屆 USENIX Hot Topics in Cloud Computing 研討會) 年份2010
閱讀原始論文 PDF 所有論文

一口氣講完 — 整篇論文的濃縮

MapReduce 與 Dryad 讓一般商用叢集變得好用,但兩者都建立在非循環資料流模型上,每個 job 都得重新從磁碟載入輸入資料。對本論文鎖定的兩類工作負載而言,這是致命傷:迭代式機器學習會對同一份資料掃描數十次,而互動式分析中每一條 Hive 或 Pig 查詢都要付出數十秒的 MapReduce 啟動與磁碟 I/O 成本。Spark 的答案是 resilient distributed dataset(RDD)——一種唯讀、切成多個分割區的集合,它的 handle 本身就帶著足夠資訊,可以從可靠儲存中的資料重算任何遺失的分割區,因此使用者可以用 cache 提示把它釘在記憶體裡反覆重用,而不需要複本或 checkpoint。圍繞 RDD,論文再加上兩種受限的共享變數(廣播變數與只能累加的 accumulator),以及一套 Scala 整合:把使用者的閉包序列化後送到 worker 執行,並改造 Scala 直譯器,使互動式叢集 shell 成為可能。在 20 台 EC2 節點上,29 GB 資料的 logistic regression 從 Hadoop 的每次迭代 127 秒降到第一次之後每次僅 6 秒,約快 10 倍;39 GB 的 Wikipedia dump 也能以 0.5 到 1 秒的延遲互動查詢。

在這篇論文之前 — 它所降落的世界

到 2010 年,MapReduce 模型已經成為在不可靠的商用叢集上運算的預設方式,Dryad 與 Map-Reduce-Merge 則進一步擴充了它支援的資料流形狀。這些系統的交換條件都一樣:使用者寫出一張非循環的運算子圖,系統則負責考量資料位置的排程、負載平衡與容錯,全程不需使用者介入。但只要應用會重複碰同一份資料,這個交換條件就嚴重漏水。跑梯度下降的 Hadoop 使用者必須把每一次迭代寫成一個獨立 job,每次都從 HDFS 重讀整份資料;透過 Pig 或 Hive 做臨時 SQL 探索的分析師,每條查詢都要等上數十秒,因為每條查詢都是一個從磁碟讀資料的全新 MapReduce job。看似顯而易見的替代方案——分散式共享記憶體(DSM)——雖然已被研究二十年,卻是用 checkpoint 來容錯,失敗時整支程式必須回捲,而且即使什麼都沒壞也要付出額外成本;Twister 雖能讓靜態資料跨迭代留在記憶體,卻完全沒有容錯,而且只允許一個 map 函式與一個 reduce 函式。

問題 — 當時真正壞掉的地方

  • MapReduce、Dryad 及其變體都建立在非循環資料流模型上,無法有效表達那些會跨多個平行操作重複使用同一份工作集的應用。
  • 迭代式機器學習演算法會反覆對同一份資料套用同一個函式來最佳化參數,若把每次迭代寫成獨立的 MapReduce 或 Dryad job,就必須每次都從磁碟重新載入資料,帶來顯著的效能損失。
  • 透過 Pig、Hive 這類 SQL 介面做臨時性探索查詢,每條查詢都要承受數十秒的延遲,因為每條查詢都是獨立的 MapReduce job,是從磁碟讀資料,而不是讀取已經載入叢集記憶體中的資料集。
  • 分散式共享記憶體在通用性上足以解決這個問題,但既有的 DSM 系統靠 checkpoint 容錯:失敗時程式必須回捲到 checkpoint,而不是只重算遺失的部分,而且即使沒有任何節點故障也要付出額外開銷。
  • Twister 是把 MapReduce 迭代化最接近的前作,它讓靜態資料留在長生命週期的 map task 中,但完全沒有實作容錯,而且只提供一個 map 函式與一個 reduce 函式,程式無法定義多個資料集並在它們之間交替執行操作。
  • 當時沒有任何高效的通用程式語言能夠以互動方式在叢集上處理大型資料集,因為在直譯器中逐行輸入的閉包,沒有辦法帶著它捕捉到的狀態送到 worker 機器上執行。

核心概念 — 主要貢獻,以及它們為何成立

Resilient distributed dataset

RDD 是一個唯讀、切分到多台機器上的物件集合,任何分割區遺失都能被重建。關鍵在於這些元素不必真的存在於任何實體儲存上:RDD 的 handle 帶著足夠的資訊,可以從可靠儲存中的資料把整個資料集算出來,所以這個物件與其說是容器,不如說是一份食譜。正因如此,把工作集主要放在記憶體裡才是安全的——記憶體不見了,並不會失去任何無法重新生成的東西。作者也明講 RDD 並不是通用的共享記憶體抽象,而是刻意受限、落在表達力與擴展性/可靠性之間的甜蜜點。

用 lineage 取代 checkpoint

每個 RDD 物件都持有指向父資料集的指標,以及父資料集被如何轉換的描述,因此這條資料集物件鏈本身就是復原用的紀錄。當某個分割區遺失時,Spark 只針對該分割區、對著父資料集重播那串轉換來重新推導出來,而不是把整個運算回捲到 checkpoint。由於各分割區彼此獨立,多個遺失的分割區可以在不同節點上平行重建,而且在沒有節點故障時完全沒有額外開銷。lineage 之所以能被便宜地記錄下來,正是因為程式模型受限:只有粗粒度且具決定性的轉換,每個資料集只要幾個位元組的中繼資料就能描述整段推導過程。

持久化只是提示,不是保證

RDD 預設是延遲(lazy)且短暫(ephemeral)的:分割區在平行操作需要時才依需求具現化,用完就從記憶體丟掉。cache 動作讓資料集維持延遲,但提示系統在第一次算出來之後應該留在記憶體中;save 動作則強制求值並寫入分散式檔案系統。cache 這個提示並沒有約束力,而這正是重點:如果叢集記憶體不足以容納所有分割區,Spark 就在用到時重算,程式因此只會變慢而不會壞掉。作者把這個設計比喻成某種鬆散版的虛擬記憶體,並把整體目標定義為讓使用者在儲存成本、存取速度、遺失機率與重算成本之間自行取捨。

兩種受限的共享變數

傳給 map、filter、reduce 的閉包,其捕捉到的自由變數平常會被複製到每個 worker,但對兩種常見樣式而言這既浪費又不合適。廣播變數包住一份大型唯讀資料,例如查表或評分矩陣,並保證它只會被複製到每個 worker 一次,而不是隨每個閉包一起打包送出,而且可以跨多個平行操作重複使用,不必綁定在單一 job 上。accumulator 則是 worker 只能以結合律運算「加」進去、而且只有 driver 能讀取的變數,只要型別具備 add 運算與 zero 值就能定義。正是這種只加不減的語意,讓 accumulator 容易做到容錯:被重新執行的 task 所貢獻的更新可以被丟棄或恰好套用一次。

用真正的語言運送閉包

Spark 以 Scala 實作,提供類似 DryadLINQ 精神的函數式介面,但它不是捕捉運算式樹,而是直接運送已編譯的閉包。Scala 閉包本身就是普通的 Java 物件,用 Java serialization 就足以把一段運算送到另一台機器,這也是整套系統能以很小的實作量疊在 Mesos 之上的原因。回報就是可程式性:因為 Scala 的 for 語法其實是 foreach 的語法糖,而 accumulator 又支援多載的 += 運算子,論文中的 logistic regression 範例只比序列版本多三行。加上型別推論,這些程式讀起來就像一般的 Scala 集合程式碼,儘管每個操作其實都是分散式 job。

叢集上的互動式 shell

Scala 直譯器會為使用者輸入的每一行編譯出一個類別,其中含有一個 singleton 物件,保存該行的變數,並在建構子中執行該行程式碼。Spark 對這個機制做了兩處修改,讓同一套做法能在叢集上運作,使用者因此可以在提示字元下定義 RDD、函式、變數與類別,並直接用在平行操作中。作者宣稱 Spark 是第一個讓高效通用程式語言能以互動方式在叢集上處理大型資料集的系統。搭配已快取的 RDD,互動體驗發生了本質改變:39 GB 的資料集載入一次之後,後續即使全量掃描也能在一秒內回應,感覺就像在操作本機資料。

運作方式 — 具體的機制

建立與轉換 RDD

取得 RDD 的方式恰好只有四種:從 HDFS 這類共享檔案系統讀檔;把 driver 中的 Scala 集合 parallelize,也就是切成若干片段送到各節點;轉換既有的 RDD;或改變既有 RDD 的持久化設定。最基本的轉換是 flatMap,它把每個元素送進型別為 A 到 List of B 的使用者函式,語意等同於 MapReduce 中的 map;型別為 A 到 B 的 map,以及依述詞篩選元素的 filter,都可以用 flatMap 表達。持久化則透過兩個動作改變:cache 讓資料集維持延遲但標記為應保留在記憶體,save 則求值並寫入分散式檔案系統,之後的操作都改讀已存檔的版本。這一切都發生在 driver 程式中,由它實作應用的高階控制流程,並平行啟動各項操作。

平行操作與 driver

有三個操作會觸發實際運算:reduce 以結合律函式合併元素並把結果送回 driver,collect 把所有元素送回 driver,foreach 則單純為了副作用(例如更新共享變數)對每個元素套用函式。在其中之一被呼叫之前,什麼都不會具現化;在文字搜尋範例中,errs 與 ones 從未被具現化,reduce 被呼叫時每個 worker 以串流方式掃描輸入區塊、就地算出中間元素、先做本地 reduce,最後只把本地計數送回 driver。此階段的 Spark 不支援 MapReduce 那種分組 reduce,所有 reduce 結果都落在單一的 driver 行程,不過各節點仍會先做本地縮減。作者以先前有研究在沒有平行縮減的情況下實作了十種機器學習演算法為由來辯護,並規劃日後加入 shuffle 轉換。

RDD 介面與 lineage 鏈

在內部,每種 RDD 都實作同一套三個操作的介面:getPartitions 回傳分割區 ID 清單,getIterator(partition) 走訪單一分割區,getPreferredLocations(partition) 則回報該分割區最好在哪裡計算。這些資料集物件串成一條記錄 lineage 的鏈,因此日誌計數範例會產生 HdfsTextFile、帶著 contains 述詞的 FilteredDataset、CachedDataset,以及帶著映射函式的 MappedDataset,每個都指向自己的父物件。不同型別的 RDD 差別只在如何實作那三個方法:HdfsTextFile 的分割區是 HDFS 的區塊 ID,偏好位置就是區塊所在位置,getIterator 會開啟區塊串流;MappedDataset 沿用父物件的分割區與偏好位置,只是在 iterator 中套上映射函式。CachedDataset 的 getIterator 會先找本地是否有已轉換分割區的快取副本,其偏好位置一開始等同父物件,但只要某個分割區在某節點被快取,就會更新成偏好該節點。容錯是這個設計的自然結果:節點掛掉時,它負責的分割區就從父資料集重新讀取並重算,最後被快取到別的節點上。

排程與運送閉包

Spark 跑在 Mesos 這個細粒度的叢集資源管理器上,因此能與 Hadoop、MPI 的 Mesos 版本共用同一座叢集與資料,也大幅降低了實作工作量。平行操作被呼叫時,Spark 為每個分割區建立一個 task,並用 delay scheduling 盡量把 task 送到該分割區的偏好位置;task 在 worker 上啟動後就呼叫 getIterator 開始讀取。運送 task 就等於運送閉包——包括定義資料集用的閉包,以及傳給 reduce 等操作的閉包——Spark 仰賴 Scala 閉包本質上是可被 Java serialization 搬移的 Java 物件這一點。但 Scala 內建的閉包實作並不理想:閉包物件有時會參照到函式主體其實沒用到的外層變數,把它們一起拖過網路,因此 Spark 對閉包類別的 bytecode 做靜態分析,找出這些未使用的變數,並在序列化前把對應欄位設為 null。

廣播變數的實作

兩種共享變數都以具備自訂序列化格式的類別實作,這正是它們與一般被捕捉變數傳輸方式不同的原因。當廣播變數 b 以值 v 建立時,v 會被寫入共享檔案系統中的一個檔案,而 b 的序列化形式只不過是指向該檔案的路徑。當 worker 第一次查詢 b 的值時,Spark 先檢查本地快取,未命中才從檔案系統讀取,因此不論有多少閉包或操作參照它,每個 worker 只付一次傳輸成本。最初的實作用 HDFS 來做這件事;論文指出以 HDFS 或 NFS 做的樸素廣播會讓廣播時間隨節點數線性成長,因此作者另外做了一套應用層 multicast 系統,並正在開發效率更好的串流式廣播。

accumulator 的實作

每個 accumulator 在建立時都會取得唯一 ID,其序列化形式只帶著這個 ID 加上該型別的 zero 值,因此 worker 從來不會收到累積中的狀態。在 worker 上,系統用執行緒區域變數為每個執行任務的執行緒各建立一份 accumulator 副本,並在 task 開始時歸零,同機器上的並行 task 因此不會互相干擾。task 結束後,worker 送一則訊息給 driver,內含該 task 對各個 accumulator 所做的更新。driver 對每個操作的每個分割區只套用一次更新,這條不變式正是在 task 因故障被重新執行、或分割區依 lineage 重算時,避免重複計數的關鍵。

直譯器整合

原版 Scala 直譯器會為使用者輸入的每一行編譯一個類別,其中的 singleton 物件保存該行的變數與函式,並在建構子中執行該行程式碼;之後對 x 的參照會編譯成透過 Line1.getInstance().x 的呼叫。Spark 改了兩件事。第一,讓直譯器把它定義的類別輸出到共享檔案系統,worker 節點再透過自訂的 Java class loader 載入,如此 worker 就能執行叢集啟動時根本不存在的程式碼。第二,修改產生的程式碼,讓每一行的 singleton 物件直接參照前面各行的 singleton 物件,而不是走靜態的 getInstance 方法,使閉包在被序列化時能捕捉到它所參照 singleton 的當下狀態;若不這麼做,像在提示字元下把 x 設為 7 這樣的更新就永遠傳不到 worker。

論文證明了什麼 — 量測數據與證明

  • 在 20 台四核心 m1.xlarge EC2 節點上對 29 GB 資料跑 logistic regression:Hadoop 每次迭代需 127 秒,因為每次迭代都是獨立的 MapReduce job;Spark 第一次迭代 174 秒(作者推測是用 Scala 而非 Java 所致),之後每次僅 6 秒,整體最多快 10 倍。
  • 在 10 次迭代的 logistic regression 實驗中途讓一個節點崩潰,平均使整體時間增加 50 秒(21%);遺失節點上的分割區會在其他節點平行重算並快取,但因為 HDFS 區塊大小設為 128 MB,每個節點只有 12 個區塊,復原過程無法用滿叢集所有核心,所以恢復時間偏高。
  • 在 30 節點 EC2 叢集上以 5000 部電影、15000 名使用者跑 alternating least squares:用廣播變數把評分矩陣 R 快取在 worker 記憶體中,相較於每次迭代重送 R,效能提升 2.8 倍;若不這麼做,重送 R 的時間會主導整個 job 的執行時間。
  • 使用以 HDFS 或 NFS 實作的樸素廣播時,廣播時間隨節點數線性成長,限制了 ALS job 的擴展性,這也是作者另外實作應用層 multicast 系統的原因。
  • 用改造過的 Scala 直譯器把 39 GB 的 Wikipedia dump 載入 15 台 m1.xlarge EC2 機器的記憶體中:第一次查詢約需 35 秒,與跑一個 Hadoop job 相當,但後續查詢即使掃描全部資料也只要 0.5 到 1 秒。
  • 在可程式性而非速度方面,論文指出 Spark 版的 logistic regression 程式與同演算法的序列版本只差三行,靠的是 accumulator 加上 Scala 的 for 語法自動展開為平行的 foreach。

限制與取捨 — 論文自承的,以及後來被發現的

  • 論文自承:Spark 不支援 MapReduce 那種分組 reduce,所有 reduce 結果都收斂到單一 driver 行程;在討論一節列為未來工作的 shuffle 轉換做出來之前,group-by 與 join 都無法實作。
  • 論文自承:cache 動作只是提示,工作集大於叢集記憶體時會悄悄退化成重算;而且當時只提供記憶體快取與寫入分散式檔案系統兩種持久化層級,記憶體內複本以及在儲存成本與重建成本之間取捨的通用旋鈕都仍是未來工作。
  • 論文自承:故障復原既不免費也不即時,在實測的故障實驗中吃掉 21% 的執行時間,而且受限於 128 MB 的粗粒度區塊;此外 Spark 的第一次迭代其實比 Hadoop 的每次迭代還慢(174 秒對 127 秒),所以只有在資料被重複使用時才划算。
  • 論文自承:RDD 明確不是通用的共享記憶體抽象,它是唯讀的、不支援細粒度寫入;作者也把「形式化刻畫 RDD 的性質及其對各類工作負載的適用性」列為待解項目,並全文將此實作描述為早期階段的原型。
  • 後續工作揭露:由於沒有 checkpoint 機制,長時間迭代 job 的 lineage 鏈會無限增長,復原成本也隨之上升;而此文對相依關係的處理過於粗略,排程器沒有可推理的資訊。2012 年 NSDI 的 RDD 論文補上了 narrow 與 wide dependency 的區分、以 stage 為單位的排程,以及 checkpoint,正是為了填補這些缺口。

它後來變成什麼 — 繼承這個想法的系統

這篇四頁半的研討會論文是 Apache Spark 的種子,而 Spark 後來取代 Hadoop MapReduce,成為大規模批次處理的預設引擎。它的直系續作、2012 年 NSDI 的 resilient distributed datasets 論文完整保留了 lineage 的想法,並補上 narrow 與 wide dependency 的分類、以 stage 為單位的 DAG 排程,以及針對長 lineage 鏈的 checkpoint。討論一節承諾的 shuffle 轉換後來確實落地,使 group-by 與 join 成為可能,也讓 Shark 乃至 Spark SQL、DataFrame 與 Catalyst 查詢最佳化器得以出現;這裡草擬的互動式直譯器則變成 spark-shell,以及 Databricks 商業化的 notebook 工作流程。其他抽象同樣延續下來:廣播變數與 accumulator 在現代 Spark 中幾乎原封不動,而「把資料留在記憶體重複使用」的論證也被帶進做迭代學習的 MLlib、做圖演算法的 GraphX,以及 Spark Streaming——後者的 discretized stream 本質上就是切成時間視窗的 RDD,並以 lineage 復原取代 upstream backup。更廣義地說,「一個帶著足夠資訊能自我重建的資料集 handle,是比複本或 checkpoint 更好的容錯原語」這個核心主張,也出現在 Flink、Dask 等後續資料流系統,以及 Delta Lake 這條 lakehouse 脈絡中——從持久化輸入做具決定性的重算,至今仍是它們的復原機制。

論文原文 — 逐字引用

“RDDs achieve fault tolerance through a notion of lineage: if a partition of an RDD is lost, the RDD has enough information about how it was derived from other RDDs to be able to rebuild just that partition.”

§1

“We note that our cache action is only a hint: if there is not enough memory in the cluster to cache all partitions of a dataset, Spark will recompute them when they are used.”

§2.1

“Although RDDs are not a general shared memory abstraction, they represent a sweet-spot between expressivity on the one hand and scalability and reliability on the other hand, and we have found them well-suited for a variety of applications.”

§1

術語 — 依本篇論文的用法

Resilient distributed dataset(RDD)
一個唯讀、切分到多台機器上的物件集合,任一分割區遺失都能被重建。它的元素不必存在於實體儲存上;handle 本身就帶著足以從可靠儲存中的資料算出整個資料集的資訊。
Lineage
由資料集物件串成的鏈,記錄每個 RDD 指向父物件的指標,以及父物件是如何被轉換的。Spark 重播這條鏈只重算遺失的分割區,而不是靠 checkpoint 與回捲。
Working set(工作集)
應用會跨多個平行操作重複使用的那份資料,例如迭代式機器學習或反覆執行的互動查詢。這正是非循環資料流系統處理不好、而 Spark 專門為之設計的工作負載類型。
平行操作(parallel operation)
把閉包送到 worker 上、藉此觸發 RDD 運算的動作:reduce 以結合律函式合併元素並回傳給 driver,collect 把所有元素送回 driver,foreach 則為副作用而對每個元素執行函式。
Driver 程式
使用者的主程式,負責實作應用的高階控制流程、定義 RDD 與共享變數,並在叢集上啟動平行操作。所有 reduce 與 collect 的結果都回到這裡。
廣播變數(broadcast variable)
包住一份大型唯讀值的封裝,保證該值只會被複製到每個 worker 一次,而不是隨每個閉包一起打包。它的序列化形式只是共享檔案系統中的一個檔案路徑,而且可跨多個平行操作重複使用。
Accumulator(累加器)
worker 只能以結合律運算加入、且只有 driver 能讀取的共享變數,只要型別具備 add 運算與 zero 值即可定義。只加不減的語意讓它容易做到容錯,driver 對每個分割區的更新只套用一次。
cache 動作
一種持久化設定的變更,讓 RDD 維持延遲求值,但提示系統在第一次算出後應保留在記憶體中。它僅僅是提示:叢集記憶體不足時,Spark 會在用到那些分割區時重新計算。
偏好位置(preferred locations)
由 getPreferredLocations 回傳、以分割區為單位的放置提示,供 delay scheduling 把 task 送到資料所在之處。對已快取的資料集,這些位置一開始沿用父物件,某分割區在某節點被快取後即更新為該節點。

在時間軸上的位置 — 這篇論文在整段故事中的座標

在時間軸上查看