データの集計と結合を効率化する

Databricks における Spark SQL を使ったデータ変換

Disha Mukherjee

Lead Data Engineer

3つの問い、1つのデータセット

recraft: half: カラフルなグラフとビジネスレポートを複数の画面で確認するデータアナリスト、透明な背景に配置

 

$$

  • どのカテゴリが最も売上に貢献しているか?
  • 上位顧客は誰か?
  • 部門情報で結果を補完するには?
Databricks における Spark SQL を使ったデータ変換

`groupBy()` と `agg()` の仕組み

 

$$

  • groupBy() - キー列で行を分割
  • agg() - 各グループに並列で関数を適用
  • どちらも遅延評価 - アクションで実行
  • agg() 1回 = データスキャン1回

recraft: half: カラフルなデータ行が3つのグループに分割され、各グループの合計値が下に表示されている、透明な背景に配置

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: 2つのデータベーステーブルが光るリンクで接続されている、透明な背景に配置

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Databricks における Spark SQL を使ったデータ変換

標準の左結合

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 を使ったデータ変換

隠れたコスト - シャッフル

recraft: half: データパケットが矢印でサーバーとノードのネットワーク上を流れている、透明な背景に配置

 

  • 一致するキーは同じマシンに集める必要がある
  • Spark は行をネットワーク越しに転送して整合させる
  • シャッフル = 大規模時のボトルネック
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 を使ったデータ変換

ブロードキャスト結合 - 解決策

# Wrap the small table in 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転送されない
  • 1つの変更で大規模処理が数分から数秒に
Databricks における Spark SQL を使ったデータ変換

練習しましょう!

Databricks における Spark SQL を使ったデータ変換

Preparing Video For Download...