開始使用免費開始

建立欄位

在本章中,你將學會如何使用 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。

若要覆寫既有欄位,只要把該欄位名稱當作第一個引數傳入即可!

請記得,名為 sparkSparkSession 已經在你的工作環境中。

本練習屬於課程

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(____)
編輯並執行程式碼