Aggregare e unire i dati in modo efficiente

Trasformazione dei dati con Spark SQL in Databricks

Disha Mukherjee

Lead Data Engineer

Tre domande, un dataset pulito

recraft: half: Un* data analyst che esamina grafici colorati e report aziendali su più schermi, isolato su sfondo trasparente

 

$$

  • Quali categorie generano più ricavi?
  • Chi sono i clienti top?
  • Come arricchire i risultati con il contesto dei reparti?
Trasformazione dei dati con Spark SQL in Databricks

Come funzionano groupBy() e agg()

 

$$

  • groupBy() - partiziona le righe per colonna chiave
  • agg() - applica funzioni a ogni gruppo in parallelo
  • Entrambe sono lazy - un'azione avvia il calcolo
  • Una agg() = una scansione dei dati

recraft: half: Righe di dati colorate divise in tre gruppi etichettati con un valore totale calcolato sotto ciascun gruppo, isolato su sfondo trasparente

Trasformazione dei dati con Spark SQL in Databricks

Ricavi per categoria - codice

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())
)
Trasformazione dei dati con Spark SQL in Databricks

Ricavi per categoria - output

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       |
+-----------+-----------------+-----------------+---------------+
Trasformazione dei dati con Spark SQL in Databricks

Arricchire con una tabella dimensione

recraft: half: Due tabelle di database collegate da un collegamento o ponte luminoso, isolate su sfondo trasparente

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Trasformazione dei dati con Spark SQL in Databricks

Left join standard

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          |
+-----------+-----------+----------+------------------+
Trasformazione dei dati con Spark SQL in Databricks

Il costo nascosto: shuffle

recraft: half: Pacchetti di dati che fluiscono in una rete di server e nodi con frecce che indicano il movimento, isolati su sfondo trasparente

 

  • Le chiavi corrispondenti devono finire sulla stessa macchina
  • Spark sposta righe in rete per allinearle
  • Shuffle = collo di bottiglia in scala
Trasformazione dei dati con Spark SQL in Databricks

Spark UI

recraft: full: Un cruscotto di monitoraggio analytics scuro con barre orizzontali colorate che mostrano fasi dei job, barre di avanzamento e metriche di throughput, isolato su sfondo trasparente

 

  • Jobs e Stages - cosa è stato eseguito e per quanto
  • Shuffle Read/Write - quanta data si è mossa
  • Utile per diagnosticare colli di bottiglia nelle pipeline
Trasformazione dei dati con Spark SQL in Databricks

Leggere il piano di query con .explain()

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - ispeziona il piano di esecuzione
  • Nodi 15 e 16 - nessun routing Arguments registrato
Trasformazione dei dati con Spark SQL in Databricks

Broadcast join - la soluzione

# Incapsula la tabella piccola 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
Trasformazione dei dati con Spark SQL in Databricks

Prima e dopo - verifica del piano

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]
  • La tabella dim ha istruzioni di instradamento
  • La tabella grande df_valid non si muove mai
  • Una modifica = da minuti a secondi su larga scala
Trasformazione dei dati con Spark SQL in Databricks

Esercitiamoci!

Trasformazione dei dati con Spark SQL in Databricks

Preparing Video For Download...