Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 05.03 · 14 мин
Продвинутый
Pandas UDFArrowVectorized UDFpandas_udfSeries to Series

Pandas UDF: Arrow-based векторизация

Проблема, которую решает Pandas UDF

Как мы выяснили в предыдущем уроке, Python UDF выполняет сериализацию на каждую строку. При 1 миллиарде строк это 1 миллиард циклов serialize -> socket -> deserialize.

Pandas UDF решают эту проблему через батчевую обработку: вместо строк данные передаются пакетами (batches) по тысячам строк за раз, используя Apache Arrow columnar format.

Per-Row (Python UDF) vs Per-Batch (Pandas UDF, Arrow)
Python UDF (per-row)
JVM
──[row 1]──→
Python
──[result 1]──→
JVM
JVM
──[row 2]──→
Python
──[result 2]──→
JVM
JVM
──[row 3]──→
Python
──[result 3]──→
JVM
... (1 миллиард socket-вызовов)
Pandas UDF (per-batch, Arrow)
JVM
──[batch 10,000 rows, Arrow]──→
Python
(1 socket-вызов, columnar, zero-copy)
Python
──[batch 10,000 results, Arrow]──→
JVM
... (100,000 socket-вызовов вместо 1 миллиарда)

Результат: Pandas UDF в 5-20x быстрее обычного Python UDF.

Apache Arrow: ключ к производительности

Apache Arrow — это in-memory columnar format, который позволяет:

  1. Columnar layout: данные одного столбца хранятся непрерывно в памяти, что оптимально для vectorized операций
  2. Zero-copy reads: Python может читать Arrow-данные без копирования — JVM и Python разделяют один участок памяти
  3. Batch transfer: тысячи строк передаются одним блоком вместо per-row сериализации
Per-row (cloudpickle) vs Per-batch (Arrow)
Per-row (cloudpickle)
Row 1: {a:1, b:2}→ serialize
Row 2: {a:2, b:4}→ serialize
Row 3: {a:3, b:6}→ serialize
...
N serializations
Per-batch (Arrow)
Column A: [1, 2, 3...]
Column B: [2, 4, 6...]
1 serialize call for entire batch

Типы Pandas UDF

1. Series to Series (Scalar)

Самый частый тип. Принимает pd.Series (столбец), возвращает pd.Series того же размера:

import pandas as pd

from pyspark.sql.functions import pandas_udf

@pandas_udf("double")
def normalize_zscore(salary: pd.Series) -> pd.Series:
    """Z-score нормализация зарплаты"""
    return (salary - salary.mean()) / salary.std()

df.withColumn("salary_z", normalize_zscore(col("salary")))

Когда использовать: поэлементные трансформации, математические вычисления, строковые операции с библиотеками Python (regex, transliterate).

2. Iterator of Series to Iterator of Series

Вариант для stateful обработки или когда нужно загрузить модель один раз:

from typing import Iterator

@pandas_udf("double")
def predict_churn(batches: Iterator[pd.Series]) -> Iterator[pd.Series]:
    """ML-инференс: загружаем модель один раз, применяем к каждому батчу"""
    import joblib
    model = joblib.load("/models/churn_model.pkl")  # загрузка один раз!

    for batch in batches:
        features = batch.values.reshape(-1, 1)
        predictions = model.predict(features)
        yield pd.Series(predictions)

df.withColumn("churn_score", predict_churn(col("total_spent")))

Когда использовать: ML-инференс (загрузка модели один раз), соединение с внешним сервисом, любая stateful инициализация.

3. Grouped Map (applyInPandas)

Самый мощный тип. Получает весь DataFrame группы, возвращает DataFrame (может изменить схему):

def fill_missing_by_department(df: pd.DataFrame) -> pd.DataFrame:
    """Заполняем пропуски средним по отделу"""
    df["salary"] = df["salary"].fillna(df["salary"].mean())
    df["bonus"] = df["bonus"].fillna(0)
    return df

result = employees.groupBy("department").applyInPandas(
    fill_missing_by_department,
    schema="department string, name string, salary double, bonus double"
)

Когда использовать: сложные pandas-операции над группами, custom aggregations, pivot/reshape внутри группы, оконные функции, недоступные в Spark.

TIP

Выбор типа Pandas UDF:

  • Нужно трансформировать столбец → Series to Series
  • Нужно загрузить модель/ресурс один раз → Iterator of Series
  • Нужно обработать группу как pandas DataFrame → Grouped Map (applyInPandas)

4. Grouped Aggregate

Кастомные агрегации на уровне группы:

@pandas_udf("double")
def weighted_avg(values: pd.Series, weights: pd.Series) -> float:
    """Взвешенное среднее"""
    return (values * weights).sum() / weights.sum()

