以 workflows 建立生產級管線

在 Databricks 使用 Spark SQL 進行資料轉換

Disha Mukherjee

Lead Data Engineer

為何選擇 Delta Lake?

一座安全發光的數位保險庫,內有整理好的資料表與防護盾,扁平現代風格

 

$$

  • ACID 交易 → 回滾失敗寫入
  • 綱要強制 → 阻擋型別不符
  • 版本控管 → 查詢任一過往狀態
在 Databricks 使用 Spark SQL 進行資料轉換

寫入 Delta

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • 一個新的 Delta 資料表會出現在 Unity Catalog
在 Databricks 使用 Spark SQL 進行資料轉換

Notebook 任務

task1_ingest:載入並清理

  • 載入 CSV → 進行清理 → 寫成 Delta 表
在 Databricks 使用 Spark SQL 進行資料轉換

Notebook 任務

task2_metrics:依類別的營收

  • 讀取已清理表 → 計算指標 → 寫入新表
在 Databricks 使用 Spark SQL 進行資料轉換

Notebook 任務

task3_customers:依消費額排名

  • 讀取已清理表 → 排名客戶 → 存成新表
在 Databricks 使用 Spark SQL 進行資料轉換

建立作業

Jobs 與 Pipelines 介面,顯示完成的三步任務 DAG,任務依賴從 task1_ingest 到 task2_metrics 到 task3_customers

  • 可手動執行、排程作業並設定觸發條件
在 Databricks 使用 Spark SQL 進行資料轉換

執行作業

作業執行 DAG:task1_ingest 與 task2_metrics 成功為綠,task3_customers 失敗為紅

在 Databricks 使用 Spark SQL 進行資料轉換

執行作業

檢查第三個任務的錯誤

  • 錯誤:customer-id 應為 Customer_ID
在 Databricks 使用 Spark SQL 進行資料轉換

執行作業

成功的圖形檢視

時間軸檢視

在 Databricks 使用 Spark SQL 進行資料轉換

什麼是 Lakeflow?

$$

$$

比較:Imperative(Jobs)| Declarative(Lakeflow)

 

  • Jobs → 逐步由我們管理
  • Lakeflow → 宣告表格應包含的內容
  • Databricks 負責順序、重試與運算資源
在 Databricks 使用 Spark SQL 進行資料轉換

@dlt.table 模式

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
在 Databricks 使用 Spark SQL 進行資料轉換

管線執行

Lakeflow 管線 DAG:transactions_bronze 100K 筆,transactions_silver 33K 筆,category_revenue_gold 6 筆,皆為綠色實體化檢視

  • Bronze(100K 筆)→ Silver(33K 筆)→ Gold(6 筆)
在 Databricks 使用 Spark SQL 進行資料轉換

該用 Notebooks、Jobs,或 Lakeflow?

$$

層級:Notebooks、Databricks Jobs、Lakeflow Pipelines

 

$$

  • Notebooks → 探索、雛形製作
  • Databricks Jobs → 多步驟、可排程的管線
  • Lakeflow → 完全受管、宣告式
在 Databricks 使用 Spark SQL 進行資料轉換

一起來練習吧!

在 Databricks 使用 Spark SQL 進行資料轉換

Preparing Video For Download...