跳至主要內容
論文精煉 · Distributed storage

Cassandra - A Decentralized Structured Storage System(Cassandra:一套去中心化的結構化儲存系統)

把 Dynamo 的無主環狀分散架構接上 Bigtable 的 column family 資料模型,扛下每天數十億次寫入的生產級儲存系統。

作者Avinash Lakshman、Prashant Malik,任職於 Facebook 發表於LADIS 2009(ACM SIGOPS 大規模分散式系統與中介軟體研討會);後刊於 ACM SIGOPS Operating Systems Review 44(2), 2010 年份2008–2010
閱讀原始論文 PDF 所有論文

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

Facebook 的 Inbox Search 需要一套能吃下每天數十億次寫入、在機器持續故障時仍不中斷、而且橫跨美國東西岸資料中心的儲存系統。Cassandra 的作法是把兩條血脈接起來:分散層走 Dynamo 路線,用 consistent hashing 環分割資料、複製到 N 個節點,讀寫都以 quorum 為準;資料層走 Bigtable 路線,是由 column family 與 super column 組成的多維映射。每次寫入先循序附加到 commit log,再更新記憶體中的資料結構;記憶體結構滿了就整批循序寫成一個不可變更、自帶索引的磁碟檔案,之後交由背景的壓縮合併程序做 merge sort。磁碟上的檔案永不修改,因此伺服器在讀寫路徑上幾乎完全免鎖。叢集成員與控制狀態靠 Scuttlebutt anti-entropy gossip 擴散,節點是否存活則交給 Phi accrual failure detector 判斷,它輸出的是連續的懷疑程度而不是非生即死的布林值。在生產環境中,Inbox Search 以 150 個節點存放超過 50TB 資料,讀取延遲中位數分別是 15.69 毫秒與 18.27 毫秒。

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

2008 年前後,Facebook 用數萬台伺服器服務數億使用者,而 Inbox Search 必須建立在原本躺在 MySQL 裡的 7TB 訊息資料之上。複寫式關聯式資料庫能給強一致性,但正如 Gray 與 Helland 早就指出的,代價是可擴充性與可用性,而且一旦發生網路分割就沒辦法繼續服務。Amazon 的 Dynamo 已經證明,一個會 gossip 的 consistent hashing 環加上交給客戶端處理衝突,可以在故障中維持可用;但它用 vector clock 偵測衝突,等於每次寫入都得先做一次讀取,在寫入遠多於讀取的場景裡負擔太重。Google 的 Bigtable 則示範了如何給應用一套稀疏、有序、以 column family 組織的資料模型,但它的持久性靠 GFS、協調靠有 Chubby 撐著的 master。Cassandra 正是長在這個縫隙裡:要 Dynamo 的可用性,要 Bigtable 的資料模型,底下不掛任何分散式檔案系統。

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

  • Inbox Search 要求系統吸收極高的寫入吞吐量,每天數十億次寫入,而且必須隨使用者數成長而擴充:上線時約一億使用者,論文寫作時已超過兩億五千萬。
  • 在由數千個元件組成的基礎設施裡,任何時刻都有數量雖小但不可忽視的伺服器與網路元件正在故障,因此儲存系統必須把故障視為常態而非例外,而且不能有任何單點故障。
  • 使用者由地理上分散的資料中心服務,所以每一列資料都必須跨資料中心複寫,一方面壓低搜尋延遲,一方面在停電、散熱失效、網路中斷或天災導致整座資料中心失效時仍能存活。
  • 傳統的複寫式關聯式資料庫把重心放在保證強一致性,作者指出這讓它們在可擴充性與可用性上受限,而且無法在網路分割發生時繼續運作。
  • Dynamo 用 vector clock 偵測更新衝突,寫入時必須順帶做一次讀取來維護版本時間戳,作者認為這在必須承受極高寫入吞吐量的環境中限制太大。
  • 現成的 gossip 式失效偵測器會隨叢集規模劣化:在一次 100 個節點的實驗中,偵測到一個失效節點所需的時間長達兩分鐘左右,作者直言這在他們的環境中根本不可用。

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

Dynamo 的分散架構配上 Bigtable 的資料模型

