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

Hive - A Warehousing Solution Over a Map-Reduce Framework(Hive:建構於 Map-Reduce 框架之上的資料倉儲方案)

把 SQL 式資料倉儲蓋在 Hadoop 上:HiveQL 編譯成 map-reduce DAG,並搭配系統目錄、分割、bucket 與可插拔 SerDe。

作者Ashish Thusoo、Joydeep Sen Sarma、Namit Jain、Zheng Shao 等人(Facebook Data Infrastructure Team) 發表於VLDB 2009,法國里昂(demonstration 論文) 年份2009
閱讀原始論文 PDF 所有論文

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

Facebook 的資料量已經超過商用資料倉儲能負擔的成本,但 Hadoop 只提供最底層的 map-reduce API,每一份報表都得靠工程師手寫一支 Java 程式。Hive 在 Hadoop 之上蓋了一層資料倉儲:table、partition 與 hash bucket 實際上就是 HDFS 的目錄與檔案,語言是類 SQL 的 HiveQL,再加上一個保存 schema、儲存位置與序列化類別的 Metastore 系統目錄。Driver 會把每一道 HiveQL 交給編譯器,先剖析、再依 metastore 做型別檢查,接著建立邏輯運算子樹、套用述詞下推與分割裁剪等規則式改寫,最後在 repartition 與 union all 這兩種標記處把樹切成一張 map-reduce 工作的 DAG。因為任意檔案格式都能透過可插拔的 SerDe 讀取,任意邏輯也能以串流式 map-reduce 腳本嵌進查詢,這層 SQL 從來不會變成困住資料的牢籠。論文寫作當時,Facebook 這套系統已有數千張表、超過 700 TB 資料,服務一百多位使用者、每天五千多道查詢。

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

到了 2009 年,業界為了商業智慧所蒐集的資料規模成長飛快,傳統資料倉儲方案的價格已高到難以負擔,Facebook 於是把日誌處理搬到 Hadoop,也就是跑在一般商用硬體上的開源 map-reduce 實作。問題出在介面:map-reduce 是非常底層的程式模型,每冒出一個新的商業問題,就得有工程師從頭撰寫、除錯,之後還要長期維護一支專用程式。這些程式難以重用,把檔案的實體配置寫死在程式碼裡,也讓只懂 SQL、不懂 Java 的分析師完全碰不到資料。當時已有其他 map-reduce 前端,最有名的是 Pig 與微軟的 Scope,但它們是資料流語言、沒有系統目錄,schema 與儲存格式散落在各個腳本中,而不是集中在一個共享、可查詢的地方。Hive 就是為了讓 Hadoop 叢集表現得像 Facebook 已經買不起的那座資料倉儲而寫的。

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

  • 面對業界為商業智慧所蒐集的資料量,傳統資料倉儲方案的成本已高到難以承受。
  • map-reduce 程式模型非常底層,逼得開發者必須撰寫難以維護、也難以重用的客製程式。
  • 分析師就算能用 SQL 把問題講清楚,若沒有工程師幫忙翻成 map-reduce 程式,就完全無法在 Hadoop 資料上跑起來。
  • 倉儲資料以各式各樣的序列化格式進來,只有單一固定 IO 函式庫的系統,根本讀不了它被要求查詢的那些檔案。
  • Pig、Scope 這類同性質的 map-reduce 前端沒有系統目錄,schema 與統計資訊無處可放,也就談不上資料探索與查詢最佳化。
  • 對同一份輸入做好幾種不同的彙總時,天真的做法是每道查詢各掃一次,而在倉儲規模下掃描成本正是總成本的大宗。

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

HiveQL:編譯成 map-reduce 的宣告式語言

