有效率地彙總與連接資料

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

Disha Mukherjee

Lead Data Engineer

三個問題,一份乾淨資料集

recraft: half: 一位資料分析師在多個螢幕上檢視彩色圖表與商務報告,透明背景獨立

 

$$

  • 哪些類別帶來最高營收?
  • 誰是頂尖客戶
  • 如何用部門脈絡豐富結果?
在 Databricks 使用 Spark SQL 進行資料轉換

groupBy() 與 agg() 的運作方式

 

$$

  • groupBy()-依鍵欄分割列
  • agg()-對每個群組並行套用函式
  • 兩者皆為惰性-需動作才會計算
  • 一次 agg() 呼叫=一次資料掃描

recraft: half: 彩色資料列被分成三個已標示的群組,每組下方計算總值,透明背景獨立

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

類別營收-程式碼

category_revenue = (
    df_valid
    .groupBy("Category")

.agg( F.round(F.sum("Transaction_Amount"), 2).alias("total_revenue"), F.count("Transaction_Amount").alias("transaction_count"), F.round(F.avg("Transaction_Amount"), 2).alias("avg_transaction"), )
.orderBy(F.col("total_revenue").desc())
)
在 Databricks 使用 Spark SQL 進行資料轉換

類別營收-輸出

category_revenue.show(truncate=False)
+-----------+-----------------+-----------------+---------------+
|Category   |total_revenue    |transaction_count|avg_transaction|
+-----------+-----------------+-----------------+---------------+
|Clothing   |338266605.24     |6742             |50173.04       |
|Dining     |337631718.64     |6719             |50250.29       |
|Electronics|331102843.72     |6614             |50060.91       |
|Savings    |329346609.07     |6596             |49931.26       |
|Groceries  |329003451.57     |6537             |50329.43       |
|Unknown    |880591.27        |15               |58706.08       |
+-----------+-----------------+-----------------+---------------+
在 Databricks 使用 Spark SQL 進行資料轉換

用維度表豐富資料

recraft: half: 兩個資料庫資料表以發光連結或橋樑相接,透明背景獨立

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
在 Databricks 使用 Spark SQL 進行資料轉換

標準 left join

df_joined = df_valid.join(df_dim, on="Category", how="left")

df_joined.select( "Customer_ID", "Category", "Department", "Transaction_Amount" ).show(5, truncate=False)
+-----------+-----------+----------+------------------+
|Customer_ID|Category   |Department|Transaction_Amount|
+-----------+-----------+----------+------------------+
|CUST003    |Electronics|Tech      |5752.36           |
|CUST009    |Clothing   |Retail    |28959.12          |
|CUST010    |Savings    |Finance   |72098.18          |
|CUST011    |Savings    |Finance   |49771.05          |
|CUST020    |Groceries  |Food      |69825.14          |
+-----------+-----------+----------+------------------+
在 Databricks 使用 Spark SQL 進行資料轉換

隱藏成本:shuffle

recraft: half: 數據封包穿越連接的伺服器與節點網路,箭頭顯示流向,透明背景獨立

 

  • 相同鍵必須落在同一台機器
  • Spark 會在網路間傳送列以對齊
  • Shuffle=大規模瓶頸
在 Databricks 使用 Spark SQL 進行資料轉換

Spark UI

recraft: full: 暗色系分析監控儀表板,含彩色水平長條圖顯示作業階段、進度條與資料吞吐量指標,透明背景獨立

 

  • JobsStages-追蹤執行內容與耗時
  • Shuffle Read/Write-查看資料移動量
  • 有助診斷管線瓶頸
在 Databricks 使用 Spark SQL 進行資料轉換

.explain() 讀取查詢規劃

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain()檢視執行規劃
  • 節點 15 與 16-未記錄路由參數
在 Databricks 使用 Spark SQL 進行資料轉換

Broadcast join-解法

# 將小表包在 F.broadcast()
df_broadcast = df_valid.join(
    F.broadcast(df_dim),
    on="Category",
    how="left"
)

print(f"Broadcast joined rows: {df_broadcast.count():,}")
Broadcast joined rows: 33,223
在 Databricks 使用 Spark SQL 進行資料轉換

前後對照-驗證規劃

df_broadcast.explain(mode="formatted")
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...            
(15) PhotonShuffleExchangeSink
Arguments: SinglePartition
(16) PhotonShuffleMapStage
Arguments: EXECUTOR_BROADCAST, [id=#11839]
  • 維度表已有路由指示
  • 大表 df_valid 不需搬移
  • 一處更動=大規模時由數分鐘變數秒
在 Databricks 使用 Spark SQL 進行資料轉換

一起來練習吧!

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

Preparing Video For Download...