使用者自訂函式

使用 PySpark 清理資料

Mike Metzger

Data Engineering Consultant

定義方式…

「使用者自訂函式」或 UDFs

  • Python 方法
  • 透過 pyspark.sql.functions.udf 包裝
  • 存成變數
  • 像一般 Spark 函式呼叫
使用 PySpark 清理資料

字串反轉 UDF

先定義一個 Python 方法

def reverseString(mystr):
    return mystr[::-1]

包裝成 UDF 並存為變數

udfReverseString = udf(reverseString, StringType())

在 Spark 中使用

user_df = user_df.withColumn('ReverseName', 
                 udfReverseString(user_df.Name))
使用 PySpark 清理資料

無參數範例

def sortingCap():
    return random.choice(['G', 'H', 'R', 'S'])
udfSortingCap = udf(sortingCap, StringType())
user_df = user_df.withColumn('Class', udfSortingCap())
Name Age Class
Alice 14 H
Bob 18 S
Candice 63 G
使用 PySpark 清理資料

一起來練習吧!

使用 PySpark 清理資料

Preparing Video For Download...