Hive 提供類 SQL 的宣告式語言,支援 select、project、join、彙總、union all,以及 from 子句中的子查詢;DDL 可指定序列化格式、分割欄位與 bucket 欄位,DML 則有 load 與 insert。編譯器會把每一道敘述轉成執行計畫,使用者只說要什麼,至於要拆成幾個 map-reduce 工作則由 Hive 決定。這等於把宣告式查詢的老論點移植到批次執行引擎上:shuffle 邊界放哪裡、要讀哪些分割,這些實體決策從分析師手上收回到編譯器手上。在 Facebook 立刻兌現的效益是,會寫 SQL 的人不再需要一位 Java 工程師當中介。

Metastore:一個名副其實的系統目錄

Hive 維護一份系統目錄,內容包含 database、table 與 partition,記錄欄位與型別、擁有者、儲存位置、bucket 資訊,以及該份資料的序列化/反序列化類別名稱。論文明白指出,正是這個 metastore 讓 Hive 像 Oracle 或 DB2 那樣是一套傳統資料倉儲方案,而與 Pig、Scope 這類 map-reduce 前端區隔開來。有了目錄,schema 只在建表時宣告一次、之後每次引用都重複使用,因此型別檢查、select star 展開與分割裁剪都能在編譯期完成。它同時讓整座倉儲變得可瀏覽:使用者不必打開任何一個檔案,就能知道有哪些表、裡面裝了什麼。

以目錄結構承載 table/partition/bucket

一張 table 就是一個 HDFS 目錄,一個 partition 就是以欄位值命名的子目錄,例如 /wh/T/ds=20090101/ctry=US,而一個 bucket 就是分割目錄底下、由某欄位雜湊決定的單一檔案。換句話說,整個實體資料模型都編碼在路徑名稱裡,Hive 與 Hadoop 不需要任何索引就能解讀。這正是分割裁剪如此便宜的原因:排除某個述詞範圍等同於不去列出某個目錄,而且沒有額外結構需要維護。Bucket 把同一招再往下推一層,讓抽樣查詢可以整個檔案跳過,也讓 join 有一份預先雜湊好的配置可以利用。

SerDe:可插拔序列化與讀時綱要

Hive 不要求資料先載入某種專有格式,而是讓每張表關聯到一種序列化格式並把這層關聯寫進目錄;內建格式會運用壓縮與延遲反序列化,使用者也能用 Java 撰寫自訂的 serialize 與 de-serialize 方法來支援新格式。由於 SerDe 類別是中繼資料而不是寫在查詢裡的程式碼,編譯器與執行引擎會在編譯期與執行期自動採用。於是檔案可以註冊成 external table,涵蓋 HDFS、NFS 或本機目錄,就地被查詢而不必搬動。這就是所謂讀時綱要(schema on read)的姿態,也是 Hadoop 倉儲與先載入再查詢的關聯式倉儲最根本的分野。

Multi-table insert:共用一次掃描

一道 HiveQL 敘述可以在 from 子句後面接好幾個 insert 子句,把針對同一份輸入的多道查詢一起寫出來。Hive 的最佳化方式是共用輸入資料的掃描,在論文的例子中,status updates 與 profiles 那個昂貴的 join 只做一次,卻同時產出 gender summary 與 school summary。在一個讀輸入就是主要成本、而且每個階段邊界都要寫回 HDFS 的系統上,把掃描成本攤到多張輸出表,是能拿到的最大紅利之一。這其實是使用者手動宣告版的多查詢最佳化(multi-query optimization),而論文把它的一般形式列為未來工作。

為 SQL 說不出口的事留下逃生口

HiveQL 允許以 Java 撰寫的使用者自訂欄位轉換函式(UDF)與彙總函式(UDAF),更激進的是允許使用者用任何語言嵌入自訂 map-reduce 腳本,透過 MAP 與 REDUCE 子句、以列為單位的串流介面從標準輸入讀列、往標準輸出寫列。論文自己的例子就用 Python 寫的 meme-extractor 當 mapper、用 top10.py 當 reducer,原因正是 Hive 當時還沒有 rank 彙總函式。這等於承認一個年輕的 SQL 方言不可能涵蓋所有需求,並把這個缺口變成插拔點而不是一堵牆。論文也坦白指出,這份彈性的代價是每一列都要在字串之間來回轉換。