Cassandra 明白地說自己是既有技術的綜合,而不是新演算法:分散層完全是 Dynamo 那一套,consistent hashing 環、複製因子 N、preference list、quorum 讀寫、gossip 成員管理;紀錄層則是 Bigtable 那一套,一個由 column family 組成的分散式多維映射,再用 Super column family 多加一層巢狀。這個嫁接之所以成立,是因為兩層解決的是互不相干的問題:環決定哪些機器持有某個 key、以及它們掛掉時系統怎麼反應;column family 則決定應用在同一個 key 底下如何排列與排序資料。應用得到的保證是每個 row key 在每個複本上的操作都是原子的,不論一次讀寫涉及多少欄,而且可以指定欄位依名稱或依時間排序,Inbox Search 直接利用後者讓結果天生就是時間序。與 Bigtable 不同的是,底下沒有分散式檔案系統,持久性就靠每個節點自己的本機檔案系統。

保序的環,用負載來重新平衡

雜湊函數的輸出值域被視為一個固定的環狀空間,每個節點取一個隨機 token 當作自己在環上的位置;一個 key 落在自己位置順時針方向遇到的第一個節點上,那個節點就是它的 coordinator,並負責自己與前驅節點之間的那一段區間。Cassandra 刻意採用保序的雜湊函數,讓 key 的順序在分割之後仍然成立,同時保留 consistent hashing 最重要的性質:節點加入或離開只影響它的直接鄰居。隨機 token 會造成資料與負載分布不均,也完全不理會機器效能的異質性;在這一點上 Cassandra 與 Dynamo 分道揚鑣:它不是讓一個節點在環上佔多個虛擬位置,而是分析環上的負載資訊,把負載輕的節點在環上搬動,去分擔負載重的節點,也就是 Chord 論文描述的路線。作者選這條路的理由很實際:設計與實作都好處理,而且負載平衡的決策非常確定。

全部只做循序寫的寫入路徑

一次寫入就是往專屬磁碟上的 commit log 做一次循序附加,只有在這次附加成功之後才會更新記憶體中的資料結構;當該結構超過依資料量與物件數計算出的門檻,就整批循序寫到某顆普通磁碟上,同時產生以 row key 為基礎的索引,索引與資料檔一起持久化。這些檔案寫出之後永不修改,所以真正負責整併版本、回收空間的是背景的壓縮合併程序,也就是從 Bigtable 借來的 compaction。這條路徑之所以快,是因為寫入路徑上每一次磁碟操作都是循序的,正好是廉價磁碟最擅長的事;而不可變更也意味著讀取者永遠不必和寫入者搶鎖。作者直接寫道,Cassandra 的伺服器實例在讀寫操作上幾乎完全免鎖,這正是它避開 B-tree 式資料庫實作中那些並行控制問題的原因。

不需要先讀再寫的 quorum

任何一個節點都能接下對某個 key 的讀寫請求;它算出該 key 的複本集合,把寫入送往所有複本並等待 quorum 數量的確認,而讀取則依客戶端要求的一致性,或送往最近的複本,或送往全部複本再等 quorum 個回應。調和方式是看時間戳:路由狀態機挑出最新的回應,並對落後的複本排程一次修復,因此寫入之前不需要先讀一次來產生版本標記,而這正是作者想避開的 Dynamo 成本。寫入可以設定成同步或非同步,寫入遠多於讀取的系統就採非同步複寫。面對節點失效與網路分割時的持久性保證,來自放寬 quorum 要求,而不是死等所有複本到齊。

架在 gossip 之上的累積式失效偵測

成員資訊與其他控制狀態透過 Scuttlebutt 這個 anti-entropy gossip 協定擴散,選它是因為它對 CPU 與 gossip 通道的使用都極有效率。在它之上,每個節點為每個對象維護一個 gossip 訊息到達間隔的滑動視窗,估出分布之後輸出一個連續的懷疑程度 Phi,而不是上線或下線的布林值:在 Phi = 1 就判定可疑,之後被遲到的心跳打臉的機率約 10%,Phi = 2 約 1%,Phi = 3 約 0.1%。這樣做的價值在於,門檻變成一個明確的準確度對速度的旋鈕,而且會自動隨網路與伺服器負載調整,不再是那種叢集一長大就得重調的固定逾時。作者對原版偵測器做的修改是把到達間隔近似成 Exponential 分布而非 Gaussian 分布,因為這更貼近 gossip 通道的行為,並認為這是首次把累積式失效偵測器用在 gossip 環境中。

感知機架與資料中心的複寫,加上一個 leader

