一顆磁碟配一顆處理器
Gamma 不去打造特殊的平行讀出儲存子系統,而是讓一顆一般磁碟搭配一顆一般處理器,再用網路把這些配對連起來。五十顆每秒 2 MB 的磁碟能提供與單一每秒 100 MB 子系統相同的總頻寬,但網路不必再承載這個頻寬,因為選擇等會大幅縮減資料量的工作在磁碟端就先做完了。這個設計還讓 I/O 頻寬可以一組處理器加磁碟地逐步擴充,也更容易吸收磁碟技術的進步。這正是後來被稱為 shared-nothing 的架構,而 Gamma 是它第一個被實際量測的實作。
第一台真正跑起來的 shared-nothing 平行資料庫:一顆處理器配一顆磁碟,關聯全部分割,查詢以自我排程的資料流執行。
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 法則」,搬過去的資料多半根本用不到。
Gamma 不去打造特殊的平行讀出儲存子系統,而是讓一顆一般磁碟搭配一顆一般處理器,再用網路把這些配對連起來。五十顆每秒 2 MB 的磁碟能提供與單一每秒 100 MB 子系統相同的總頻寬,但網路不必再承載這個頻寬,因為選擇等會大幅縮減資料量的工作在磁碟端就先做完了。這個設計還讓 I/O 頻寬可以一組處理器加磁碟地逐步擴充,也更容易吸收磁碟技術的進步。這正是後來被稱為 shared-nothing 的架構,而 Gamma 是它第一個被實際量測的實作。
Gamma 中的每一個關聯都水平分割到系統中所有磁碟上,因此不存在「未分割的表」,每次掃描自動就是平行掃描。查詢語言提供四種切分策略:round robin、雜湊、由使用者指定各站台鍵值範圍的範圍分割,以及用平行合併排序算出均勻分布的範圍分割。關鍵在於,與 VSAM 或 Tandem file system 不同,Gamma 完全不要求切分屬性與站台內 tuple 的排列順序有任何關係,所以一個銀行關聯可以依帳號切分以追求吞吐量,同時在分行代號上建立叢集索引來服務彙總查詢。把切分與站台內叢集脫鉤,正是讓實體設計能同時服務兩種存取模式的原因。
每個運算子都寫得像是只跑在單一處理器上:讀進一串 tuple、吐出一串 tuple。平行完全由輸出端一個很小的資料結構注入,也就是 split table,它把每個輸出 tuple 導出的值對應到目的行程的位址。由於運算子程式碼完全不知道 split table 裡有什麼,同一份循序的 join 或 select 只要換一張表就能以任意平行度執行,運算子之間的重新分割除了轉送之外幾乎不花額外成本。把平行從運算子程式碼中抽離、收攏進一張路由表,是這篇論文最可重複利用的想法。
排程器啟動一個運算子行程後,該行程先回報自己的身分,接著就不再需要任何監督:持續讀取輸入串流、施加自己的函式、透過 split table 轉送結果,直到偵測到輸入結束為止。關閉輸出串流會順帶對下游行程送出 end-of-stream 訊息,最後再以一則控制訊息向排程器回報完成。也就是每個運算子在每台處理器上只需三則控制訊息,兩則啟動、一則結束,且與流過多少 tuple 完全無關,這正面回應了 DIRECT 那種隨關聯大小乘積成長的訊息成本。其餘一切都只是資料在行程之間流動,沒有任何集中控制。
對兩個來源關聯套用完全相同的雜湊 split table,可讓所有 join 屬性值相同的 tuple 落到同一個站台,於是兩個大關聯的 join 就分解成許多小 bucket 的獨立 join。Join 運算子先跑 building 階段,把第一個關聯吃進記憶體內的雜湊表,回報完成,等排程器確認所有站台都建完之後,才進入掃描第二個關聯的 probing 階段。兩階段之間的排程器屏障是唯一額外的控制互動,使得執行一次 hash join 的總控制成本為每站台五則訊息。若外側關聯本來就依 join 屬性切分,它根本不必傳輸,只需把內側關聯依外側各片段的範圍重新分配即可。
每個 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。
Gamma 使用三種 split table。雜湊型對 join 或切分屬性套用雜湊函式得到索引值,例如四台處理器時是 0 到 3,再查出目的地的處理器編號與埠號。範圍型以各分割範圍的上界作為鍵值,用在永久關聯採範圍分割時,也用在查詢樹葉節點的切分屬性剛好就是該關聯的水平切分屬性時;此時表格會用來源關聯自己的邊界值初始化,使每個片段都在本地處理、完全不必傳輸。第三種完全忽略鍵值,把 tuple 以 round robin 分送,這是結果關聯的預設做法。split table 中還會插入一組 bit vector filter 陣列,在轉送前先丟掉那些不可能參與 join 的 tuple。
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 現在應該落在哪個站台。
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.”
“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.”
“To utilize the I/O bandwidth available in such a design, all relations in Gamma are horizontally partitioned across all disk drives.”