Prestandaförbättringar

Datarensning med PySpark

Mike Metzger

Data Engineering Consultant

Förklara Sparks exekveringsplan

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>
Datarensning med PySpark

Vad är shuffling?

Shuffling innebär att data flyttas mellan olika workers för att slutföra en uppgift

  • Döljer komplexitet för användaren
  • Kan vara långsamt
  • Sänker det totala genomflödet
  • Är ofta nödvändigt, men försök att minimera det
Datarensning med PySpark

Hur begränsar man shuffling?

  • Begränsa användningen av .repartition(num_partitions)
    • Använd .coalesce(num_partitions) i stället
  • Var försiktig med .join()
  • Använd .broadcast()
  • Kanske inte behöver begränsas
Datarensning med PySpark

Broadcasting

Broadcasting:

  • Skickar en kopia av ett objekt till varje worker
  • Minskar onödig kommunikation mellan noder
  • Kan kraftigt påskynda .join()-operationer

Använd metoden .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Datarensning med PySpark

Nu kör vi en övning!

Datarensning med PySpark

Preparing Video For Download...