運作方式 — 具體的機制

HDFS 上的實體配置

每張表對應一個 HDFS 目錄,資料列依目錄中記錄於系統目錄的格式序列化後存成檔案。分割欄位根本不存在資料檔裡,而是編碼在子目錄路徑上:表 T 若以 ds 與 ctry 分割,ds=20090101 且 ctry=US 的資料就落在 /wh/T/ds=20090101/ctry=US 底下的檔案中。在一個分割之內,資料還可依某欄位的雜湊切成 bucket,每個 bucket 具體化為分割目錄裡的一個檔案。欄位型別可以是基本型別(整數、浮點數、一般字串、日期、布林),也可以是可巢狀的集合型別 array 與 map,使用者還能以程式方式自訂型別。

Metastore 的內容與它為何不放在 HDFS

metastore 收納三類物件:Database 作為表的命名空間(未指定時用名為 default 的資料庫);Table 記錄欄位與型別、擁有者、儲存資訊(資料位置、資料格式、bucket 資訊)、SerDe 實作類別及其支援資訊,還可放任意使用者鍵值資料,將來用來存放表的統計資訊;Partition 則可擁有自己的欄位、SerDe 與儲存資訊,作為日後支援綱要演進(schema evolution)的伏筆。它的儲存系統必須為隨機存取與更新的線上交易而最佳化,而 HDFS 是為循序掃描而生、並不適合,因此 metastore 跑在 MySQL、Oracle 這類關聯式資料庫,或 local、NFS、AFS 這類檔案系統上。好處是只碰中繼資料的 HiveQL 敘述可以用極低延遲完成。代價則是資料與中繼資料分居兩個系統,Hive 必須自己明確維持兩者的一致性。

Driver、Thrift server 與請求路徑

外部介面包含命令列(CLI)、網頁 UI,以及 JDBC 與 ODBC API,全部透過 Thrift server 進入 Hive,該伺服器對外提供一組執行 HiveQL 的極簡客戶端 API;由於 Thrift 能為多種語言產生客戶端,同一支 Java 伺服器可同時支撐 Java 的 JDBC、C++ 的 ODBC,以及 php、perl、python 的腳本驅動程式。Driver 收到敘述後會建立 session handle,用來追蹤執行時間、輸出列數等統計,並管理該敘述在編譯、最佳化與執行整個生命週期中的狀態。編譯器回傳 DAG 之後,driver 依拓撲順序把個別 map-reduce 工作送進執行引擎。執行引擎目前就是 Hadoop,而 Hive 的其他所有元件都會與 metastore 互動。

編譯器前端:parser 與語意分析器

Parser 把查詢字串轉成剖析樹,語意分析器再把剖析樹轉成以 block 為單位的內部查詢表示。過程中它會向 metastore 取得輸入表的 schema 資訊,據此驗證欄位名稱、把 select star 展開成實際欄位清單,並做型別檢查、必要時插入隱含型別轉換。系統目錄與查詢就是在這一步交會:沒有 metastore,這個階段根本無從檢查起。若是 DDL 敘述,產生的計畫只含中繼資料操作;若是 LOAD 敘述,只含 HDFS 操作;因此只有 insert 與查詢會繼續往下產生 map-reduce 工作。

邏輯計畫與規則式最佳化器

邏輯計畫產生器把內部表示轉成一棵邏輯運算子樹,最佳化器接著對它做多趟改寫:把共用相同 join key 的多個 join 併成單一多路 join,於是壓成一個 map-reduce 工作;在 join、group by 與自訂 map-reduce 運算子前插入 repartition 運算子(即 ReduceSinkOperator),標示 map 階段結束、reduce 階段開始的邊界;提早裁剪欄位、把述詞往表掃描端下推,以減少運算子之間傳輸的資料量;對分割表裁掉查詢用不到的分割,對抽樣查詢裁掉用不到的 bucket。使用者另可對最佳化器下提示(hint):為高基數的分組彙總加入部分彙總(partial aggregation)運算子、為分組彙總的資料傾斜加入 repartition 運算子,或把 join 從 reduce 階段改到 map 階段執行。

