Îmbunătățiri de performanță

Curățarea datelor cu PySpark

Mike Metzger

Data Engineering Consultant

Planul de execuție 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>
Curățarea datelor cu PySpark

Ce este shuffling-ul?

Shuffling înseamnă redistribuirea datelor între noduri pentru a finaliza o sarcină

  • Ascunde complexitatea față de utilizator
  • Poate fi lent
  • Reduce randamentul general
  • Este adesea necesar, dar trebuie minimizat
Curățarea datelor cu PySpark

Cum se limitează shuffling-ul?

  • Limitați utilizarea .repartition(num_partitions)
    • Folosiți .coalesce(num_partitions) în schimb
  • Atenție la apelurile .join()
  • Utilizați .broadcast()
  • Nu întotdeauna este necesar să îl limitați
Curățarea datelor cu PySpark

Broadcasting

Broadcasting:

  • Furnizează o copie a unui obiect fiecărui nod
  • Previne comunicarea excesivă între noduri
  • Poate accelera semnificativ operațiile .join()

Utilizați metoda .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Curățarea datelor cu PySpark

Să exersăm!

Curățarea datelor cu PySpark

Preparing Video For Download...