Ответ
UDF (User-Defined Function) в Apache Spark — это пользовательская функция, которую можно зарегистрировать и использовать в SQL-запросах или DataFrame API для выполнения операций, не предусмотренных встроенными функциями Spark.
Зачем нужны UDF? Для реализации сложной, специфичной для предметной области логики, которую нельзя выразить стандартными функциями Spark (например, парсинг специализированного формата строки, сложные математические расчеты).
Типы UDF в Spark:
-
Обычные UDF (Scalar UDF): Принимают одну строку входных данных и возвращают одно значение.
from pyspark.sql.functions import udf from pyspark.sql.types import StringType # Определяем функцию на Python def reverse_string(s): return s[::-1] # Создаем и регистрируем UDF reverse_udf = udf(reverse_string, StringType()) # Использование в DataFrame API df.withColumn('reversed_name', reverse_udf(df['name'])).show() # Регистрация для использования в Spark SQL spark.udf.register("sql_reverse", reverse_string, StringType()) spark.sql("SELECT sql_reverse(name) FROM users").show() -
Pandas UDF (Vectorized UDF): Используют Apache Arrow для пакетной передачи данных между JVM и Python, что значительно повышает производительность по сравнению с обычными UDF за счет векторных операций.
import pandas as pd from pyspark.sql.functions import pandas_udf from pyspark.sql.types import DoubleType @pandas_udf(DoubleType()) def multiply_by_two(series: pd.Series) -> pd.Series: # Операция над целым столбцом (Series) Pandas return series * 2 df.withColumn('doubled_value', multiply_by_two(df['value'])).show()
Важные ограничения:
- Производительность: Обычные (невекторизованные) Python UDF выполняются медленнее, чем встроенные функции Spark (написанные на Scala/JVM), так как данные сериализуются, передаются в Python-процесс и десериализуются.
- Оптимизация: Catalyst Optimizer не может "заглянуть" внутрь UDF, поэтому некоторые оптимизации (предикат pushdown) могут не работать.
Рекомендация: Всегда сначала ищите решение через встроенные функции Spark. Используйте Pandas UDF вместо обычных, где это возможно. Для сложной логики, критичной к производительности, рассмотрите написание собственных агрегаторов (Aggregator) или использование Scala.