การใช้ Broadcasting ในการ Join ของ Spark
การ join ตารางใน Spark จะถูกแบ่งไปประมวลผลบน worker หลายตัวในคลัสเตอร์ หากข้อมูลไม่ได้อยู่ในเครื่องเดียวกัน Spark จำเป็นต้องทำ shuffle ซึ่งส่งผลเสียต่อประสิทธิภาพ ในแบบฝึกหัดนี้จะใช้ฟีเจอร์ broadcast ของ Spark เพื่อส่งสำเนาข้อมูลที่ระบุไปยัง ทุก node
เคล็ดลับสำคัญ:
- ควร broadcast DataFrame ที่มีขนาดเล็กกว่า เพราะยิ่ง DataFrame ใหญ่ ก็ยิ่งใช้เวลาในการส่งไปยัง worker node มากขึ้น
- สำหรับ DataFrame ขนาดเล็กมาก อาจไม่จำเป็นต้องใช้ broadcasting เพราะ Spark สามารถปรับแต่งการทำงานให้เองได้
- หากดูแผนการรันคิวรี การปรากฏของ broadcastHashJoin หมายความว่าตั้งค่า broadcasting สำเร็จแล้ว
DataFrame flights_df และ airports_df พร้อมใช้งานแล้ว
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
การทำความสะอาดข้อมูลด้วย PySpark
คำแนะนำการฝึกหัด
- Import เมธอด
broadcast()จากpyspark.sql.functions - สร้าง DataFrame ใหม่ชื่อ
broadcast_dfโดย joinflights_dfกับairports_dfโดยใช้ broadcasting - แสดงแผนการรันคิวรีและสังเกตความแตกต่างจากแบบเดิม
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
# 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.____()