跳至主要內容
論文精煉 · Parallel DBMS

Parallel Database Systems: The Future of High Performance Database Processing

以 shared-nothing 硬體、資料分割與 split/merge 資料流執行環境,讓關聯式查詢取得近乎線性的 speedup 與 scaleup。

作者David J. DeWitt(University of Wisconsin-Madison 資訊科學系)與 Jim Gray(Digital Equipment Corporation,San Francisco Systems Center) 發表於CACM 35(6),1992 年 6 月 年份1986–1990s
閱讀原始論文 PDF 所有論文

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

到了 1992 年,專用硬體的資料庫機器已經宣告失敗,但完全由通用處理器、記憶體與磁碟組成的平行資料庫系統,卻正在最大型的工作負載上取代大型主機。DeWitt 與 Gray 解釋了原因:關聯式運算子吃進、吐出的都是格式一致的 tuple 串流,因此一句 SQL 可以編譯成一張資料流圖,既能做 pipeline 平行,更能在分割(declustered)後的關聯上做效益高得多的分割平行。關鍵技巧在於任何運算子都不必重寫:既有的循序 scan、sort、join 程式碼只需外包兩個管線運算子,split(把每筆輸出 tuple 路由到某個目的行程)與 merge(把多條串流匯成單一循序輸入埠),並內建流量控制。建議的硬體是 shared-nothing,每個磁碟與記憶體由單一處理器擁有,網路上流動的只有問題與答案而不是資料頁;相對地,shared-memory 與 shared-disk 都要付出干擾(interference)代價,規模因此卡在數十顆處理器。論文以 Teradata、Tandem、Gamma、Bubba 以及跑在 nCUBE 上的 Oracle 為證,說明這個設計可以近乎線性地擴充到數百顆處理器,而且便宜到讓 Grosch 定律在資料庫領域徹底失效。

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

1983 年 Boral 與 DeWitt 曾寫過一篇著名的批判,認為資料庫機器是一個時代已經過去的點子,而當時證據站在他們那一邊:十年的研究追逐過 CCD 記憶體、磁泡記憶體、head-per-track 磁碟與光碟,沒有一項兌現承諾。預測指出處理器速度的成長會遠快於磁碟吞吐量,批評者因此認為多處理器很快就會被 I/O 卡死。同一時間,大型主機的設計者做不出單一機器來服務上千個同時使用者,或掃描 TB 級的關聯式資料庫。Encore、Intel、NCR、nCUBE、Sequent、Tandem、Teradata 與 Thinking Machines 陸續推出便宜的微處理器機器,而建立在高速區域網路上的訊息式 client-server 作業系統也從研究玩具變成主流架構。Teradata 從 1978 年起就低調地出貨高度平行的 SQL 機器,Tandem 與一批新創隨後跟上,所以 1992 年要問的已經不是平行資料庫行不行,而是它為什麼行。

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

  • 大型主機設計者無法造出單一機器,同時滿足關聯式資料庫服務大量同時使用者、以及搜尋 TB 級資料庫所需的 CPU 與 I/O 需求。
  • 在多數電腦設計中,每加入一顆處理器都會讓其他處理器慢一點;論文指出即使干擾只有 1%,最大 speedup 也只有 37,天真的多處理器根本擴不上去。
  • 十年來的專用資料庫機器硬體,包括 CCD 記憶體、磁泡記憶體、head-per-track 磁碟與光碟,都沒有兌現承諾,連帶讓整個資料庫機器路線失去信譽。
  • 當時預測磁碟吞吐量只會成長一倍,處理器速度卻會成長得快得多,因此若不正面處理 I/O 瓶頸,多處理器系統終將受制於 I/O。
  • 為單處理器寫的舊軟體搬到任何多處理器上都得不到 speedup 或 scaleup,必須重寫;而對絕大多數應用領域而言,重寫的成本高到不可行。
  • 即使正確地平行化,仍有三個通用的線性障礙:啟動(startup)成本,也就是必須拉起上千個行程的時間;共用資源上的干擾(interference);以及偏斜(skew),當各步驟大小的變異數超過平均值時,整份工作的完成時間由最慢的那一步決定。

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

