建立欄位
在本章中,你將學會如何使用 Spark 的 DataFrame 類別所提供的方法來執行常見的資料操作。
先來看看如何做欄位層級的運算。在 Spark 中,你可以使用 .withColumn() 方法。它需要兩個引數:第一個是新欄位的名稱(字串),第二個是新欄位本身。
新欄位必須是 Column 類別的物件。建立它很容易,只要用 df.colName 從你的 DataFrame 取出一個欄位即可。
更新 Spark DataFrame 和在 pandas 裡操作有些不同,因為 Spark 的 DataFrame 是「不可變」(immutable)的。也就是說它不能直接被修改,因此不能就地更新欄位。
所以,這些方法都會回傳一個新的 DataFrame。若要覆寫原本的 DataFrame,你必須像下面這樣把回傳的 DataFrame 指派回去:
df = df.withColumn("newCol", df.oldCol + 1)
以上程式碼會建立一個和 df 具有相同欄位、外加一個新欄位 newCol 的 DataFrame,其中每個值都等於對應的 oldCol 值再加 1。
若要覆寫既有欄位,只要把該欄位名稱當作第一個引數傳入即可!
請記得,名為 spark 的 SparkSession 已經在你的工作環境中。
本練習屬於課程
PySpark 基礎
練習說明
- 使用
spark.table(),引數為"flights",從.catalog建立一個包含flights資料表內容的 DataFrame,並將其存為flights。 - 使用
flights.show()顯示flights的前幾列。請檢查輸出:air_time欄位包含航班的飛行時間(分鐘)。 - 更新
flights,新增名為duration_hrs的欄位,內容為每段航班的飛行時間(小時)。你需要用每小時的分鐘數去除以air_time。
動手互動練習
試著完成這個範例程式碼,體驗一下這個練習。
# Create the DataFrame flights
flights = spark.table(____)
# Show the head
____.____()
# Add duration_hrs
flights = flights.withColumn(____)