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

GAMMA - A High Performance Dataflow Database Machine

第一台真正跑起來的 shared-nothing 平行資料庫:一顆處理器配一顆磁碟,關聯全部分割,查詢以自我排程的資料流執行。

作者David J. DeWitt、Robert H. Gerber、Goetz Graefe、Michael L. Heytens 等(另有 Krishna B. Kumar、M. Muralikrishna)- University of Wisconsin 計算機科學系 發表於VLDB 1986(第十二屆 Very Large Data Bases 國際會議論文集,京都,1986 年 8 月),頁 228-237 年份1986–1990s
閱讀原始論文 PDF 所有論文

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

Gamma 是一台完整運作的關聯式資料庫機器,由 20 台 VAX 11/750 以 80 Mbit/s 的 token ring 串接而成,其中八台各掛一顆 160 MB 磁碟。系統中所有關聯都以四種切分策略之一水平分割到全部磁碟上,因此掃描天生就是平行的,而且原始磁碟頻寬完全不需要經過網路。查詢樹被編譯成一串運算子行程,每個行程讀入 tuple 串流、處理後再透過 split table 這張把雜湊值或範圍值對應到目的行程的小表格把結果送出去;除了啟動運算子的兩則控制訊息與結束時的一則之外,整個執行過程完全自我排程。Join 採用雜湊分割,切成 building 與 probing 兩個階段,並把 bit vector filter 回送給產生內側關聯的行程,另有一套處理雜湊表溢位的再分割機制。實測結果是選擇與 join 從一顆磁碟擴到八顆都接近線性加速,單處理器的絕對時間足以和商用的 IDM500 競爭,論文還誠實留下一個關於網路介面壅塞的負面結果。

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

到了 1980 年代中期,資料庫機器這個領域只做出屈指可數的研究原型與三項商用產品,沒有任何一台證明高度平行的關聯式機器真能造得出來;其中商業上最成功的 Britton-Lee IDM500 根本完全沒有用到平行。作者自己先前的原型 DIRECT 經評估後被發現有嚴重設計缺陷:它把平行當成索引的替代品,於是把稀缺的 I/O 頻寬拿去掃描那些有索引就能跳過的資料。更糟的是,DIRECT 執行平行 join 所需的控制訊息數量與兩個輸入關聯大小的乘積成正比,訊息處理時間反而蓋過了真正的運算時間。同一時期的硬體趨勢極為殘酷:十年間單晶片 CPU 效能至少提升兩個數量級,商用磁碟頻寬卻只成長約三倍,I/O 因此成為所有既有機器設計的死結。當時流行的解法是把大量小磁碟串成高頻寬的儲存子系統,但那只是把問題往後推,因為資料接著仍得穿過一條同等頻寬的互連網路,而按照「90-10 法則」,搬過去的資料多半根本用不到。

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

  • 過去十年間處理器速度與磁碟頻寬嚴重脫節,CPU 效能大約提升兩個數量級,I/O 頻寬卻只成長三倍,這使得好幾種資料庫機器設計直接失去意義。
  • DIRECT 把平行當作索引的替代品,但索引的本質正是讓系統免於搜尋資料庫的一大塊,因此當 I/O 頻寬成為關鍵資源時,這種做法導致災難性的效能。
  • DIRECT 平行 join 演算法的控制成本與兩個輸入關聯大小的乘積成正比,即使訊息傳遞是透過共享記憶體實作,傳送與處理訊息的時間仍然支配了運算與 I/O 時間。
  • 就算真的做出有效頻寬達每秒 100 MB 的大量儲存子系統也沒用,除非前面那條互連網路同樣能吃下每秒 100 MB,而搬運過去的資料絕大多數根本用不上。
  • 當時沒有任何原型或產品證明高度平行的關聯式資料庫機器能被實際量測:Teradata 未發表任何效能數據且屢次拒絕作者的測試請求,DELTA 只公布了排序引擎的數字而且比超級迷你電腦上的商用排序套件還慢,MBDS 則完全沒有複雜運算的結果。
  • 既有的分割式檔案系統如 VSAM 與 Tandem file system 規定:檔案若依某個鍵分割,各站台上也必須依該鍵排序存放,等於強迫實體叢集方式跟著切分方式走。

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

一顆磁碟配一顆處理器

