การเขียน Spark configurations
หลังจากที่ได้ตรวจสอบ Spark configurations บน cluster แล้ว ขั้นตอนต่อไปคือการปรับแต่งค่าบางอย่างเพื่อให้ Spark ทำงานได้ตามต้องการ จากนั้นจะนำเข้าข้อมูลเพื่อยืนยันว่าการเปลี่ยนแปลงมีผลกับ cluster จริง
ในตอนแรก Spark configuration ถูกตั้งไว้ที่ค่าเริ่มต้น 200 partitions
ออบเจกต์ spark พร้อมใช้งานแล้ว ไฟล์ชื่อ departures.txt.gz พร้อมสำหรับการนำเข้า และ DataFrame เริ่มต้นที่มีแถวที่ไม่ซ้ำกันจาก departures.txt.gz ถูกเก็บไว้ในตัวแปร departures_df
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
การทำความสะอาดข้อมูลด้วย PySpark
คำแนะนำการฝึกหัด
- เก็บจำนวน partitions ใน
departures_dfไว้ในตัวแปรbefore - เปลี่ยนค่า configuration
spark.sql.shuffle.partitionsเป็น 500 partitions - สร้าง
departures_dfขึ้นใหม่โดยอ่านแถวที่ไม่ซ้ำกันจากไฟล์ departures - แสดงจำนวน partitions ก่อนและหลังการเปลี่ยนแปลง configuration
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
# Store the number of partitions in variable
before = departures_df.____
# Configure Spark to use 500 partitions
____('spark.sql.shuffle.partitions', ____)
# Recreate the DataFrame using the departures data file
departures_df = spark.read.csv('departures.txt.gz').____
# Print the number of partitions for each instance
print("Partition count before change: %d" % ____)
print("Partition count after change: %d" % ____)