使用 Lakeflow 管線建立管線
生產資料管線需要可靠性、可維護性及明確的資料品質強制執行。 作為資料工程師,你很可能會花大量時間撰寫程式碼來處理增量處理、協調相依性,以及驗證資料品質。 Azure Databricks 中的 Lakeflow 管線透過讓你定義資料應該呈現的樣貌,而非定義如何逐步處理資料,來解決這些挑戰。
在這個單元中,你將學習如何利用宣告式方法建立資料管線、定義串流資料表與實體化檢視,並應用資料品質期望來對資料施加限制。
了解宣告式方法
傳統的資料管線要求你撰寫命令式程式碼,明確規定每個處理步驟。 你負責增量處理邏輯、管理檢查點恢復,以及協調資料表間的相依關係。 而 Lakeflow pipelines,則是宣告 想要的最終狀態,框架則負責執行細節。
宣告式方法為生產管線帶來三大主要優勢:
自動協調:框架分析資料表間的相依關係,並決定正確的執行順序。 如果你定義一張銀色表格,從青銅表格讀取,框架會自動先處理青銅表格。
增量處理:串流資料表每筆記錄只處理一次。 實體化視圖會在可能時自動識別並僅處理已變更的資料,避免昂貴的全面重算。
內建重試邏輯:當失敗發生時,框架會以最細緻的層級重試——先是 Spark 任務層級,接著是流程層級,最後是管線層級。
假設一個情境,你建立客戶交易的分析管線。 與其編寫數百行包含檢查點管理和狀態處理的結構化串流程式碼,不如宣告資料表,讓架構來管理複雜性。
定義用於資料擷取的串流資料表
串流資料表是一個針對僅限附加的資料處理而最佳化的 Delta 資料表。 每個輸入記錄 僅處理一次,使串流資料表非常適合從雲端儲存或訊息佇列等來源擷取資料。
以下 SQL 範例建立一個串流表,透過 Auto Loader 從雲端儲存擷取 JSON 檔案:
CREATE OR REFRESH STREAMING TABLE customers_bronze
AS SELECT * FROM STREAM read_files(
"/Volumes/raw_data/customers",
format => "json"
);
關鍵字 STREAM 告訴管線將來源視為 串流資料集。 每次管線更新時,只有新檔案會被處理並附加到資料表中。
你也可以用 Python 達到同樣的效果:
from pyspark import pipelines as dp
@dp.table
def customers_bronze():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.inferColumnTypes", "true")
.load("/Volumes/raw_data/customers")
)
裝飾器 @dp.table 將此函式標記為串流表定義。 函式回傳一個串流資料框架,框架會隨著每次更新逐步處理。
建立具體化的轉換視圖
雖然串流資料表可以處理僅限附加的擷取,但具體化檢視更適合需要對可能變更的聚合或聯結維度進行合理的轉換。 實體化檢視會快取查詢結果,並自動與 上游變更同步。
以下範例建立一個實體化的視圖,連結客戶資料與交易資料:
CREATE OR REPLACE MATERIALIZED VIEW customer_transactions
AS SELECT
c.customer_id,
c.customer_name,
c.region,
t.transaction_date,
t.amount
FROM customers c
INNER JOIN transactions t ON c.customer_id = t.customer_id;
與串流資料表不同,實體化檢視處理包括聚合與連接等複雜查詢。 當上游資料表變動時,框架會決定更新實體化視圖的最有效方式——通常只處理受影響的資料,而非重新計算所有資料。
在 Python 中,您會使用 @dp.materialized_view 裝飾項目:
from pyspark import pipelines as dp
@dp.materialized_view
def regional_sales():
customers_df = spark.read.table("customers")
transactions_df = spark.read.table("transactions")
return (
customers_df.join(transactions_df, on="customer_id", how="inner")
.groupBy("region")
.agg({"amount": "sum"})
)
小提示
當您的來源資料僅限附加且您需要低延遲處理時,請使用串流資料表。 當你需要聚合、複雜連接,或必須處理原始資料的更新與刪除時,請使用實體化檢視。
套用資料品質標準
生產流程需要確保資料品質。 期望 是宣告式的限制,用來在記錄通過管道時進行驗證。 你定義什麼是有效數據,框架會追蹤指標並在違規發生時採取行動。
當紀錄驗證失敗時,框架會根據你指定的動作,以不同方式回應:
| 動作 | 行為 |
|---|---|
EXPECT (預設值) |
保留無效紀錄,並以指標追蹤違規次數 |
ON VIOLATION DROP ROW |
從輸出中移除無效紀錄 |
ON VIOLATION FAIL UPDATE |
停止目前流程;其他流程繼續執行。 在重新處理之前需要人工干預。 |
用於 ON VIOLATION FAIL UPDATE 關鍵約束條件,若任何違規顯示存在嚴重資料問題,需調查後才能繼續處理。
使用 Lakeflow 管線編輯器開發管線
Lakeflow 管線編輯器提供整合開發環境,用於建立與測試管線。 當你在 Azure Databricks 建立新的 ETL 管線時,編輯器會提供預設的資料夾結構,並分別存放轉換原始碼、探索筆記本和工具模組的目錄。
備註
本單元的程式碼範例為 管線定義,而非互動筆記本程式碼。 你無法直接在一般筆記本中執行這些語句。
pyspark.pipelines模組與串流表的 DDL 語句僅在 Lakeflow 管線執行時中提供。 若要執行您的管線程式碼,請使用試執行功能進行驗證,然後透過作業&管線介面執行管線。
編輯器透過多項功能支援迭代開發:
Genie Code:用自然語言描述您的管線,讓這個代理體驗為你創建、更新與除錯管線程式碼——從資料發現、程式碼產生到管線執行及解決資料品質問題。
試跑:在不處理資料的情況下驗證你的管線程式碼,讓你在執行前發現語法錯誤和缺少的相依性。
選擇性執行:執行單一檔案或單一資料表定義,而非整個管線,促進開發過程中更快的迭代。
互動式 DAG:視覺化資料表間的相依關係圖,選擇多個資料表進行目標刷新,並檢視執行指標。
資料預覽:直接在編輯器中取樣串流資料表與實體化視圖,以驗證轉換邏輯。
透過宣告式方法與整合工具,您可以建立可維護、可觀察且可以確保資料品質一致性的生產級資料管線。 這個框架處理增量處理與協調的複雜性,讓你能專注於定義對組織重要的業務邏輯。