Zlepšení výkonu

Cleaning Data with PySpark

Mike Metzger

Data Engineering Consultant

Vysvětlení plánu provádění v Sparku

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>
Cleaning Data with PySpark

Co je shuffling?

Shuffling označuje přesun dat mezi pracovními uzly za účelem dokončení úlohy

  • Skrývá složitost před uživatelem
  • Může být pomalý
  • Snižuje celkovou propustnost
  • Je často nutný, ale je třeba ho minimalizovat
Cleaning Data with PySpark

Jak omezit shuffling?

  • Omezte použití .repartition(num_partitions)
    • Použijte místo toho .coalesce(num_partitions)
  • Buďte opatrní při volání .join()
  • Používejte .broadcast()
  • Omezení nemusí být vždy nutné
Cleaning Data with PySpark

Broadcasting

Broadcasting:

  • Poskytuje kopii objektu každému pracovnímu uzlu
  • Zabraňuje nadměrné komunikaci mezi uzly
  • Může výrazně zrychlit operace .join()

Použijte metodu .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Cleaning Data with PySpark

Pojďme si procvičit!

Cleaning Data with PySpark

Preparing Video For Download...