實體計畫產生:把運算子樹切成工作

實體計畫產生器把邏輯計畫轉成 map-reduce 工作的 DAG:對每一個標記運算子(repartition 與 union all)開一個新的 map-reduce 工作,再把夾在標記之間的邏輯計畫片段指派給這些工作的 mapper 與 reducer。同一個工作內部,ReduceSinkOperator 以下的運算子樹由 mapper 執行、以上的由 reducer 執行,而重新分割這件事本身則由執行引擎完成,不是某個運算子做的。圖 2 顯示 multi-table insert 例子被編譯成三個工作:第一個做 join 並寫出 tmp1 與 tmp2 兩個 HDFS 暫存檔,第二與第三個工作分別讀取它們產出 gender 與 school 的彙總結果。這層具體化同時也是同步點,因為第二與第三個工作必須等第一個工作跑完。

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

  • 論文寫作當時,Facebook 的 Hive 倉儲已有數千張表、超過 700 TB 資料,並被一百多位使用者大量用於報表與臨機分析(第 1 節)。
  • 同一套倉儲每天支撐超過 5000 道查詢,且在 Facebook 內外都有活躍的使用者與開發者社群(第 5 節)。
  • 在依 Pavlo 等人(SIGMOD 2009)基準測試所做的初步實驗中,團隊把 Hadoop 本身的效能相對已發表數字改善了 20%,主要靠改用更快的 Hadoop 資料結構,例如以 Text 取代 String。
  • 同樣的查詢改用 HiveQL 撰寫,相對於那份手工最佳化的 Hadoop 實作有 20% 額外負擔,作者據此認為 Hive 的效能與該比較研究中的 Hadoop 程式相當。
  • multi-table insert 的例子會編譯成剛好三個 map-reduce 工作的 DAG(圖 2),其中一次 join 掃描餵出兩個 HDFS 暫存檔,再由兩個彙總工作消費。
  • 只存取中繼資料物件的 HiveQL 敘述能以極低延遲執行,因為 metastore 背後是關聯式資料庫或一般檔案系統,而不是 HDFS(第 3.1 節)。

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

  • 論文自承:HiveQL 目前只接受 SQL 的子集,不支援更新或刪除既有表中的資料列,也沒有 rank 彙總函式,這正是範例必須靠外部 Python reduce 腳本算 top-10 memes 的原因。
  • 論文自承:最佳化器是規則不多的天真規則式最佳化器,沒有成本式(cost-based)計畫選擇,使用者得自己下提示才能拿到 map 端 join、部分彙總或傾斜處理;成本式與適應式最佳化被列為未來工作。
  • 論文自承:自訂 map-reduce 腳本的串流介面必須把每一列在字串之間來回轉換;而且儲存是列式的,欄式儲存與更聰明的資料放置在當時都還只在探索階段。
  • 論文自承:由於 metastore 刻意放在 HDFS 之外的交易型儲存上,Hive 必須明確維持中繼資料與資料檔之間的一致性;這道分裂在日後確實造成了併發寫入與分割清單過期等真實問題。
  • 後續研究揭露:每個階段邊界都要具體化到 HDFS,下游工作只能乾等上游跑完,因此連小查詢也要花上好幾分鐘;而論文提供的效能數字僅屬初步且為自行量測。這道延遲下限正是 Tez、Spark、Impala 與 Presto 想要拆掉的東西。

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

