定義しましたか?使ってください!

PySpark 入門

Benjamin Schmidt

Data Engineer

繰り返し可能なタスクのための UDF

UDF (ユーザー定義関数): PySparkデータフレームを使用してデータを扱うためのカスタム関数

UDFの利点:

  • 般的なタスクを再利用して繰り返す
  • Sparkに直接登録され、共有できます
  • PySpark DataFrames(小規模データセット向け)
  • 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{{5}} にカスタムPythonロジックを適用できます

PySpark 入門

pandas UDF

  • コードとデータの高コストな変換を排除
  • Spark セッションに登録される必要がない
  • 非常に大規模なデータセットで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

  • 比較的小規模なデータセットに最適
  • データクリーニングのような簡単な変換
  • 変更は行レベルではなく、カラムレベルで発生します
  • SparkSession に登録する必要がありudf()

pandas UDF

  • 比較的大きなデータセット
  • 単純なデータクリーニングを超える複雑な操作
  • 列レベルの変更よりも行レベルの特定の変更
  • Spark セッションの外部から呼び出すことができる
PySpark 入門

練習しましょう!

PySpark 入門

Preparing Video For Download...