Gamma 不去打造特殊的平行讀出儲存子系統,而是讓一顆一般磁碟搭配一顆一般處理器,再用網路把這些配對連起來。五十顆每秒 2 MB 的磁碟能提供與單一每秒 100 MB 子系統相同的總頻寬,但網路不必再承載這個頻寬,因為選擇等會大幅縮減資料量的工作在磁碟端就先做完了。這個設計還讓 I/O 頻寬可以一組處理器加磁碟地逐步擴充,也更容易吸收磁碟技術的進步。這正是後來被稱為 shared-nothing 的架構,而 Gamma 是它第一個被實際量測的實作。

所有關聯一律水平分割

Gamma 中的每一個關聯都水平分割到系統中所有磁碟上,因此不存在「未分割的表」,每次掃描自動就是平行掃描。查詢語言提供四種切分策略:round robin、雜湊、由使用者指定各站台鍵值範圍的範圍分割,以及用平行合併排序算出均勻分布的範圍分割。關鍵在於,與 VSAM 或 Tandem file system 不同,Gamma 完全不要求切分屬性與站台內 tuple 的排列順序有任何關係,所以一個銀行關聯可以依帳號切分以追求吞吐量,同時在分行代號上建立叢集索引來服務彙總查詢。把切分與站台內叢集脫鉤,正是讓實體設計能同時服務兩種存取模式的原因。

split table 作為平行化的基本元件

每個運算子都寫得像是只跑在單一處理器上:讀進一串 tuple、吐出一串 tuple。平行完全由輸出端一個很小的資料結構注入,也就是 split table,它把每個輸出 tuple 導出的值對應到目的行程的位址。由於運算子程式碼完全不知道 split table 裡有什麼,同一份循序的 join 或 select 只要換一張表就能以任意平行度執行,運算子之間的重新分割除了轉送之外幾乎不花額外成本。把平行從運算子程式碼中抽離、收攏進一張路由表,是這篇論文最可重複利用的想法。

自我排程的資料流執行

排程器啟動一個運算子行程後,該行程先回報自己的身分,接著就不再需要任何監督:持續讀取輸入串流、施加自己的函式、透過 split table 轉送結果,直到偵測到輸入結束為止。關閉輸出串流會順帶對下游行程送出 end-of-stream 訊息,最後再以一則控制訊息向排程器回報完成。也就是每個運算子在每台處理器上只需三則控制訊息,兩則啟動、一則結束,且與流過多少 tuple 完全無關,這正面回應了 DIRECT 那種隨關聯大小乘積成長的訊息成本。其餘一切都只是資料在行程之間流動,沒有任何集中控制。

兩階段的分割式 hash join

對兩個來源關聯套用完全相同的雜湊 split table,可讓所有 join 屬性值相同的 tuple 落到同一個站台,於是兩個大關聯的 join 就分解成許多小 bucket 的獨立 join。Join 運算子先跑 building 階段,把第一個關聯吃進記憶體內的雜湊表,回報完成,等排程器確認所有站台都建完之後,才進入掃描第二個關聯的 probing 階段。兩階段之間的排程器屏障是唯一額外的控制互動,使得執行一次 hash join 的總控制成本為每站台五則訊息。若外側關聯本來就依 join 屬性切分,它根本不必傳輸,只需把內側關聯依外側各片段的範圍重新分配即可。

放進 split table 的 bit vector filter

每個 join 行程用外側關聯建雜湊表的同時,也把 join 屬性值雜湊進一個 bit vector filter。building 階段結束時,各行程把自己的 filter 送給排程器,排程器收齊後再轉發給負責產生內側關聯的行程,在那裡以陣列形式安裝進 split table。凡是不可能配對成功的內側 tuple,就在生產端直接丟棄,根本不會被放上網路。這其實就是用一個廉價的近似集合來做 semijoin,把 join 的選擇率直接換算成省下來的通訊量。

運作方式 — 具體的機制

硬體與系統軟體基礎

原型由 20 台 VAX 11/750 組成,每台配兩 MB 記憶體,以 Proteon 為該團隊打造的 80 Mbit/s token ring 相連,另有一台跑 Berkeley UNIX 的 VAX 擔任 host。二十台中有八台掛上 160 MB 的 Fujitsu 磁碟存放資料庫,其餘為無磁碟節點,可供 join 與 spool 工作使用。這些處理器跑的是專為資料庫工作撰寫的作業系統 NOSE,提供共享記憶體的輕量級行程、避免 convoy 現象的非搶占式排程,以及以計時器為基礎的一位元停等正向確認協定,並用 deltaT 機制重建序號。檔案、記錄、索引與掃描服務來自 Wisconsin Storage System,而 WiSS 的頁面格式直接內嵌 NOSE 的跨處理器訊息標頭,因此一個頁面可以從磁碟讀出後直接送給另一台處理器,不必把 tuple 複製到外送訊息樣板中。