以 speedup 與 scaleup 當標尺

論文主張評價平行資料庫要看兩個比值,而不是硬體的規格峰值。線性 speedup 指的是硬體放大 N 倍,同一份工作就快 N 倍,計算方式是小系統耗時除以大系統耗時。線性 scaleup 指的是放大 N 倍的系統在同樣時間內完成放大 N 倍的工作,也就是比值等於 1;它又分兩種:交易 scaleup(N 倍的用戶端對 N 倍大的資料庫送出 N 倍多的小請求,正是 TPC 基準測試採用的放大方式)與批次 scaleup(同一句查詢跑在 N 倍大的資料庫上)。把成功定義成可被推翻的量測,才讓作者能明確點名破壞線性的三件事:startup、interference 與 skew。

shared-nothing 贏得架構之爭

論文沿用 Stonebraker 的分類,比較 shared-memory(所有處理器共用一塊全域記憶體與所有磁碟)、shared-disk(各自擁有私有記憶體,但每顆處理器都能直接定址所有磁碟)與 shared-nothing(每塊記憶體與磁碟由單一處理器擁有並擔任其伺服器,處理器之間只靠訊息溝通)。shared-nothing 用最小化共用來最小化干擾:原始的記憶體與磁碟存取都在本地完成,只有過濾後大幅縮小的結果會經過網路,網路上跑的是問題與答案而不是資料頁。shared-memory 則必須打造頻寬等於所有處理器與磁碟總和的互連網路,而且量測顯示,資料庫工作負載下這些大型私有快取的填入與清空會明顯拖慢處理器。shared-disk 的失敗方式不同:任何更新都得先宣告意圖、等所有其他處理器確認,再讀寫整個實體資料頁,成本遠高於交換小小的高階問答訊息。

關聯式查詢就是資料流圖

每個關聯式運算子都以關聯為輸入、以關聯為輸出,因此運算子可以任意組合;一句 SQL 其實只是 scan、sort、聚合、join、insert、update、delete 這組節點所構成的圖的語法糖。由於 SQL 是非程序式語言,決定查詢怎麼執行的是系統而不是程式設計師,這代表為單處理器寫的 SQL 應用程式可以原封不動地在 shared-nothing 機器上平行執行。這正是為什麼在其他領域擋住平行化的舊軟體障礙,唯獨在資料庫領域是例外。作者也點出歷史上的反諷:Codd 提出關聯式模型是為了程式設計生產力與資料獨立性,平行處理則是意料之外的紅利。

真正有效的是分割平行而非 pipeline

論文直言,資料流圖提供的兩種平行中,pipeline 是弱的那一種。關聯式 pipeline 通常很短,長度十的鏈已屬罕見;sort 與聚合是阻塞型運算子,吃完全部輸入之前吐不出任何結果,根本無法 pipeline;而且往往某個運算子的成本遠高於其他運算子,這本身就是一種 skew,會限制 pipeline 的加速上限。分割式執行採取的是分而治之:把運算子的輸入與輸出分割到多個處理器與磁碟上,將一件大工作變成許多互相獨立的小工作。此時 speedup 隨分割數成長,而不是隨查詢計畫的深度成長,這也是論文說「分割的資料是分割式執行的關鍵」的原因。

declustering 就是實體設計

把一個關聯的 tuple 散佈到多顆磁碟,是分割式執行的前提,而且不必任何特殊硬體就能取得優於 RAID 式條帶化的 I/O 頻寬。三種基本策略各有取捨:round-robin 在每筆查詢都全表掃描時最理想,但關聯式的關聯查找會被迫打到每一顆磁碟;hash 分割能把分割屬性上的等值查找導到單一磁碟,代價是資料被隨機打散而非群聚;range 分割保留了群聚性、也很適合範圍述詞,但有資料偏斜以致執行偏斜的風險,這點 hash 與 round-robin 反而比較耐受。Bubba 進一步依每筆 tuple 的存取頻率(heat)來做 range 分割,讓各分割在被存取的頻率(temperature)上取得平衡,而不是在 tuple 數量(volume)上取得平衡。分割也不是切越細越好:超過某個程度後,每個節點啟動查詢的成本會佔實際執行時間相當可觀的比例,再切下去反而拉長回應時間。

