Pipeline อย่างรวดเร็ว
ก่อนที่จะประมวลผลข้อมูลที่ซับซ้อนขึ้น ผู้จัดการต้องการดูตัวอย่าง pipeline อย่างง่ายที่ครอบคลุมขั้นตอนพื้นฐาน ในตัวอย่างนี้ จะต้องนำเข้าไฟล์ข้อมูล กรองบางแถว เพิ่มคอลัมน์ ID แล้วส่งออกเป็นข้อมูล JSON
มีการกำหนด context ของ spark ไว้แล้ว พร้อมกับไลบรารี pyspark.sql.functions ที่ตั้งชื่อแทนเป็น F ตามรูปแบบที่นิยมใช้กัน
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
การทำความสะอาดข้อมูลด้วย PySpark
คำแนะนำการฝึกหัด
- นำเข้าไฟล์
2015-departures.csv.gzไปยัง DataFrame โดย header ถูกกำหนดไว้แล้ว - กรอง DataFrame ให้เหลือเฉพาะเที่ยวบินที่มีระยะเวลามากกว่า 0 นาที ให้ใช้ index ของคอลัมน์ ไม่ใช่ชื่อคอลัมน์ (ใช้
.printSchema()เพื่อดูชื่อและลำดับของคอลัมน์) - เพิ่มคอลัมน์ ID
- บันทึกไฟล์เป็นเอกสาร JSON ชื่อ
output.json
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
# Import the data to a DataFrame
departures_df = spark.____(____, header=____)
# Remove any duration of 0
departures_df = departures_df.____(____)
# Add an ID column
departures_df = departures_df.____('id', ____)
# Write the file out to JSON format
____.write.____(____, mode='overwrite')