df.groupBy("department").agg(
    weighted_avg(col("salary"), col("experience_years")).alias("weighted_salary")
)

Сравнение производительности: Python UDF vs Pandas UDF

# Тестовый датасет: 50 миллионов строк
df = spark.range(50_000_000).withColumn("value", (rand() * 1000).cast("double"))

# Python UDF: ~6 минут
@udf(returnType=DoubleType())
def py_sqrt(x):
    import math
    return math.sqrt(x) if x and x >= 0 else None

# Pandas UDF: ~35 секунд (10x быстрее)
@pandas_udf("double")
def pd_sqrt(x: pd.Series) -> pd.Series:
    import numpy as np
    return np.sqrt(x)

Pandas UDF быстрее по двум причинам:

  1. Arrow батчи вместо per-row сериализации (сокращает transfer overhead)
  2. NumPy vectorized операции вместо Python for-loop (сокращает compute overhead)
Spark4.0

Spark 4.0 улучшает Arrow-интеграцию для Pandas UDF: оптимизированный батчинг, уменьшенный memory overhead при передаче данных между JVM и Python, и улучшенная поддержка сложных типов (StructType, ArrayType) через Arrow.

Ограничения и подводные камни

1. Размер батча влияет на memory

Spark контролирует размер батча через spark.sql.execution.arrow.maxRecordsPerBatch (по умолчанию 10,000). Слишком большой батч может вызвать OOM в Python worker:

# Если каждая строка содержит 1KB данных:
# 10,000 строк × 1KB = 10MB per batch (OK)
# 1,000,000 строк × 1KB = 1GB per batch (OOM!)

spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "10000")

2. Schema должна быть точной

Pandas UDF требует точного указания выходной схемы. Несоответствие типов вызовет runtime error:

# ОШИБКА: функция возвращает int, но объявлен double
@pandas_udf("double")
def wrong_type(x: pd.Series) -> pd.Series:
    return x.astype(int)  # вернёт int, а ожидается double!

3. Grouped Map может вызвать data skew

applyInPandas собирает все данные группы в один executor. Если одна группа содержит 80% данных (например, город “Москва”), один executor получит несоразмерную нагрузку.

Анти-паттерн: Pandas UDF, обрабатывающий по одной строке

# ПЛОХО: Pandas UDF, который не использует векторизацию
@pandas_udf("string")
def bad_pandas_udf(names: pd.Series) -> pd.Series:
    results = []
    for name in names:    # Python for-loop!
        if name:
            results.append(name.strip().title())
        else:
            results.append(None)
    return pd.Series(results)
# ХОРОШО: Vectorized операция через pandas/str accessor
@pandas_udf("string")
def good_pandas_udf(names: pd.Series) -> pd.Series:
    return names.str.strip().str.title()  # vectorized!

Первый вариант получает данные батчем через Arrow (хорошо), но обрабатывает их Python for-loop (плохо). Второй вариант использует vectorized pandas operations, которые выполняются на C-уровне.

Проверка знанийKnowledge check
Почему Pandas UDF с Iterator of Series предпочтительнее обычного Series-to-Series для ML-инференса?
ОтветAnswer
Iterator of Series позволяет загрузить ML-модель ОДИН раз (в начале итерации), а затем применять её ко всем батчам данных. Обычный Series-to-Series вызывается для каждого батча заново, и если модель загружается внутри функции, она будет загружаться повторно для каждого батча (тысячи раз). С Iterator of Series: model = load_model() выполняется однократно, затем for batch in batches: yield model.predict(batch) переиспользует загруженную модель.
Проверка знанийKnowledge check
Что произойдёт, если Pandas UDF использует Python for-loop вместо vectorized pandas-операций для обработки данных в батче?
ОтветAnswer
Pandas UDF всё равно получит данные батчем через Arrow (а не per-row как Python UDF), что сокращает overhead сериализации. Однако Python for-loop внутри функции будет в 10-100x медленнее, чем vectorized pandas/NumPy операции. Вы получите benefit от Arrow transfer, но потеряете benefit от vectorized computation. Результат: быстрее Python UDF (из-за Arrow), но значительно медленнее правильно написанного Pandas UDF (из-за for-loop).

Polymorphic Python UDTFs (Spark 4.0)

Spark4.0

В дополнение к Pandas UDF, Spark 4.0 ввёл Python User-Defined Table Functions (UDTFs) — функции, которые принимают вход (skalary, columns или таблицу целиком) и возвращают таблицу из произвольного числа строк. Это принципиально новый instrument: UDF возвращает скаляр, UDTF возвращает rows.