每筆資料複製到 N 台主機,N 是以 instance 為單位設定的;coordinator 除了把落在自己區間的 key 存在本機,還會複製到另外 N-1 個節點,選誰由應用挑選的複寫策略決定。Rack Unaware 直接取環上的 N-1 個後繼節點,而 Rack Aware 與 Datacenter Aware 則刻意把複本擺到不同機架、不同資料中心,這也表示複本位置不再能單靠環本身推導出來。因此 Cassandra 透過 Zookeeper 選出一個 leader;節點加入叢集時向 leader 詢問自己是哪些區間的複本,leader 則盡力維持一個不變量:沒有任何節點負責超過 N-1 個區間。區間的中繼資料同時快取在各節點本機、也以容錯方式存在 Zookeeper 裡,讓當掉又復原的節點知道自己原本負責什麼;而把一個 key 的 preference list 分散到由高速網路連接的多個資料中心,正是 Facebook 能夠整座資料中心失效卻不中斷服務的原因。

運作方式 — 具體的機制

請求路由狀態機

叢集中任何節點都能當請求的入口。它的路由狀態機會走過五個狀態:找出擁有該 key 資料的節點;把請求送出去並等待回應;若在設定的逾時內沒等到回覆就讓請求失敗、回報客戶端;依時間戳判斷哪個回應最新;最後對任何沒有最新資料的複本排程一次修復。所有系統控制訊息走 UDP,而複寫與請求路由這類應用訊息走 TCP,底層是以 non-blocking I/O 建構的網路層。訊息處理管線與工作管線依 SEDA 架構切成多個階段,分割模組、成員與失效偵測模組、儲存引擎模組全部以 Java 從頭實作。

節點加入環的 bootstrap

節點第一次啟動時會為自己在環上的位置挑一個隨機 token,並把這個對應同時持久化到本機磁碟與 Zookeeper,接著透過 gossip 傳出去;由於每個節點最後都知道所有節點的 token,任何節點都能把一個 key 的請求直接導到正確的節點。要加入既有叢集的節點會讀設定檔中列出的 seeds,也就是幾個初始聯絡點,這些聯絡點也可以由 Zookeeper 這類設定服務提供。每一則訊息都帶著該 Cassandra instance 的叢集名稱,所以設定錯誤而想加入別的叢集的節點會被擋下來。由於 Facebook 環境中的節點停機通常只是暫時性的、極少代表永久離開,節點的加入與移除被設計成明確的管理操作,由管理者用命令列工具或瀏覽器下達,這樣暫時性故障就不會觸發分割區重新配置或對聯絡不到的複本做修復。

擴充叢集與資料串流

新節點會被指派一個能替負載重的節點分憂的 token,結果是把原本某個節點負責的區間切開。bootstrap 由操作人員從系統中任何一個節點發起,用命令列工具或 Cassandra 的網頁儀表板都可以。讓出資料的節點以 kernel 對 kernel 的複製技巧把資料串流給新節點,營運經驗顯示單一來源節點的傳輸速率約為 40 MB/sec。作者也提到正在改成讓多個複本一起參與 bootstrap 傳輸,把工作平行化,概念上類似 Bittorrent。

失效偵測的運作迴圈

每個節點為每個對象維護一個 gossip 訊息到達間隔的滑動視窗,判定其分布後算出 Phi。呼叫者拿 Phi 和門檻比較來決定是否把對方視為已下線,而這個判斷不只用於成員管理,也用來在各種操作中避免嘗試與聯絡不到的節點通訊。由於 Phi 是懷疑程度而非布林值,不同子系統可以在同一個底層訊號上選擇不同的信心水準。Facebook 採用略為保守的 Phi = 5,在 100 個節點的叢集中平均約 15 秒就能偵測到失效。

commit log、記憶體資料表與落地

寫入會附加到機器上專屬磁碟的 commit log,因為所有 commit log 寫入都是循序的,可以把磁碟吞吐量壓到最大;記憶體中的資料結構只在日誌寫入成功後才更新。commit log 超過可設定的大小就會換新,生產環境中設為 128MB 效果很好。每個 commit log 都有一個固定長度的位元向量標頭,長度比系統可能擁有的 column family 數量還多;每當某個 column family 的記憶體結構被寫到磁碟,就把它對應的位元設起來,而每次換新 commit log 時,會檢查它與所有更早日誌的位元向量,確認資料都已落地就把這些日誌刪掉。另外還有 fast sync 模式,會把 commit log 寫入與記憶體結構的落地都改成緩衝式,作者明說這代表機器當掉時有遺失資料的可能。

