Poprawa wydajności

Czyszczenie danych w PySpark

Mike Metzger

Data Engineering Consultant

Analiza planu wykonania 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>
Czyszczenie danych w PySpark

Czym jest shuffling?

Shuffling to przenoszenie danych między węzłami w celu wykonania zadania

  • Ukrywa złożoność przed użytkownikiem
  • Może być powolny
  • Obniża ogólną przepustowość
  • Często konieczny, ale należy go minimalizować
Czyszczenie danych w PySpark

Jak ograniczyć shuffling?

  • Ogranicz użycie .repartition(num_partitions)
    • Zamiast tego użyj .coalesce(num_partitions)
  • Ostrożnie stosuj .join()
  • Używaj .broadcast()
  • Nie zawsze trzeba go ograniczać
Czyszczenie danych w PySpark

Broadcasting

Broadcasting:

  • Dostarcza kopię obiektu do każdego węzła
  • Ogranicza zbędną komunikację między węzłami
  • Może znacznie przyspieszyć operacje .join()

Użyj metody .broadcast(<DataFrame>)

from pyspark.sql.functions import broadcast
combined_df = df_1.join(broadcast(df_2))
Czyszczenie danych w PySpark

Czas na ćwiczenia!

Czyszczenie danych w PySpark

Preparing Video For Download...