Améliorations de performance

Nettoyer des données avec PySpark

Mike Metzger

Data Engineering Consultant

Expliquer le plan d'exécution de Spark

voter_df = df.select(df['VOTER NAME']).distinct()
voter_df.explain()
== Physical Plan ==
*(2) HashAggregate(keys=[VOTER NAME#15], functions=[])
+- Exchange hashpartitioning(VOTER NAME#15, 200)
   +- *(1) HashAggregate(keys=[VOTER NAME#15], functions=[])
      +- *(1) FileScan csv [VOTER NAME#15] Batched: false, Format: CSV, Location: 
      InMemoryFileIndex[file:/DallasCouncilVotes.csv.gz], 
      PartitionFilters: [], PushedFilters: [], 
      ReadSchema: struct<VOTER NAME:string>
Nettoyer des données avec PySpark

Qu'est-ce que le shuffling ?

Le « shuffling » consiste à déplacer des données entre divers travailleurs pour terminer une tâche

  • Cache la complexité à l'utilisateur
  • Peut être lent à terminer
  • Réduit le débit global
  • Souvent nécessaire, mais à minimiser
Nettoyer des données avec PySpark

Comment limiter le shuffling ?

  • Limiter l'usage de .repartition(num_partitions)
    • Utiliser plutôt .coalesce(num_partitions)
  • Être prudent avec .join()
  • Utiliser .broadcast()
  • Il n'est pas toujours nécessaire de le limiter
Nettoyer des données avec PySpark

Broadcasting

Le « broadcasting » :

  • Fournit une copie d'un objet à chaque travailleur
  • Évite des communications inutiles/excessives entre nœuds
  • Peut accélérer fortement les opérations .join()

Utiliser la méthode .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Nettoyer des données avec PySpark

Passons à la pratique !

Nettoyer des données avec PySpark

Preparing Video For Download...