磁碟格式與讀取路徑

所有資料以主鍵建索引,每個資料檔切成一連串區塊,每個區塊最多 128 個 key,並由一個區塊索引標定,記錄每個 key 在區塊內的相對位移與其資料大小;這份索引在落地時產生、寫到磁碟,同時也留在記憶體中以便快速存取。讀取一律先查記憶體中的資料結構,因為那裡永遠握有某個 key 的最新資料,找不到才對磁碟上的資料檔做 I/O,並且由新到舊逐一查看,一命中就回傳。每個資料檔都附帶一個彙整其 key 的 bloom filter,同樣保存在記憶體中,先問過它,不可能含有該 key 的檔案就完全不會被打開。由於一個 key 在 column family 底下可能有非常多欄,欄位序列化寫出時會每 256K 區塊邊界產生一次欄位索引,讓讀取能直接跳到正確的區塊,而不必掃過磁碟上每一欄;這個邊界可設定,但 256K 在生產工作負載上表現很好。

壓縮合併

由於落地的檔案會隨時間累積,背景程序會把多個檔案合成一個,本質上就是對一堆已排序的資料檔做 merge sort,和 Bigtable 的 compaction 非常類似。Cassandra 只合併大小相近的檔案:論文明言絕不會出現 100GB 的檔案和小於 50GB 的檔案被合併在一起的情況,這也就限制了單次合併的成本。系統會週期性執行一次 major compaction,把所有相關資料檔壓成一個大檔。作者承認壓縮合併是磁碟 I/O 密集的操作,也提到可以做很多最佳化,避免影響進來的讀取請求。

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

  • Facebook Inbox Search 的部署在一個 150 節點的叢集上存放約 50TB 以上的資料,節點分散於美國東西岸的資料中心。
  • Inbox Search 在生產環境量到的讀取延遲:互動搜尋最小 7.69 毫秒、中位數 15.69 毫秒、最大 26.13 毫秒;詞彙搜尋最小 7.78 毫秒、中位數 18.27 毫秒、最大 44.41 毫秒。
  • 把累積式失效偵測器的 Phi 門檻設在略為保守的 5,在 100 節點叢集中偵測失效的平均時間約 15 秒;相較之下,同一個實驗中團隊試過的舊式 gossip 失效偵測器要花大約兩分鐘。
  • 節點之間的 bootstrap 資料傳輸使用 kernel 對 kernel 的複製技巧,實測單一來源節點的速率為 40 MB/sec。
  • 首次上線前的匯入用 Map/Reduce 作業,把 Facebook MySQL 基礎設施中一億以上使用者、共 7TB 的收件匣資料建成反向索引,再經由背景通道把序列化後的資料送進 Cassandra,使 Cassandra instance 只受限於網路頻寬。
  • 論文報告的生產環境參數:commit log 在 128MB 換新、欄位索引每 256K 區塊邊界產生一次、資料檔區塊最多 128 個 key;Inbox Search 於 2008 年 6 月為約一億使用者上線,論文寫作時已服務超過兩億五千萬使用者。

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

  • 論文自承:原子性僅限於單一 key、單一複本,沒有跨 key 的交易。有些應用要求交易支援,主要是為了維護次要索引,而作者只說他們正在研究如何提供這類原子操作。
  • 論文自承:壓縮、跨 key 的原子性與次要索引支援全被列在結論的未來工作裡,因此想要倒排索引的應用得自己建、自己維護,就像 Inbox Search 用 super column 做的那樣。
  • 論文自承:commit log 的 fast sync 模式會把日誌寫入與記憶體結構落地都改成緩衝式,作者明說這代表機器當掉時可能遺失資料,也就是持久性是可調的參數而非保證。
  • 論文自承:雖然 Cassandra 被描述成完全去中心化的系統,作者也學到某種程度的協調是必要的,因此用 Zookeeper 做 leader 選舉與區間指派,而成員變更仍是明確的管理指令,而非系統自動反應。
  • 後續才暴露的問題:保序的分割器加上由人工搬動 token,實際上造成熱點與痛苦的重新平衡,Cassandra 最終改以隨機分割器為預設,並採用 Dynamo 式的 virtual node。純以時間戳調和複本也意味著後寫者勝出,在時鐘偏移下並行更新可能悄悄遺失,這個取捨論文完全沒有提及;此外評估本身相當單薄,只有單一應用的讀取延遲,沒有任何吞吐量或擴充性曲線。

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