把平行封裝進 split 與 merge

系統不去改寫 scan、sort、join 的平行版本,而是保留既有的循序實作,只改動它們之間的管線。每個運算子有一組輸入埠與一個輸出埠;merge 運算子把多條平行串流併成一條循序串流餵進某個埠,split 運算子則依 tuple 的屬性值把每筆輸出送往多個目的行程之一。這個映射可以是範圍述詞、hash、round-robin、整串流複製,甚至任意程式;由於平行完全住在這兩個運算子裡,日後新增的任何關聯式運算子都自動具備平行能力。split 與 merge 還內建緩衝與流量控制:當 split 的輸出緩衝區塞滿時,它會擋住上游運算子直到下游索取更多資料,整張圖因此自我調速。

hash join 作為平行 join

傳統的 sort-merge join 先把兩邊依 join 屬性排序再合併,因此繼承了排序的 n log n 成本;一旦資料偏斜,某些排序分割會遠大於其他分割,資料偏斜就轉成執行偏斜,限制了 speedup 與 scaleup。hash join 改成先把兩個關聯依 join 屬性做 hash 分割,把 A 的某個分割建成記憶體中的 hash table,再掃描 B 的對應分割逐筆探測,命中就把兩筆串接後送進輸出串流。它的成本是線性而非 n log n,把一個大 join 拆成許多互相獨立的小 join,而且更耐偏斜,因此除非輸入本來就已排序,否則都優於 sort-merge join。只要 hash 函數夠好、偏斜不太嚴重,各 bucket 大小差異就很小,join 便能取得線性 speedup 與 scaleup;作者以此為證,說明真正值得投入的是更好的平行演算法而不是更特殊的硬體。

運作方式 — 具體的機制

先把每個關聯 decluster 到磁碟上

載入時,每個關聯都要指定一種分割策略與一組磁碟片段,通常一個參與的處理器一份。round-robin 把第 i 筆 tuple 放到第 i mod n 顆磁碟;hash 分割對選定屬性套用 hash 函數決定磁碟;range 分割把連續的屬性區間(例如姓名 a 到 c 放一顆、d 到 g 放下一顆)對應到不同磁碟;Gamma 另外提供混合了 hash 與 range 特性的 hybrid-range 分割。提高分割程度會縮短循序掃描時間,因為同時被讀取的磁碟變多;也會縮短關聯查找時間,因為每個節點上的 tuple 變少、要搜尋的索引也變小。實體設計因此變成逐關聯選擇策略、屬性與分割度的問題,而論文明白指出當時沒有任何自動化工具能幫忙做這個決定。

把 SQL 編譯成循序運算子構成的圖

查詢最佳化器產生一張邏輯查詢圖:以「將 A 與 B 依 A.x = B.y 做 join 後插入 C」這句查詢為例,樹上有兩個 scan 節點(各對應一個輸入關聯)、一個 join 節點與一個 insert 節點。每個節點都是普通的循序關聯式運算子,具備一組編號的輸入埠與單一輸出埠。平行化是另一個獨立步驟:在這棵樹的節點之間插入 split 與 merge 運算子,而不是去改寫任何運算子。Tandem、Gamma 與 Volcano 採用的正是這個做法,也因此同一個查詢計畫形狀可以套用在任意的平行度上。

merge:多條平行串流併成一個埠

假設關聯 A 被分割成 A0、A1、A2 三個片段。平行查詢執行器會建立三個 scan 行程,各自指向一個片段,並要求三者把輸出送到同一個 merge 節點。merge 運算子產生單一輸出串流,可以送給應用程式、顯示在終端機上,或餵給下一個關聯式運算子;下游完全不知道自己的輸入其實來自三台機器。因此 merge 就是讓循序的消費者能坐在分割式生產者之上、卻不必修改任何程式碼的那個運算子。

split:一條串流分送到多個目的地