Главная новинка 4.0 — polymorphic UDTFs, у которых выходная схема определяется не статической аннотацией, а вычисляется самой функцией в зависимости от входных аргументов.

Базовый UDTF: фиксированная схема

from pyspark.sql.functions import udtf

@udtf(returnType="word: string, length: int")
class TokenizeUDTF:
    def eval(self, text: str):
        for word in text.split():
            yield (word, len(word))

# Использование как table-valued function
spark.sql("SELECT * FROM TokenizeUDTF('hello world spark')").show()
# +------+------+
# |  word|length|
# +------+------+
# | hello|     5|
# | world|     5|
# | spark|     5|
# +------+------+

eval() — generator-метод, который yields rows. Каждая входная строка может породить N выходных строк (включая 0).

Polymorphic UDTF: динамическая схема

В polymorphic режиме UDTF реализует classmethod analyze(), который получает мета-информацию об аргументах (тип каждого, литеральное значение если scalar) и возвращает StructType для выходной схемы. Spark вызывает analyze() один раз во время planning — результат используется как outputSchema конкретного вызова.

from pyspark.sql.functions import udtf, AnalyzeArgument, AnalyzeResult
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

@udtf
class SplitColumnsUDTF:
    @staticmethod
    def analyze(text_arg: AnalyzeArgument, n_arg: AnalyzeArgument) -> AnalyzeResult:
        # text_arg.value может быть None (column reference) или конкретное значение (literal)
        n = n_arg.value  # сколько колонок-частей вернуть -- известно на planning-фазе
        fields = [StructField(f"part_{i}", StringType(), True) for i in range(n)]
        fields.append(StructField("remainder", StringType(), True))
        return AnalyzeResult(schema=StructType(fields))

    def eval(self, text: str, n: int):
        parts = text.split(",", n)
        # Pad до n частей и положить остаток в remainder
        head = parts[:n] + [None] * max(0, n - len(parts))
        remainder = parts[n] if len(parts) > n else None
        yield tuple(head + [remainder])

# n=2 -> схема (part_0, part_1, remainder)
spark.sql("SELECT * FROM SplitColumnsUDTF('a,b,c,d,e', 2)").show()
# +------+------+---------+
# |part_0|part_1|remainder|
# +------+------+---------+
# |     a|     b|    c,d,e|
# +------+------+---------+

# n=4 -> схема (part_0, part_1, part_2, part_3, remainder)
spark.sql("SELECT * FROM SplitColumnsUDTF('a,b,c,d,e', 4)").show()
# +------+------+------+------+---------+
# |part_0|part_1|part_2|part_3|remainder|
# +------+------+------+------+---------+
# |     a|     b|     c|     d|        e|
# +------+------+------+------+---------+

Один и тот же UDTF возвращает разную outputSchema в зависимости от literal-аргумента n. Это и есть polymorphism: одна функция, разные структурные типы выхода.

Lazy evaluation и terminate()

UDTF — lazy evaluation: код в eval() не выполняется до тех пор, пока action не запустит план. Это касается и polymorphic-варианта: analyze() вызывается на planning, eval() — только при materialization.

Опциональный метод terminate() вызывается ПОСЛЕ обработки всех входных строк и может yield-ить дополнительные rows — удобно для подведения итогов:

@udtf(returnType="event_type: string, count: int")
class EventCounterUDTF:
    def __init__(self):
        self._counts = {}

    def eval(self, event_type: str):
        self._counts[event_type] = self._counts.get(event_type, 0) + 1
        # eval ничего не yield -- агрегируем in-memory

    def terminate(self):
        # finalization: отдаём агрегацию когда вход исчерпан
        for event_type, count in self._counts.items():
            yield (event_type, count)

Когда UDTF vs alternatives

СценарийИнструментПричина
Скаляр в скалярBuilt-in / Pandas UDFSimpler, лучше оптимизация
1 row -> N rows фиксированной схемыUDTF (фиксированная)explode() работает только для arrays
1 row -> N rows динамической схемыPolymorphic UDTFНевозможно через explode/Pandas UDF
Группа -> новый DataFrameapplyInPandas (Grouped Map)UDTF работает per-row, не per-group
Парсинг JSON в variant-колонкиPolymorphic UDTFСхема зависит от schema-inferring аргумента
INFO

Polymorphic UDTF лучше всего подходят для случаев, где схема выхода известна на planning-фазе (зависит от literal-аргументов), но не известна на write-time UDTF-кода. Классические применения: schema-aware JSON-парсеры, конфигурируемые pivot-функции, динамическое разбиение строк по конфигу разделителей.

Что дальше?

В следующем уроке мы разберём Scala UDF — как UDF на JVM-языке обходит проблему сериализации, обеспечивая производительность ~2x от встроенных функций.

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 5