SQL และ Parquet
ไฟล์ Parquet เหมาะอย่างยิ่งสำหรับใช้เป็นแหล่งข้อมูลสำหรับคิวรี SQL ใน Spark แม้จะสามารถรันคิวรีแบบเดียวกันผ่านฟังก์ชัน Python ของ Spark ได้โดยตรง แต่บางครั้งการใช้คิวรี SQL ควบคู่กับ Python ก็ทำได้ง่ายกว่า
ในแบบฝึกหัดนี้ เราจะอ่านไฟล์ Parquet ที่สร้างไว้ในแบบฝึกหัดที่แล้ว แล้วลงทะเบียนเป็นตาราง SQL จากนั้นจะรันคิวรีทดสอบกับตารางนั้น (ซึ่งก็คือไฟล์ Parquet นั่นเอง)
ออบเจกต์ spark และไฟล์ AA_DFW_ALL.parquet พร้อมใช้งานให้แล้วโดยอัตโนมัติ
แบบฝึกหัดนี้เป็นส่วนหนึ่งของหลักสูตร
การทำความสะอาดข้อมูลด้วย PySpark
คำแนะนำการฝึกหัด
- นำเข้าไฟล์
AA_DFW_ALL.parquetไปยังflights_df - ใช้เมธอด
createOrReplaceTempViewเพื่อตั้งชื่อตารางเป็นflights - รันคิวรี Spark SQL กับตาราง
flights
แบบฝึกหัดเชิงโต้ตอบแบบลงมือทำ
ลองทำแบบฝึกหัดนี้โดยเติมโค้ดตัวอย่างนี้ให้สมบูรณ์
# Read the Parquet file into flights_df
flights_df = spark.read.____(____)
# Register the temp table
flights_df.____('flights')
# Run a SQL query of the average flight duration
avg_duration = spark.____('SELECT avg(flight_duration) from flights').collect()[0]
print('The average flight time is: %d' % avg_duration)