儲存組織與多處理器索引

partition 指令指定切分策略,若是範圍分割還要給出鍵值邊界:partition employee on emp_id (100, 300, 1000) 會把 emp_id 不超過 100 的放到處理器 1,100 到 300 放處理器 2,300 到 1000 放處理器 3,其餘放處理器 4。若使用者提不出範圍,Gamma 會先以 round robin 載入,再用平行合併排序依切分屬性排序、重新分配以拉平各站台的 tuple 數,最後把各站台的最大鍵值回傳給 host。之後即可在每個片段上建立一般的叢集與非叢集索引。只要用了任一種範圍策略,Gamma 還會額外建一個多處理器索引:磁碟與其處理器本身就是一棵主要叢集索引的節點,而根頁面存放在 host 上的 schema 中;查詢最佳化器讀取這個根頁面,就能判斷像 q 介於 A 與 C 之間這種查詢只需送往處理器 1。

行程結構

Catalog Manager 是 host 上的常駐程式,把所有概念層與內部層的 schema 資訊存在 UNIX 檔案中,資料庫開啟時載入記憶體,並負責維持各使用者本地快取副本的一致性,內部另有鎖管理器保護 catalog。每個活躍使用者對應一個 Query Manager,負責 schema 快取、剖析、查詢最佳化與編譯。每個跨站台查詢由一個 Scheduler 行程控制,負責啟動編譯後查詢樹各節點的運算子行程;排程器刻意跑在機器內部而非 host 上,因為兩台查詢處理器之間的訊息速度是查詢處理器與 host 之間的兩倍,後者必須穿過 UNIX。每個參與的處理器上都有 Operator Process 執行對應的運算子;此外還有集中式的 Deadlock Detection Process,負責從各鎖管理器蒐集 wait-for 圖片段、找環並選出犧牲者,以及負責從查詢處理器蒐集日誌片段以支援提交、中止與回復的 Log Manager。

split table 與 tuple 路由

Gamma 使用三種 split table。雜湊型對 join 或切分屬性套用雜湊函式得到索引值,例如四台處理器時是 0 到 3,再查出目的地的處理器編號與埠號。範圍型以各分割範圍的上界作為鍵值,用在永久關聯採範圍分割時,也用在查詢樹葉節點的切分屬性剛好就是該關聯的水平切分屬性時;此時表格會用來源關聯自己的邊界值初始化,使每個片段都在本地處理、完全不必傳輸。第三種完全忽略鍵值,把 tuple 以 round robin 分送,這是結果關聯的預設做法。split table 中還會插入一組 bit vector filter 陣列,在轉送前先丟掉那些不可能參與 join 的 tuple。

Hash join 的控制流程

Join 的啟動方式與其他運算子一致,只多出一次控制互動。building 階段中,每個 join 行程把第一個來源關聯的 tuple 吃進記憶體內的雜湊表與 bit vector filter,然後向排程器回報建表完成。排程器必須聽到所有 join 行程都回報完畢,才會送出啟動 probing 階段的訊息;該階段中每個行程讀取第二個關聯的 tuple,以 join 屬性值去探查自己的雜湊表,最後每個行程再送一則訊息回報完成。實際上 building 與 probing 在控制層面被視為兩個獨立的運算子,這正是啟動與控制一次 hash join 的淨成本為每站台五則訊息的原因,也說明排程器必須同時安排兩個輸入串流的生產者,讓它們與兩個階段的時機吻合。

雜湊表溢位處理

若 building 階段中某些 bucket 長得太大,記憶體內的雜湊表就會溢位。本地的 join 運算子會收窄用來建表的 tuple 分割維度,切成兩個子分割:一個繼續灌入雜湊表,另一個傾倒到磁碟上的溢位檔(有可能在遠端),而已經進表卻屬於溢位子分割的 tuple 會被移出。運算子回報 building 完成時,會一併告知排程器它用了哪種再分割方案,排程器便能改寫 probing 端關聯的 split table,把對應的溢位子分割直接 spool 到磁碟、完全繞過 join 運算子。等非溢位的子分割 join 完之後,排程器再對那些被 spool 出去的溢位子分割遞迴套用 join。這套方法只在單一 join 屬性值所對應的 tuple 總大小超過可用記憶體時失效,此時改用以雜湊為基礎的 nested loops join 變形。

選擇與更新運算子

