Agrégation et jointure efficaces des données

Transformation des données avec Spark SQL dans Databricks

Disha Mukherjee

Lead Data Engineer

Trois questions, un ensemble de données propre

recraft: half: Un ou une analyste de données examine des graphiques colorés et des rapports d'affaires sur plusieurs écrans, isolé sur fond transparent

 

$$

  • Quelles catégories génèrent le plus de revenus ?
  • Qui sont les meilleurs clients ?
  • Comment enrichir les résultats avec le contexte du service ?
Transformation des données avec Spark SQL dans Databricks

Fonctionnement de groupBy() et agg()

 

$$

  • groupBy() - partitionne les lignes par colonne clé
  • agg() - applique des fonctions à chaque groupe en parallèle
  • Les deux sont paresseuses - une action déclenche le calcul
  • Un appel à agg() = une lecture des données

recraft: half: Des rangées de données colorées divisées en trois groupes distincts étiquetés avec une valeur totale calculée sous chaque groupe, isolé sur fond transparent

Transformation des données avec Spark SQL dans Databricks

Revenus par catégorie - code

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())
)
Transformation des données avec Spark SQL dans Databricks

Revenus par catégorie - résultat

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       |
+-----------+-----------------+-----------------+---------------+
Transformation des données avec Spark SQL dans Databricks

Enrichir avec une table de dimensions

recraft: half: Deux tables de base de données reliées par un lien ou pont lumineux, isolé sur fond transparent

$$

+-----------+----------+
|Category   |Department|
+-----------+----------+
|Clothing   |Retail    |
|Dining     |Food      |
|Electronics|Tech      |
|Groceries  |Food      |
|Savings    |Finance   |
+-----------+----------+
Transformation des données avec Spark SQL dans Databricks

Jointure gauche 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          |
+-----------+-----------+----------+------------------+
Transformation des données avec Spark SQL dans Databricks

Le coût caché : le shuffle

recraft: half: Des paquets de données circulent à travers un réseau de serveurs et de nœuds connectés avec des flèches montrant le mouvement, isolé sur fond transparent

 

  • Les clés correspondantes doivent se trouver sur la même machine
  • Spark déplace des lignes sur le réseau pour les aligner
  • Shuffle = goulot d'étranglement à grande échelle
Transformation des données avec Spark SQL dans Databricks

Interface Spark (Spark UI)

recraft: full: Un tableau de bord sombre de surveillance analytique avec des diagrammes à barres horizontales colorées montrant des étapes de tâches, des barres de progression et des mesures de débit de données, isolé sur fond transparent

 

  • Jobs et Stages - suivent ce qui a tourné et la durée
  • Shuffle Read/Write - volume de données déplacées
  • Utile pour diagnostiquer les goulots d'étranglement d'un pipeline
Transformation des données avec Spark SQL dans Databricks

Lire le plan de requête avec .explain()

df_joined.explain(mode="formatted")
...
+- PhotonBroadcastHashJoin LeftOuter (18)
   :- ...
   +- PhotonShuffleExchangeSource (17)
      +- PhotonShuffleMapStage (16)
         +- PhotonShuffleExchangeSink (15)
            +- LocalTableScan (13)
...
  • .explain() - examiner le plan d'exécution
  • Nœuds 15 et 16 - aucun Argument de routage consigné
Transformation des données avec Spark SQL dans Databricks

Jointure par diffusion (broadcast) - la solution

# Encapsuler la petite table avec 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
Transformation des données avec Spark SQL dans Databricks

Avant/après - valider le plan

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 table de dimensions a des instructions de routage
  • La grande table df_valid ne bouge jamais
  • Un changement = des minutes aux secondes à grande échelle
Transformation des données avec Spark SQL dans Databricks

Passons à la pratique !

Transformation des données avec Spark SQL dans Databricks

Preparing Video For Download...