เริ่มต้นใช้งานเริ่มต้นใช้งานได้ฟรี

การใช้ 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 โดย join flights_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.____()
แก้ไขและรันโค้ด