論文把選擇運算子的吞吐量視為整條管線的關卡,因為若選擇運算子餵不飽查詢樹,後續運算子能發揮的平行度就受限。Gamma 用三種互補技術:能用索引就用索引;把選擇述詞編譯成機器語言程序,讓述詞求值不必被直譯;以及採用有限度的預讀,讓一個頁面的處理與下一個頁面的 I/O 重疊。更新運算子(replace、delete、append)大致採用標準做法,唯一例外是修改到切分屬性的 replace:這種情況下不能把 tuple 寫回本地片段,而必須把改過的 tuple 推進 split table,由它決定該 tuple 現在應該落在哪個站台。

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

  • 原型是完全可運作的系統:20 台 VAX 11/750,每台兩 MB 記憶體,80 Mbit/s token ring,其中八台掛 160 MB 的 Fujitsu 磁碟,軟體跑在 NOSE 作業系統與 Wisconsin Storage System 之上。
  • 在每個關聯一萬筆、每筆 208 位元組(十三個四位元組整數加三個 52 位元組字串)、選擇率 10% 的合成資料上,無索引且針對切分屬性的選擇(S1)與針對非切分屬性的選擇(S3),從一顆磁碟到八顆磁碟的加速都相當接近線性;S3 甚至略快於 S1,因為 round robin 結果關聯的分送成本由所有處理器分攤,而 S1 是由單一處理器產生全部結果 tuple。
  • 單站台、針對切分屬性的索引選擇(S2)只在一到約三顆處理器之間有改善,之後就持平,因為增加磁碟只減少索引走訪的層數與存放結果的時間,卻不減少需讀取的葉層資料頁數,最終由那唯一產生結果的處理器成為瓶頸。
  • 與配備資料庫加速器與同等磁碟的商用 IDM500 相比,Gamma 的單處理器配置具有競爭力:IDM500 在 S1 選擇要 22.3 秒、S2 要 5.2 秒、J2 join 要 84.3 秒、J4 join 要 14.3 秒,而 Gamma 透過多處理器索引取回單筆 tuple 只需 0.14 秒。
  • 完全跑在無磁碟處理器上的 join(remote join)不但沒有變慢,反而略快於跑在有磁碟處理器上的 join(local join),原因是本地執行時 join 與 select 運算子會爭搶同一顆 CPU,而且 Gamma 在兩台處理器之間傳送循序 tuple 串流的速率幾乎與同機行程之間相同;這證明複雜運算子可以從儲存節點卸載出去。
  • 論文照實報告並診斷了一個負面結果:針對非切分屬性的索引選擇(S4)從四顆處理器加到八顆幾乎沒有改善,而且問題不在頻寬,因為重新分送 750 筆 208 位元組的結果 tuple 只有 120 萬位元,在 80 Mbit 的環網上約需兩百分之一秒;真正的原因是網路介面壅塞:逐筆 tuple 的 round robin 會讓 64 個輸出緩衝區幾乎同時被填滿,而介面只能緩衝兩個進來的封包,於是每個站台有五個封包必須重傳,還會與確認訊息互相碰撞。

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

  • 論文自承:這份評估明確只是初步且僅限單一使用者,沒有做多使用者測試,沒有評估更新運算子,彙總運算與彙總函式也尚未實作,量測的屬性值分布全為均勻;當時的開發重心是正確性而非絕對速度。
  • 論文自承:造成 S4 曲線持平的網路介面壅塞問題並未在論文中修掉,作者提出改用逐頁的 round robin 並隨機化第一個輸出頁的目的地,但沒有實測;此外測試中第七、第八顆磁碟是較舊的 14 吋機種,效能只有其他六顆的 82%,進一步扭曲了加速曲線。
  • 論文自承:join 的 building 與 probing 兩階段不重疊,因此 join 的回應時間受限於兩階段耗時之和;而且每台處理器上控制訊息對資料訊息的比例會隨處理器增加而上升,這在八顆磁碟時就已浮現——每個 join 運算子每個關聯只處理約十四個資料頁,卻要收發五則控制訊息。
  • 論文有提及但低估:當單一 join 屬性值所對應的 tuple 超過可用記憶體時,溢位機制直接失效,只能退回以雜湊為基礎的 nested loops join,而作者把非均勻屬性分布留給未來工作;資料傾斜正是雜湊分割式平行的核心弱點,後續關於傾斜處理與自適應重新分割的研究,正是從這個缺口長出來的。
  • 後續工作揭露:Gamma 刻意不提供站台自治,schema 集中管理、所有查詢由單一點發起,死結偵測器與排程器池也都集中,因此整個設計對節點故障與彈性成員變動毫無對策;後續的 Gamma 論文必須先把 VAX 與 token ring 硬體換成 32 節點的 Intel iPSC/2 hypercube,才能認真研究可擴展性。

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

