前言

這篇接續 Airbnb 的語意層:開發者體驗,前一篇談的是我們的語意層 Minerva 一路走來的工程決策。以 Airflow DAG 數量來看,Minerva 是 Airbnb 最大的資料框架,而且經常佔用公司一半以上的回填(backfill)運算量。這篇要談的是運算的部分:我們如何讓一萬多個 Minerva 來源在兩萬多個 Airflow DAG 上保持最新狀態。

運算

語意在語意層裡定義好之後,我們要把這些定義具體化(materialize)成實體資料,查詢層才能讀取。成千上萬個來源和定義隨時都在變動,每週還有一百多個 PR 在更新這些定義,因此我們需要一套能在這種規模下管理變更的框架。

首先,我們把整個語意層建模成一張圖(graph)。每個節點是一個註冊在 Minerva 裡的資料集,如果某個節點是由另一個節點衍生出來的,兩者之間就有一條邊。有了這個結構,只要某個定義被更新,我們就能找出圖上有哪一塊會跟著改變。運算框架接著判斷變更範圍,並自動回填受影響的資料集。

對每一個來源,Minerva 會持續調和(reconcile)資料集的期望狀態與實際狀態。兩者之間的落差就是一份回填工作,框架會把它派發出去。資料集被更新時,我們會追蹤它們的狀態(state)。運算會用這些狀態來觸發下游處理,查詢層則靠它們確保只查詢到最新的資料集。當某個資料集過時或損壞時,我們的清理程序(janitor process)會負責生命週期的善後。

以下各節會逐一深入說明。本篇後續會把「資料集」和「來源」當成同義詞交替使用。

偵測變更

Minerva 圖上的每個資料集都有一個資料版本(data version),也就是把 YAML 檔中所有描述該資料集語意的欄位取 MD5 雜湊後得到的值。每個帶版本的資料集都會在資料倉儲中具體化成一張 Iceberg 資料表。當資料生產者修改設定檔時,例如加上一個維度篩選條件、或改變某個欄位的投影方式,資料版本就會更新,等於告訴 Minerva 要重新計算實體資料,好符合新的定義。以上講的是單一資料集的情況。當資料集彼此相依時,事情就有趣多了。

Minerva 裡有些資料集是從其他資料集衍生出來的。假設某位使用者從事實表定義了一個事件來源,另一位使用者從維度表定義了一個維度來源,Minerva 可以把兩者預先 join 成一個叫做維度集(dimension set)的資料集,或者把事件預先彙總到某個特定的顆粒度,產生一個彙總來源(rollup source)。這些預先 join 和預先彙總的資料集都相依於那個事件來源,所以事件來源一變,它們也必須跟著更新。

我們用鏈式資料版本(chained data version)來管理這些相依關係。下游來源的資料版本會把上游來源的資料版本納入為輸入之一。在上面的例子裡,事件來源的資料版本會被併入維度集和彙總來源各自的資料版本中。事件來源一改,這兩個下游版本也跟著改。一次 YAML 修改,可能在整張圖上掀起一波更新。

這個設計能把變更順利地往下傳遞,但它的變更偵測有時過於激進。假設我們在事件來源上新增一個欄位,那些從來沒引用過該欄位的下游來源其實不需要重算,可是它們照樣會被回填,因為我們知道的只有「上游版本變了」這件事。變通做法是有的,例如把某個來源釘在固定的資料版本上,但一遇到大規模變更就會複雜到難以收拾。這是我希望我們當初能多投入一些的地方,因為更好的變更偵測可以顯著減少 Minerva 需要回填的量。

有些工具在這方面走得更遠。SQLMesh 就在精細的變更偵測上投注了大量心力,透過比對正規化後的 SQL AST,引入了破壞性變更與非破壞性變更這樣的概念。我認為這是資料轉換工具該走的方向。

狀態調和(Reconciliation)

變更偵測講完之後,我們可以聚焦在單一資料集上,看看運算框架是怎麼調和這些變更的。由於 Airbnb 大多數資料集都是依日期分區(date-partitioned)的,調和演算法非常倚重分區層級的操作。

調和演算法

每天,每個資料集都會盤點哪些輸入分區已經存在、哪些輸出分區已經寫好。從兩者的差集產生一份計畫,把缺少的日期分區打包成一個個批次視窗,平行回填。

實務上,這個調和過程是以「每個資料集、每一天」為單位,跑在一個 Airflow DAG 裡。改用一個集中式的控制平面去逐一巡過所有資料集其實也行得通,因為調和演算法本身並不在意工作是怎麼被排程的。同一套演算法要應付好幾種情境,而每種情境派發工作的方式都不太一樣。

例行執行

最常見的情境是某個來源的定義完全沒變,每天有新的輸入資料進來,Minerva 的任務就是把它具體化。演算法會發現所有輸入分區都存在,而輸出分區只到前一天為止,於是產生一個只涵蓋一個日期分區的批次。那個分區寫完,資料集就是最新的了。

分區是讓增量運算得以成立的關鍵。不過有些資料集會有延遲到達的事件,例如取消或更改訂單,這類資料的歷史在事後很久都還可能變動,沒辦法用增量方式處理。我們對這些資料集的做法是每天從頭重跑一次,等於每次執行都是一次完整的歷史回填。這成本很高,所以我們後來擴充了演算法,讓使用者可以指定一個輸出視窗,只重跑最近 X 天,而不是整段歷史。這在實務上行得通,因為大多數使用者並不需要回溯那麼久以前的資料。