split 運算子持有一張表,把輸出 tuple 屬性上的述詞映射到形如(cpu 編號、行程編號、埠編號)的目的地三元組。在論文的例子中,每個關聯 A 的 scan 上的 split 會把落在 A-H 的 tuple 送到 cpu 5 / 行程 3 / 埠 0,I-Q 送到 cpu 7 / 行程 8 / 埠 0,R-Z 送到 cpu 2 / 行程 2 / 埠 0;而每個關聯 B 的 scan 上的 split 用相同的區間,但改送到同樣那三個行程的埠 1。其他 split 也可以複製整條串流、依 round-robin 分割,或依 hash 分割,分割函數甚至可以是任意程式。緩衝與流量控制就住在 split 裡:輸出緩衝區塞滿時它會擋住上游的關聯式運算子,直到下游索取更多資料,避免圖的某一段跑得離其他部分太遠。

完整範例:一個分割式 join

以那句 insert into C select from A, B 的查詢為例,假設三個行程執行 join,三個行程掃描 A 的片段、兩個行程掃描 B 的片段。所有 A 的 scan 套用同一個 split,因此 join 行程 0 會在埠 0 上收到由三個 A scan 產生、經 merge 併成一條的全部 A-H tuple;所有 B 的 scan 套用對應的 split,同一個 join 行程便在埠 1 上收到 B 的全部 A-H tuple。每個 join 行程於是只看到兩條普通的循序輸入串流,可以跑 hash join、sort-merge join,甚至在 tuple 到達順序合適時跑巢狀迴圈 join,完全不必察覺周圍的平行結構。三個 join 的輸出再依關聯 C 自身的分割準則被 split,並在三個 insert 節點上 merge,結果因此直接落成正確的分割形式,不需要額外一輪重新分佈。

平行 hash join 的執行流程

A 與 B 都先依 join 屬性做 hash 分割,這保證可配對的 tuple 一定落進同一組 bucket,於是完全不需要跨節點比對。接著把 A 的某個 hash 分割建成主記憶體 hash table,掃描 B 的對應分割並逐筆探測,命中就把配對的兩筆串接後寫入輸出串流;所有分割配對依序重複此步驟。正確性只依賴分割不變式,效能則取決於 bucket 大小的變異數,因此好的 hash 函數加上不太嚴重的偏斜,就能得到接近均勻的 bucket 與線性的行為。失敗模式論文說得很明白:若大量甚至全部 tuple 的 join 屬性值相同,所有 tuple 會落進同一個 bucket,而在這種病態情況下,目前沒有已知演算法能做到 speedup 或 scaleup。

實際出貨系統怎麼落實這套設計