Gamma 是把 shared-nothing 從一種主張變成一組實測數據的論文,其架構成為平行關聯式系統的標準範本,六年後由 DeWitt 與 Gray 在 CACM 的 Parallel Database Systems 一文正式定調。Teradata 的 DBC/1012、Tandem NonStop SQL、IBM DB2 Parallel Edition 與 Informix XPS 全都採用同一套核心配方:切分後的表、分割式 hash join,以及每個節點上的運算子行程。Gamma 共同作者 Goetz Graefe 後來把 split table 一般化為 Volcano 的 exchange 運算子,讓平行變成可以插進一個原本循序的查詢計畫中的單一運算子,如今 SQL Server、Greenplum、Vertica、Spark SQL、Presto 與 Trino、Snowflake 以及 Dremel 一脈的引擎全都以這種方式平行執行。兩階段 build-probe hash join 加上回送給生產端的 bit vector filter,正是今日 broadcast 與 shuffle hash join 以及 Spark、Impala、Snowflake 中執行期 Bloom filter 下推的祖先,而 Gamma 遞迴處理溢位子分割的做法則是 grace 與 hybrid hash join 落盤機制的前身。它的四種切分策略(round robin、雜湊,以及兩種範圍分割)正好就是今天每一套資料倉儲與分片式儲存仍在提供的分割選項。連論文誠實記下的失敗也留下了後代:那段網路介面壅塞的分析,是 MapReduce、Spark 以及每一種現代 exchange 實作都必須處理的 shuffle 緩衝與 incast 問題的早期案例。

論文原文 — 逐字引用

“the Gamma prototype shows how parallelism can be controlled with minimal control overhead through a combination of the use of algorithms based on hashing and the pipelining of data between processes.”

Abstract

“With the exception of these three control messages, execution of an operator is completely self-scheduling. Data flows among the processes executing a query tree in a dataflow fashion.”

§3.4

“To utilize the I/O bandwidth available in such a design, all relations in Gamma are horizontally partitioned across all disk drives.”

§6 Conclusions

術語 — 依本篇論文的用法

水平分割(horizontal partitioning / declustering)
把單一關聯的 tuple 打散到系統中所有磁碟上,使每個關聯都以片段形式存放、每次掃描天生就是平行的。Gamma 對所有關聯無一例外地套用此法,可選 round robin、雜湊或兩種範圍切分之一。
split table
掛在運算子行程輸出端的一張小表,把每個結果 tuple 導出的值(雜湊值、範圍上界,或 round robin 時什麼都不看)對應到目的行程的處理器與埠號。這是 Gamma 中唯一表達平行的地方,也正因如此運算子程式碼才能寫得像循序程式一樣。
水平切分屬性(HPA)
關聯據以切分的那個屬性。當查詢的限定屬性就是 HPA 且採用範圍分割時,最佳化器可以只把查詢送往相關站台;若限定屬性不是 HPA,查詢就必須送到每一個站台。
多處理器索引(multiprocessor index)
關聯採範圍分割時額外建立的索引,磁碟與其處理器本身就是一棵主要叢集索引的節點,根頁面與 schema 一起存放在 host 上。最佳化器查閱這個根頁面,把針對切分鍵的選擇查詢導向正確的處理器。
building 階段與 probing 階段
Gamma 分割式 hash join 的兩個半場:building 把第一個來源關聯吃進各站台記憶體中的雜湊表與 bit vector filter,probing 再把第二個關聯串流過去比對。兩者之間隔著一道排程器屏障,且在控制層面被當成兩個獨立的運算子。
bit vector filter
在 building 階段對外側關聯的 join 屬性值做雜湊而建出的近似集合,之後由排程器收齊、安裝到產生內側關聯那些行程的 split table 中。凡是沒通過 filter 的內側 tuple,在還沒上網路之前就被丟棄。
local join 與 remote join
論文用來區分「join 完全跑在有掛磁碟的處理器上」與「join 完全跑在無磁碟處理器上」的說法。實測 remote join 略快,證明複雜運算子可以從儲存節點卸載出去。
NOSE 與 WiSS
NOSE 是 Gamma 底層專門打造的作業系統,提供共享記憶體的輕量級行程、避免 convoy 現象的非搶占式排程,以及可靠的計時器式訊息協定。WiSS(Wisconsin Storage System)提供檔案、記錄、索引與掃描服務,其頁面格式內嵌 NOSE 訊息標頭,讓頁面不必複製就能直接送出。

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

在時間軸上查看