自定就能用,定義就能用!

PySpark 入門

Benjamin Schmidt

Data Engineer

用 UDF 讓任務可重複

UDFUser-Defined Function):在 PySpark DataFrame 上處理資料的自訂函式

UDF 優點:

  • 重複利用常見任務
  • 可直接向 Spark 註冊並分享
  • PySpark DataFrame(較小資料集)
  • pandas UDF(較大資料集)
PySpark 入門

定義與註冊 UDF

所有 PySpark UDF 都需要透過 udf() 註冊。

# Define the function
def to_uppercase(s):
    return s.upper() if s else None

# Register the function to_uppercase_udf = udf(to_uppercase, StringType())
# Apply the UDF to the DataFrame df = df.withColumn("name_upper", to_uppercase_udf(df["name"]))
# See the results 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...