Menggunakan broadcasting pada join di Spark
Ingat bahwa operasi join tabel di Spark dibagi ke antara worker kluster. Jika data tidak lokal, berbagai operasi shuffle diperlukan dan dapat berdampak negatif pada kinerja. Sebagai gantinya, kita akan menggunakan operasi broadcast milik Spark untuk memberikan salinan data yang ditentukan kepada setiap node.
Beberapa kiat:
- Broadcast DataFrame yang lebih kecil. Semakin besar suatu DataFrame, semakin lama waktu yang dibutuhkan untuk mentransfernya ke node worker.
- Pada DataFrame kecil, kadang lebih baik melewati broadcasting dan membiarkan Spark menentukan optimasinya sendiri.
- Jika Anda melihat rencana eksekusi kueri, adanya broadcastHashJoin menandakan Anda telah berhasil mengonfigurasi broadcasting.
DataFrame flights_df dan airports_df tersedia untuk Anda.
Latihan ini merupakan bagian dari kursus
Membersihkan Data dengan PySpark
Instruksi latihan
- Impor metode
broadcast()daripyspark.sql.functions. - Buat DataFrame baru
broadcast_dfdengan melakukan joinflights_dfdenganairports_df, menggunakan broadcasting. - Tampilkan rencana kueri dan pertimbangkan perbedaannya dengan yang asli.
Latihan interaktif langsung praktik
Cobalah latihan ini dengan melengkapi kode contoh ini.
# Import the broadcast method from pyspark.sql.functions
from ____ import ____
# Join the flights_df and airports_df DataFrames using broadcasting
broadcast_df = flights_df.____(____(airports_df), \
flights_df["Destination Airport"] == airports_df["IATA"] )
# Show the query plan and compare against the original
broadcast_df.____()