用户自定义函数

使用 PySpark 进行数据清洗

Mike Metzger

Data Engineering Consultant

定义...

"用户自定义函数"(UDF)

  • Python 方法
  • 通过 pyspark.sql.functions.udf 封装
  • 存为变量
  • 像普通 Spark 函数调用
使用 PySpark 进行数据清洗

字符串反转 UDF

定义一个 Python 方法

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

封装函数并存为变量

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 进行数据清洗

Passons à la pratique !

使用 PySpark 进行数据清洗

Preparing Video For Download...