Teradata 把處理器分成負責剖析、最佳化與協調的 Interface Processor,以及負責儲存與執行的 Access Module Processor,兩者以雙重冗餘的樹狀 Y-net 互連;第一層 hash 依主鍵挑出 AMP,第二層 hash 決定 tuple 在該 AMP 片段內的位置且片段內依 hash-key 排序,因此依鍵查找只會打到一個 AMP,快取未命中時也只需一次磁碟讀取;join 以平行 sort-merge 執行,而且不採 pipeline,而是每個運算子在所有節點上跑完之後才啟動下一個。Tandem NonStop SQL 讓應用程式與資料庫伺服器跑在同一批處理器上,以四重光纖環互連,大致配置成每 MIPS 一顆磁碟並做磁碟鏡像,關聯採 range 分割,次要索引只支援 B-tree,join 提供巢狀迴圈、sort-merge 與 hash 三種,平行化同樣靠在查詢樹中插入 split 與 merge。Tandem 在 OLTP 上的關鍵手法是平行索引維護:一個關聯通常帶五個索引、有時多達十個,把索引攤到多顆處理器與磁碟上,就能讓索引增加時的總維護時間幾乎維持不變。Gamma 跑在 32 節點的 Intel iPSC/2 Hypercube 上、每節點一顆磁碟,提供 round-robin、range、hash 與 hybrid-range 分割,以及以 B-tree 或 hash table 實作的群聚與非群聚索引,並用 split 與 merge 同時取得分割與 pipeline 兩種平行。

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

  • 支撐 shared-nothing 的干擾算式:若每加入一顆處理器都讓其他處理器慢 1%,最大可達的 speedup 只有 37,而一千顆處理器的系統實際上只發揮出單處理器系統 4% 的有效算力。
  • 實際達到的規模:當時市面上最大的 shared-memory 多處理器約僅限 32 顆處理器,而 Teradata、Tandem 與 Intel 都已出貨超過 200 顆處理器的系統,Intel 正在實作 2000 節點的 Hypercube,Teradata 的組態則可以有超過一千顆處理器與數千顆磁碟。
  • Oracle 跑在 64 節點的 nCUBE shared-nothing 系統上,成為業界標準 TPC-B 基準測試上第一個突破每秒 1000 筆交易的系統,無論峰值效能或價格效能都遠超過 Oracle 自己在傳統大型主機上的表現。
  • Tandem NonStop SQL 在 TPC-A 上線性擴充到遠超過已知最大的大型主機,價格效能比同級主機便宜三倍;Gamma、Tandem 與 Teradata 也都在複雜關聯式查詢基準測試上量到近乎線性的 speedup 與 scaleup。
  • 推翻 Grosch 定律的價格證據:大型主機當時的定價是每 MIPS 25,000 美元、每 MB 記憶體 1,000 美元,微處理器則是每 MIPS 250 美元、每 MB 記憶體 100 美元,因此把數百甚至數千台小系統組起來,能用遠低於一台中型主機的成本買到更多資料庫算力。
  • 作者為何說「近乎線性」而不是線性:因為排序成本是 n log n,問題放大一千倍會讓 n log n 增加 3000 倍,也就是在三個數量級的 scaleup 上出現 30% 的線性偏差。

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

  • 論文自承:pipeline 平行的貢獻非常有限,因為關聯式 pipeline 很少超過十層、sort 與聚合會阻塞到吃完全部輸入、又常有單一運算子的成本壓過其他運算子,所以幾乎所有效益都來自分割。
  • 論文自承:病態的資料偏斜會擊潰整套做法,當大量甚至全部 tuple 的 join 屬性值相同時,所有 tuple 會落進同一個 hash bucket,而目前沒有已知演算法能做到 speedup 或 scaleup;range 分割本身也另外暴露在資料偏斜與隨之而來的執行偏斜之下。
  • 論文自列的未解問題:ad-hoc 查詢與 OLTP 混跑(大型查詢會取得大量鎖並長時間持有,逼得系統只能選擇髒讀或多版本機制;低優先度用戶端呼叫高優先度伺服器時還會出現優先度反轉)、當時的查詢最佳化器完全不考慮平行演算法與平行計畫形狀、缺乏實體設計工具、分割只能依單一屬性,以及公用程式問題(以每秒 1 MB 重整一個 TB 的資料庫要花超過十二天,除非把公用程式做成線上、增量、平行且可復原的)。
  • 自承且部分自相矛盾之處:結論承認某些應用領域不適合關聯式模型,呼籲發展物件導向資料庫系統;而東京大學 Super Database Computer 的專用硬體排序器與 omega network 也被作者承認直接牴觸了自己「專用硬體不是好投資」的論點。
  • 後來研究揭露的問題:論文把架構之爭視為已定案,但便宜的高頻寬網路與雲端物件儲存讓 shared-disk 在 Oracle RAC、Amazon Aurora 與 Snowflake 上復活;shared-nothing 把資料綁死在特定節點,也讓彈性擴縮、重新平衡與慢節點容忍變得棘手,而且論文對長查詢執行期間的節點故障幾乎隻字未提,這正是 MapReduce 與 Spark 後來視為第一等問題的部分。

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

