Визначаєте? Використовуєте!

Вступ до PySpark

Benjamin Schmidt

Data Engineer

UDF для повторюваних завдань

UDF (User-Defined Function): власна функція для роботи з даними у датафреймах PySpark

Переваги UDF:

  • Повторне використання типових завдань
  • Реєструються безпосередньо в Spark і можуть ділитися
  • PySpark DataFrames (для менших наборів даних)
  • pandas UDFs (для більших наборів даних)
Вступ до 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 дають змогу застосовувати власну Python-логіку до датафреймів PySpark

Вступ до 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

  • Найкраще для відносно малих наборів даних
  • Прості перетворення, як-от очищення даних
  • Зміни на рівні стовпців, не рядків
  • Потрібна реєстрація в Spark Session через udf()

pandas UDF

  • Відносно великі набори даних
  • Складніші операції, ніж просте очищення даних
  • Точкові зміни на рівні рядків, а не стовпців
  • Можна викликати поза Spark Session
Вступ до PySpark

Давайте потренуємось!

Вступ до PySpark

Preparing Video For Download...