Facebook 在 2008 年把 Cassandra 開源,2009 年進入 Apache 孵化器,2010 年成為 Apache 頂層專案,一個 Facebook 內部儲存系統就這樣變成 Dynamo 加 Bigtable 這套設計的參考實作。這篇論文定下的樣板:token ring、複製因子 N 搭配可依請求調整的一致性、gossip 成員管理、read repair(路由狀態機的第五步),以及由 commit log、記憶體資料表、不可變更資料檔與壓縮合併組成的 LSM 式本機引擎,基本上就是 Apache Cassandra 至今的架構,也是 Designing Data-Intensive Applications 講 leaderless 複寫時的標準案例。後來的版本換掉了論文中那些權宜的部分:保序分割器讓位給隨機分割器,virtual node 在 Cassandra 1.2 加入,Zookeeper 被以 gossip 為基礎的協商取代,而 insert、get、delete 這三個方法的 Thrift API 也被 CQL 取代。它的後代與親戚包括 ScyllaDB,一個以 shard-per-core 思路用 C++ 重寫、且與 Cassandra 協定相容的實作;商業化的 Amazon Keyspaces 與 DataStax Enterprise;以及同樣承襲 Dynamo 血脈、走純 key-value 路線的 Riak。Phi accrual failure detector 在這篇論文之前相當冷門,之後成為叢集成員管理的標準零件,Akka Cluster 等系統都採用它。至於 Facebook 自己,後來把 Messages 與 Inbox Search 移到 HBase 上,提醒我們一件事:當初的動機應用比它的第一代儲存引擎先退場,而那個引擎反而活得更久。

論文原文 — 逐字引用

“Cassandra system was designed to run on cheap commodity hardware and handle high write throughput while not sacrificing read efficiency.”

摘要

“Cassandra morphs all writes to disk into sequential writes thus maximizing disk write throughput.”

§5.7 實作細節

“With the accrual failure detector with a slightly conservative value of PHI, set to 5, the average time to detect failures in the above experiment was about 15 seconds.”

§6 實務經驗

術語 — 依本篇論文的用法

Column family
同一個 row key 之下一組具名的欄位集合,也是記憶體與磁碟上的組織單位,因為 Cassandra 為每個 column family 各維護一份記憶體資料結構與一個資料檔。欄位以 column family : column 的慣例存取。
Super column family
巢狀在 column family 裡的 column family,以 column family : super column : column 存取。Inbox Search 就是把訊息中的詞或收件者 id 當作 super column,把個別訊息識別碼當作其中的欄位。
Coordinator
把 key 雜湊到環上,再順時針走到第一個位置較大的節點,那個節點就是 coordinator。它負責自己與前驅節點之間的環上區間,並負責把落在該區間的 key 複製到其他複本。
Preference list
借自 Dynamo 的用語,指負責某個區間的節點集合。Cassandra 會刻意讓一個 key 的 preference list 中的儲存節點跨越多個資料中心,使整座資料中心失效時服務仍不中斷。
Phi accrual failure detector
不輸出上線或下線的布林值,而是輸出連續變化的懷疑程度 Phi 的失效偵測器,其值由 gossip 訊息到達間隔的滑動視窗算出。把門檻調高就是以偵測速度換取更低的誤判機率。
Scuttlebutt
Cassandra 用來管理叢集成員、並散播其他系統控制狀態的 anti-entropy gossip 機制,選用它是因為它對 CPU 與 gossip 通道的使用都非常有效率。
Commit log
每個節點上獨佔一顆磁碟的循序持久性日誌,任何寫入都必須先寫進它,才能更新記憶體中的資料結構。日誌在可設定的 128MB 大小換新,並在其標頭位元向量顯示所含的每個 column family 都已落地後被刪除。
Compaction(壓縮合併)
把磁碟上多個不可變更的資料檔整併成較少檔案的背景程序,本質是對已排序檔案做 merge sort,且只合併大小相近的檔案;此外會週期性執行 major compaction,把所有相關檔案壓成一個。
Bloom filter
彙整某個資料檔中所有 key 的精簡摘要,與資料檔一起存放並常駐記憶體;磁碟查找前先問它,不可能含有目標 key 的檔案就完全不會被讀取。

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

在時間軸上查看