Hive 讓 SQL 成為大數據的預設介面,之後幾乎每一套 SQL-on-Hadoop 引擎不是抄它的介面,就是直接接上它的元件。其中 Hive Metastore 的壽命甚至超過 Hive 自己的執行引擎,成為事實上的目錄標準:Spark SQL、Presto 與 Trino、Impala、Drill 與 Flink SQL 都從中讀取表、分割與 SerDe 的中繼資料,AWS Glue Data Catalog 更是刻意做成 Hive metastore API 相容。Hive 式分割(column=value 的目錄命名慣例)成為物件儲存資料湖的交換配置,而 Iceberg、Delta Lake 與 Hudi 這些現代表格式,設計初衷有很大一部分正是補上這種配置做不到的事:跨分割的原子提交、檔案層級統計,以及單靠列目錄無法提供的快照隔離。SerDe 的想法則一般化成整個生態系的讀時綱要立場,而論文中說正在探索的欄式儲存,後來以 RCFile、ORC 與 Parquet 的形式落地。Hive 自己也繞著論文承認的弱點重建:從 map-reduce 換到 Tez、加入向量化執行、透過 Calcite 引入查詢最佳化的成本模型、支援 ACID 交易與 LLAP 快取。更深一層的遺產是架構形狀:宣告式語言在上、共享系統目錄在中、可替換的執行引擎在下,如今已是資料湖查詢系統的標準樣貌。

論文原文 — 逐字引用

“However, the map-reduce programming model is very low level and requires developers to write custom programs which are hard to maintain and reuse.”

§1 Introduction

“The metastore distinguishes Hive as a traditional warehousing solution (ala Oracle or DB2) when compared with similar data processing systems built on top of map-reduce like architectures like Pig [7] and Scope [2].”

§3.1 Metastore

“Hive currently has a naïve rule-based optimizer with a small number of simple rules.”

§5 Future Work

術語 — 依本篇論文的用法

Table(表)
具名的關聯,其資料存放在單一 HDFS 目錄中,依系統目錄記錄的格式序列化成檔案。Hive 也支援 external table,直接查詢已存在於 HDFS、NFS 或本機目錄的資料。
Partition(分割)
表的一種切分方式,分割欄位的值編碼在目錄路徑上而非存進資料列,例如 /wh/T/ds=20090101/ctry=US。分割決定資料在表目錄內的分布,也讓最佳化器能整個目錄地裁掉。
Bucket
分割之下再依某欄位雜湊做的切分,每個 bucket 存成分割目錄中的一個檔案。Bucket 讓抽樣查詢能整個檔案跳過,也給 join 一份預先雜湊好的配置。
SerDe
以 Java 撰寫的自訂 serialize 與 de-serialize 方法組,讓 Hive 能讀寫某種資料格式。SerDe 的實作類別存放在系統目錄中,於查詢編譯與執行時自動套用。
Metastore
Hive 的系統目錄,收錄 database、table 與 partition,內含欄位與型別、擁有者、儲存位置、bucket 與 SerDe 資訊。因為需要隨機存取更新,它跑在關聯式資料庫或一般檔案系統上,而不是 HDFS。
HiveQL
Hive 的類 SQL 宣告式語言,涵蓋 select、project、join、彙總、union all、from 子句中的子查詢,以及可指定序列化與分割選項的建表 DDL 和 load/insert DML。它會被編譯成執行計畫,而非逐列直譯。
Multi-table insert
單一 HiveQL 敘述對同一份輸入跑多道查詢,並把各自結果寫入不同的表或分割。Hive 的最佳化方式是讓所有輸出共用同一次輸入掃描。
Repartition 運算子(ReduceSinkOperator)
最佳化器在 join、group by 與自訂 map-reduce 運算子前插入的標記運算子,代表一次 shuffle。它標示 map 階段與 reduce 階段的邊界,實體計畫產生器會在每個這種標記處開啟新的 map-reduce 工作。
External table(外部表)
資料留在原處(HDFS、NFS 或本機目錄)就地查詢的表,而不是載入倉儲目錄中的資料。它是讓既有檔案免複製即可被查詢的機制。

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

在時間軸上查看