DataFrame 欄位操作

使用 PySpark 清理資料

Mike Metzger

Data Engineering Consultant

DataFrame 複習

DataFrame:

  • 由列與欄組成
  • 不可變
  • 以各種轉換操作修改資料
# 傳回名字以 "M" 開頭的列 
voter_df.filter(voter_df.name.like('M%'))

# 只選取 name 與 position voters = voter_df.select('name', 'position')
使用 PySpark 清理資料

常見 DataFrame 轉換

  • 篩選 / Where
    voter_df.filter(voter_df.date > '1/1/2019') # 或 voter_df.where(...)
    
  • 選取(Select)
    voter_df.select(voter_df.name)
    
  • withColumn
    voter_df.withColumn('year', voter_df.date.year)
    
  • drop
    voter_df.drop('unused_column')
    
使用 PySpark 清理資料

篩選資料

  • 移除 null
  • 移除異常值
  • 拆分合併來源的資料
  • 使用 ~ 取反
    voter_df.filter(voter_df['name'].isNotNull())
    voter_df.filter(voter_df.date.year > 1800)
    voter_df.where(voter_df['_c0'].contains('VOTE'))
    voter_df.where(~ voter_df._c1.isNull())
    
使用 PySpark 清理資料

欄位字串轉換

  • 位於 pyspark.sql.functions
    import pyspark.sql.functions as F
    
  • 以欄為單位套用轉換
    voter_df.withColumn('upper', F.upper('name'))
    
  • 可建立中介欄
    voter_df.withColumn('splits', F.split('name', ' '))
    
  • 可轉型為其他型別
    voter_df.withColumn('year', voter_df['_c4'].cast(IntegerType()))
    
使用 PySpark 清理資料

ArrayType() 欄位函式

與 ArrayType() 互動的多種工具函式/轉換

.size(<column>) - 回傳 ArrayType() 欄位的長度

.getItem(<index>) - 取得清單欄位在指定索引的位置元素。

使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...