這篇論文是 shared-nothing 分割式平行設計的權威宣言,此後每一套大規模平行資料庫都沿用這個架構。它的 split 與 merge 運算子是 Volcano exchange 運算子的直系祖先,Goetz Graefe 把它帶進 Microsoft SQL Server,至今仍是現代最佳化器表達平行的方式;同一個封裝概念也以 shuffle 階段的形式重現在 MapReduce、Hadoop、Spark、Presto 與 Flink 中。round-robin、hash、range 這套 declustering 分類成為實體設計的標準詞彙,並原封不動地以 sharding 之名活在 Bigtable、Dynamo、Cassandra、MongoDB、Spanner 與 CockroachDB 裡。商業脈絡則從 Teradata 與 Tandem NonStop SQL,經 IBM DB2 Parallel Edition 與 Informix XPS,一路延伸到 2000 年代的 MPP 資料倉儲 Netezza、Greenplum、Vertica 與 ParAccel,後者最終成為 Amazon Redshift。DeWitt 與 Gray 主張通用硬體的平行化推翻了 Grosch 定律,這正是今日所有雲端分析服務的經濟前提;DeWitt 本人也在 2008 至 2009 年與 Stonebraker 比較平行 DBMS 與 MapReduce 時,重新祭出這篇論文的量測紀律。歷史唯一逆轉這篇論文的地方在儲存層:Snowflake、BigQuery 與 Aurora 把運算與共享儲存層分離,等於恢復了 shared-disk 的拓樸,但仍完整保留這篇論文定義的分割式資料流執行模型。

論文原文 — 逐字引用

“Parallel database machine architectures have evolved from the use of exotic hardware to a software parallel dataflow architecture based on conventional shared-nothing hardware.(平行資料庫機器架構已經從使用奇特硬體,演進成建立在通用 shared-nothing 硬體之上的軟體平行資料流架構。)”

Abstract

“The shared-nothing design moves only questions and answers through the network.(shared-nothing 設計只讓問題與答案通過網路。)”

§2.2

“Parallelism is an unanticipated benefit of the relational model.(平行處理是關聯式模型帶來的意外紅利。)”

§2.3

術語 — 依本篇論文的用法

線性 speedup
指硬體規模或成本放大 N 倍,同一份固定工作就快 N 倍,量測方式是小系統耗時除以大系統耗時。speedup 固定問題大小,只放大硬體。
線性 scaleup
指放大 N 倍的系統能在相同時間內完成放大 N 倍的工作,也就是「小系統跑小問題的耗時」除以「大系統跑大問題的耗時」等於 1。交易 scaleup 同步放大用戶端數量、小請求數量與資料庫大小;批次 scaleup 則是同一句查詢跑在 N 倍大的資料庫上。
干擾(interference)
指每新增一個行程時,因為爭用全域記憶體、快取或互連網路等共用資源,而對其他所有行程造成的拖慢。這正是 shared-nothing 架構要極小化的障礙,因為即使只有 1% 的干擾也會讓 speedup 上限卡在 37。
偏斜(skew)
指平行各步驟的大小或成本變異數超過平均值,使整份工作的完成時間由最慢的一步決定,再增加平行度也幾乎沒有幫助。論文區分資料偏斜(多數 tuple 集中在單一分割)與執行偏斜(多數工作落在單一節點)。
shared-nothing
一種硬體架構,每塊記憶體與磁碟由單一處理器擁有並擔任該份資料的伺服器,處理器之間只透過互連網路傳送訊息溝通。由於原始的記憶體與磁碟存取都留在本地,網路上只流動過濾後的結果。
declustering(資料分割)
把單一關聯的 tuple 散佈到多顆磁碟上,每顆磁碟各自掛在自己的處理器下,使該關聯能被平行掃描。它是分割式執行在儲存層的前提,並且不需要 RAID 專用硬體就能取得多磁碟頻寬。
split 運算子
一種資料流節點,依每筆 tuple 的屬性值把某個運算子的輸出串流分割或複製成多條獨立串流,映射到指定的目的行程與埠。它同時實作緩衝與流量控制,輸出緩衝區塞滿時會擋住上游生產者。
merge 運算子
一種資料流節點,把多條平行資料串流併成單一循序串流,送進下游運算子的某個輸入埠。它與 split 搭配,使未經修改的循序關聯式運算子得以平行執行。
hash join
一種 join 演算法,先把兩個關聯依 join 屬性做 hash 分割,再把第一個關聯的某個分割建成主記憶體 hash table,並以第二個關聯的對應分割逐筆探測。它的成本是線性而非 n log n,也比 sort-merge join 更耐偏斜,除非輸入本來就已排序。

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

在時間軸上查看