離線回填

Minerva 的第一個版本裡,所有變更都直接在正式環境回填。我們很快就發現,隨著使用者變更商業定義的頻率愈來愈高,Minerva 根本跟不上,因為回填期間資料是不可用的。這造成了一連串的可用性問題,最後演變成一次重大事故,我在 Airbnb 十年回顧 裡描述過。我們學到的是:使用者需要一個隔離的環境來做回填。

幸運的是,每個帶版本的資料集本來就會具體化成自己的一張資料表,所以我們可以直接啟動回填,把新資料寫進一張與正式環境完全隔離的資料表裡。同一套調和演算法在這裡照樣適用:它看到所有輸入分區都存在、而輸出的資料表是空的,於是從最初的日期一路回填到今天。這可能要跑上一段時間,所以演算法會把工作切成互不重疊的批次平行處理,大幅縮短回填時間。

離線回填讓使用者可以不受時間壓力地重建資料集,而且資料集一旦回填完成,上線到正式環境是瞬間完成的。SQLMesh 這類工具也得出了相同的結論,並推出了 Virtual Data Environments 這樣的功能。

線上回填

最後一種情境介於例行執行和完整歷史回填之間。偶爾,使用者會發現餵給 Minerva 的輸入資料表損毀了。要修正的話,就得從修好的輸入資料表重新產出受影響的那部分 Minerva 資料,同時不能動到其他分區。這正是線上回填在做的事:使用者可以像動手術一樣,只就地重新產出那一小部分的 Minerva 分區。

我們把自我修復的演算法改造來處理這件事。當線上回填被觸發時,我們會告訴演算法某一批輸出分區「不存在」,藉此強制它們重跑。接下來演算法就走它平常的流程,建立批次視窗並派發回填工作。資料會寫回同一張正式資料表、掛在同一個資料版本下,因為這裡並沒有引入新的語意。

把這些串起來

在我離開 Airbnb 之前,團隊一直在把這些各自為政的流程整合成單一工作流程,中間切成幾個明確的階段。一個新的資料集通常會先經過試跑(dry run),接著是離線回填,最後才上線到正式環境,開始例行執行。

追蹤來源狀態

面對成千上萬個必須按正確順序處理的來源,我們得非常仔細地追蹤每一個的狀態。這件事之所以更重要,是因為回填是一個漫長的非同步過程,快慢取決於當下有多少運算資源或排程資源可用。

我們的解法是「來源狀態」,一組用來寫入和讀取來源當前狀態的 API 端點。它會追蹤來源最新的指紋(fingerprint)、資料版本,以及目前可用的分區。一旦某個來源完整回填完成,我們就向這組 API 發出更新,運算和查詢層都會從這裡讀取狀態。

在運算這端,每個資料集都排程在自己的 Airflow DAG 裡,所以我們做了一個自訂的來源狀態 sensor 來協調 DAG 之間的相依關係。它會去輪詢某個 DAG 所相依的上游資料集的狀態,確認它們是否已經就緒,只有就緒了才會放行。這能避免下游資料集在輸入還沒準備好之前就開始跑。

查詢層則是透過指紋來使用來源狀態。指紋是剛寫入的那批實體資料的雜湊值。和資料版本不同,指紋能告訴你實體資料是否真的改變了。對 Iceberg 資料表來說,它就是新 Iceberg 快照的 snapshot_id;其他引擎則改用資料表名稱加上最新分區日期的 SHA256 雜湊。有了指紋,查詢層就能拿它當快取的依據,來服務常見的查詢。

生命週期管理

新的資料集上線的同時,舊的會逐漸過時,Minerva 會跑一個清理程序把它們收掉。刪除分成兩個階段:先軟刪除,再硬刪除。軟刪除的候選名單是根據使用情況決定的,資料來自彙總後的查詢日誌。當一個來源被軟刪除時,我們會把它的資料資產從資料目錄的搜尋中隱藏起來,並暫停相關的 pipeline,但不刪除任何資料。

過了寬限期之後,才會執行硬刪除。硬刪除會同時清掉設定檔和底層資料。在設定檔這一側,我們有一組 deleter 會走過 Minerva 圖,把定義在某個來源下游的所有東西清乾淨。那份變更合併之後,對應的 Airflow DAG 會被拆掉,另一個清理程序則負責把底層的資料表和它們佔用的儲存空間刪除。

我們今天有一套相當健全的生命週期管理流程。我很希望我們在 Minerva 更早期就把它做成這樣。除非絕大部分都自動化,否則要求使用者自己做這類清理是很難推動的。在我維運 Minerva 的那段期間,靠著這些生命週期管理工具,我們找到了許多降低浪費與成本的機會。

總結

這篇談的是 Minerva 最核心的挑戰:語意定義隨時都在被建立、更新和刪除,而打造一套能讓這些資料集持續保持最新且一致的運算框架,正是整個系統的核心所在。

前面我談了這個平台的幾項關鍵能力:變更偵測、狀態調和、狀態追蹤,以及生命週期管理,還有它們是如何彼此配合、把系統撐起來的。在下一篇、也是這個系列的最後一篇裡,我們會進入查詢層,看看它如何讀取運算層產出的資料集,並在規模化的前提下把它們送到使用者手上。