Mengagregasi dan menggabungkan data secara efisien

Transformasi Data dengan Spark SQL di Databricks

Disha Mukherjee

Lead Data Engineer

Tiga pertanyaan, satu himpunan data bersih

recraft: half: Seorang analis data meninjau bagan berwarna dan laporan bisnis pada beberapa layar, terisolasi pada latar transparan

 

$$

  • Kategori mana yang mendorong pendapatan tertinggi?
  • Siapa pelanggan teratas?
  • Bagaimana kita perkaya hasil dengan konteks departemen?
Transformasi Data dengan Spark SQL di Databricks

Cara kerja groupBy() dan agg()

 

$$

  • groupBy() - membagi baris berdasarkan kolom kunci
  • agg() - menerapkan fungsi ke tiap grup secara paralel
  • Keduanya bersifat lazy - aksi memicu komputasi
  • Satu pemanggilan agg() = satu pemindaian data

recraft: half: Baris data berwarna dibagi menjadi tiga grup berlabel terpisah dengan nilai total dihitung di bawah tiap grup, terisolasi pada latar transparan

Transformasi Data dengan Spark SQL di Databricks

Pendapatan per kategori - kode

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())
)
Transformasi Data dengan Spark SQL di Databricks

Pendapatan per kategori - keluaran

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       |
+-----------+-----------------+-----------------+---------------+
Transformasi Data dengan Spark SQL di Databricks

Memperkaya dengan tabel dimensi

recraft: half: Dua tabel basis data dihubungkan oleh tautan atau jembatan bercahaya, terisolasi pada latar transparan

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Transformasi Data dengan Spark SQL di Databricks

Left join standar

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          |
+-----------+-----------+----------+------------------+
Transformasi Data dengan Spark SQL di Databricks

Biaya tersembunyi - shuffle

recraft: half: Paket data mengalir di jaringan server dan node yang terhubung dengan panah menunjukkan pergerakan, terisolasi pada latar transparan

 

  • Kunci yang cocok harus berada di mesin yang sama
  • Spark mengirim baris lewat jaringan untuk menyelaraskannya
  • Shuffle = hambatan saat skala besar
Transformasi Data dengan Spark SQL di Databricks

Spark UI

recraft: full: Dasbor pemantauan analitik gelap dengan bagan batang horizontal berwarna menampilkan tahap job, bilah progres, dan metrik throughput data, terisolasi pada latar transparan

 

  • Jobs dan Stages - lacak apa yang berjalan dan durasinya
  • Shuffle Read/Write - lihat berapa banyak data yang dipindah
  • Berguna untuk mendiagnosis bottleneck pipeline
Transformasi Data dengan Spark SQL di Databricks

Membaca rencana kueri dengan .explain()

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - inspeksi rencana eksekusi
  • Node 15 & 16 - no routing Arguments tercatat
Transformasi Data dengan Spark SQL di Databricks

Broadcast join - solusinya

# Bungkus tabel kecil dalam 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
Transformasi Data dengan Spark SQL di Databricks

Sebelum dan sesudah - verifikasi rencana

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]
  • Tabel dimensi memiliki instruksi routing
  • Tabel besar df_valid tidak pernah dipindah
  • Satu perubahan = menit menjadi detik pada skala besar
Transformasi Data dengan Spark SQL di Databricks

Ayo berlatih!

Transformasi Data dengan Spark SQL di Databricks

Preparing Video For Download...