Что такое UDF в Spark?

«Что такое UDF в Spark?» — вопрос из категории Apache Spark, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

UDF (User-Defined Function) в Apache Spark — это пользовательская функция, которую можно зарегистрировать и использовать в SQL-запросах или DataFrame API для выполнения операций, не предусмотренных встроенными функциями Spark.

Зачем нужны UDF? Для реализации сложной, специфичной для предметной области логики, которую нельзя выразить стандартными функциями Spark (например, парсинг специализированного формата строки, сложные математические расчеты).

Типы UDF в Spark:

  1. Обычные 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()
  2. 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.