你定义它?你就能用它!

PySpark 入门

Benjamin Schmidt

Data Engineer

用于可重复任务的 UDF

UDF(用户自定义函数):用 PySpark DataFrame 处理数据的自定义函数

UDF 优势:

  • 复用常见任务
  • 可直接在 Spark 中注册并共享
  • PySpark DataFrame(较小数据集)
  • pandas UDF(较大数据集)
PySpark 入门

定义并注册 UDF

所有 PySpark UDF 都需通过 udf() 函数注册。

# 定义函数
def to_uppercase(s):
    return s.upper() if s else None

# 注册函数 to_uppercase_udf = udf(to_uppercase, StringType())
# 应用于 DataFrame df = df.withColumn("name_upper", to_uppercase_udf(df["name"]))
# 查看结果 df.show()

注意:UDF 允许在 PySpark DataFrame 上应用自定义 Python 逻辑

PySpark 入门

pandas UDF

  • 消除昂贵的代码与数据转换
  • 无需向 SparkSession 注册
  • 在超大数据集上利用 pandas 能力
from pyspark.sql.functions import pandas_udf

@pandas_udf("float")
def fahrenheit_to_celsius_pandas(temp_f):
    return (temp_f - 32) * 5.0/9.0
PySpark 入门

PySpark UDF 与 pandas UDF 对比

PySpark UDF

  • 适合相对较小的数据集
  • 简单转换,如数据清洗
  • 更侧重列级而非行级变更
  • 必须用 udf() 在 Spark Session 中注册

pandas UDF

  • 适合相对较大的数据集
  • 复杂操作,超出简单清洗
  • 更支持行级的特定变更
  • 可在 Spark Session 外调用
PySpark 入门

开始练习吧!

PySpark 入门

Preparing Video For Download...