Peningkatan kinerja

Membersihkan Data dengan PySpark

Mike Metzger

Data Engineering Consultant

Menjelaskan rencana eksekusi 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>
Membersihkan Data dengan PySpark

Apa itu shuffling?

Shuffling adalah memindahkan data ke berbagai worker untuk menyelesaikan tugas

  • Menyembunyikan kerumitan dari pengguna
  • Dapat lambat diselesaikan
  • Menurunkan throughput keseluruhan
  • Sering perlu, tetapi usahakan meminimalkannya
Membersihkan Data dengan PySpark

Cara membatasi shuffling

  • Batasi penggunaan .repartition(num_partitions)
    • Gunakan .coalesce(num_partitions) sebagai gantinya
  • Berhati-hati saat memanggil .join()
  • Gunakan .broadcast()
  • Mungkin tidak perlu membatasinya
Membersihkan Data dengan PySpark

Broadcasting

Broadcasting:

  • Memberikan salinan sebuah objek ke tiap worker
  • Mencegah komunikasi antarnode yang tidak perlu/berlebihan
  • Dapat sangat mempercepat operasi .join()

Gunakan metode .broadcast(<DataFrame>)

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

Ayo berlatih!

Membersihkan Data dengan PySpark

Preparing Video For Download...