Pandas UDF (Arrow): SCALAR, GROUPED_MAP, GROUPED_AGG

Векторизованные UDF через Apache Arrow: SCALAR, GROUPED_MAP, GROUPED_AGG, Iterator - архитектура, производительность, паттерны ML inference и кастомной аналитики

core optimization

Зачем появились Pandas UDF

В предыдущем уроке мы разобрали, почему Row UDF - это «крайний случай»: каждая строка сериализуется через pickle, передаётся в Python worker через IPC сокет, обрабатывается, сериализуется обратно. На 100 миллионах строк - 200 миллионов операций сериализации.

Инженеры Databricks и Apache Arrow Core Team задались вопросом: а что, если передавать данные не по одной строке, а целыми колонками в бинарном формате, который одинаково понимают и JVM и Python? Так появились Pandas UDF (в Spark 2.3, 2018).

Главная идея: вместо строчной (row-by-row) обработки - колончатая (columnar). Вместо pickle - Apache Arrow. Вместо Python объектов - pandas.Series и pandas.DataFrame, которые уже умеют эффективно работать с Arrow-данными.

Результат: 5–20× ускорение по сравнению с Row UDF при одинаковой логике - и возможность использовать весь numpy/scipy/sklearn экосистем прямо в Spark pipeline.

Apache Arrow: революция в передаче данных

Apache Arrow - это открытый стандарт колончатого представления данных в памяти, разработанный специально для аналитических систем. Ключевой принцип: данные хранятся не по строкам ({id: 1, name: "Alice", amount: 100}, {id: 2, ...}), а по колонкам ([1, 2, 3, ...], ["Alice", "Bob", ...], [100, 200, ...]).

Зачем колончатый формат быстрее для аналитики? Потому что аналитические операции обычно работают с несколькими колонками всех строк, а не со всеми колонками одной строки. Когда данные лежат колончато - процессорный кэш заполняется нужными данными. Когда строчно - для получения одной колонки нужно «перепрыгивать» через данные других колонок.

Но главное преимущество Arrow для Pandas UDF - это общая спецификация формата. И JVM (через Arrow Java library), и Python (через pyarrow) работают с одним бинарным форматом. Это означает, что при передаче данных между ними не нужно полноценной сериализации - достаточно передать указатель на буфер в памяти или скопировать бинарные данные как есть.

Arrow execution pipeline

Ключевое отличие от Row UDF: данные передаются не строка-за-строкой, а батчами (batch). Один батч - это Arrow RecordBatch с N строками. По умолчанию N = 10 000 (настраивается через spark.sql.execution.arrow.maxRecordsPerBatch). Вместо 10 000 сериализаций - одна передача бинарного буфера.

Почему Arrow быстрее pickle

Pickle - это универсальный формат сериализации Python объектов. Он работает с любыми объектами (словари, классы, списки), но требует обхода графа объекта и записи его структуры. Строка "Alice" в pickle занимает примерно 12 байт плюс метаданные протокола.

Arrow - это фиксированный бинарный формат для табличных данных. Колонка из 10 000 строк превращается в непрерывный блок памяти (буфер). Передача этого буфера через IPC - это фактически одна операция write() в pipe.

На практике разница в скорости сериализации Arrow vs pickle:

  • Числовые колонки (int, double): Arrow ≈ 20–50× быстрее
  • Строковые колонки: Arrow ≈ 5–10× быстрее
  • Битовая маска null значений: Arrow хранит как компактный bitset

Pandas UDF vs Python Row UDF: сравнение архитектур

Разница не только в количестве IPC-вызовов. Внутри Python worker логика тоже выполняется принципиально по-другому:

  • Row UDF: for row in batch: result = func(row) - Python цикл, каждая итерация - интерпретатор
  • Pandas UDF: result_series = func(pandas_series) - одна операция над pandas.Series, которая внутри использует numpy C-код

NumPy операции (сложение, умножение, apply встроенных функций) выполняются в скомпилированном C коде без Python overhead. Это и есть векторизация - применение операции сразу ко всему массиву на уровне CPU SIMD инструкций.

Включение Arrow optimization

Arrow для Pandas UDF включён по умолчанию начиная с Spark 3.0. Для старых версий или если было отключено явно:

spark = SparkSession.builder \
    .config("spark.sql.execution.arrow.pyspark.enabled", "true") \
    .getOrCreate()

# Проверить текущее значение
print(spark.conf.get("spark.sql.execution.arrow.pyspark.enabled"))  # "true"

Размер батча - ключевой параметр для баланса memory vs throughput:

# По умолчанию: 10 000 строк на батч
# Увеличить для throughput (если строки маленькие):
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "50000")

# Уменьшить для memory safety (если строки большие или ML модель жирная):
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "1000")

Выбор размера батча - это trade-off: больший батч = меньше IPC overhead + лучше векторизация, но больше памяти Python worker. Мы вернёмся к этому в разделе про Arrow batch sizing.

SCALAR Pandas UDF: колонка → колонка

SCALAR (Series to Series) - самый распространённый тип Pandas UDF. Функция получает pandas.Series (или несколько Series) и возвращает pandas.Series того же размера.

Это прямой аналог Row UDF, но вместо обработки строк по одной - векторизованная обработка батчей. Spark вызывает функцию по одному разу на каждый батч Arrow данных.

from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType, StringType
import pandas as pd
import numpy as np

# Объявление: декоратор с returnType
@pandas_udf(DoubleType())
def log_transform(series: pd.Series) -> pd.Series:
    # Вся колонка как pd.Series - можно использовать numpy operations
    return np.log1p(series.fillna(0))

# Несколько входных колонок
@pandas_udf(DoubleType())
def weighted_score(amount: pd.Series, frequency: pd.Series, weight: pd.Series) -> pd.Series:
    # Все три аргумента - pd.Series одного размера (один батч)
    return (amount * frequency * weight).clip(0, 1)

# Строковая трансформация
@pandas_udf(StringType())
def normalize_email(emails: pd.Series) -> pd.Series:
    return emails.str.strip().str.lower()

Использование идентично Row UDF:

from pyspark.sql.functions import col

df.withColumn("log_amount", log_transform(col("amount"))) \
  .withColumn("score", weighted_score(col("amount"), col("freq"), col("weight"))) \
  .withColumn("email_clean", normalize_email(col("email")))

Когда SCALAR UDF выигрывает у Row UDF

SCALAR Pandas UDF особенно хорош когда:

  • Логика выражается через numpy/pandas операции (они уже векторизованы)
  • Используется scipy/statsmodels для математики (принимают Series/array)
  • ML inference: модель принимает numpy array и возвращает array - именно так работают sklearn, lightgbm, xgboost
import pandas as pd
import numpy as np
from scipy import stats
from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType

# Box-Cox transformation из scipy - принимает array, возвращает array
@pandas_udf(DoubleType())
def boxcox_transform(series: pd.Series) -> pd.Series:
    clean = series.dropna().clip(lower=0.001)
    transformed, _ = stats.boxcox(clean)
    # Вернуть Series с исходным индексом, NaN для пропусков
    result = pd.Series(index=series.index, dtype=float)
    result[series.notna()] = transformed
    return result

SCALAR UDF для ML inference

ML inference - один из главных use cases для SCALAR Pandas UDF. Модель загружается один раз при создании UDF (в Python worker), батчи строк передаются в model.predict() как numpy array:

import pandas as pd
import numpy as np
import joblib
from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType, ArrayType, FloatType

# Модель загружается ОДИН РАЗ при первом вызове на executor
# (Python worker кэширует переменные модуля между батчами)
_model = None

def _get_model():
    global _model
    if _model is None:
        _model = joblib.load("/shared/models/fraud_v3.pkl")
    return _model

@pandas_udf(DoubleType())
def predict_fraud_score(
    amount: pd.Series,
    merchant_cat: pd.Series,
    hour: pd.Series
) -> pd.Series:
    model = _get_model()

    # Собираем feature matrix как numpy array - нет Python цикла
    X = np.column_stack([
        amount.fillna(0).values,
        merchant_cat.fillna(-1).values,
        hour.fillna(12).values,
    ])

    scores = model.predict_proba(X)[:, 1]
    return pd.Series(scores)

df.withColumn("fraud_score",
    predict_fraud_score(col("amount"), col("merchant_cat"), col("hour"))
)

Паттерн с _model = None и _get_model() - это ленивая инициализация. Python worker (процесс) живёт на executor, пока Spark Job не завершится. Глобальные переменные в модуле сохраняются между батчами одного Job. Первый батч загрузит модель, все последующие батчи на том же worker её уже найдут в памяти.

GROUPED_MAP: группа как pandas.DataFrame

GROUPED_MAP (в Spark 3.0+ называется applyInPandas) - качественно иной паттерн. Функция получает не Series с данными одного батча, а весь pandas.DataFrame одной группы целиком. Это split-apply-combine в полную силу.

Схема работы:

Функция вызывается по одному разу на каждую группу. Она получает полный pandas DataFrame с данными группы - все строки, все колонки - и должна вернуть pandas DataFrame (с тем же или другим набором колонок).

from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, LongType
import pandas as pd
import numpy as np

# Схема возвращаемого DataFrame - обязательна для applyInPandas
output_schema = StructType([
    StructField("store_id", LongType()),
    StructField("product_id", LongType()),
    StructField("sales_normalized", DoubleType()),
    StructField("sales_rank", LongType()),
])

def normalize_within_store(store_df: pd.DataFrame) -> pd.DataFrame:
    # store_df содержит ВСЕ строки одного store_id
    # Можно делать внутри-групповую аналитику
    total = store_df["sales"].sum()
    store_df = store_df.copy()
    store_df["sales_normalized"] = store_df["sales"] / (total or 1)
    store_df["sales_rank"] = store_df["sales"].rank(ascending=False).astype(int)
    return store_df[["store_id", "product_id", "sales_normalized", "sales_rank"]]

result = df.groupBy("store_id").applyInPandas(normalize_within_store, schema=output_schema)

Split-apply-combine: как Spark материализует группы

Прежде чем вызвать функцию, Spark должен собрать все строки одной группы на одном executor. Это требует shuffle: все строки с одинаковым store_id должны оказаться в одной партиции.

По этой причине groupBy().applyInPandas() всегда включает shuffle шаг - так же, как groupBy().agg(). Если данные уже партиционированы по ключу группировки, shuffle можно избежать через repartition().

После shuffle все строки одной группы собираются в памяти Python worker как один pandas DataFrame. Это означает: если группа большая (например, топ-магазин с 10 миллионами строк) - весь этот DataFrame должен поместиться в памяти Python worker. Это главный риск GROUPED_MAP.

Кейс: обучение локальных моделей

GROUPED_MAP - идеальный паттерн для обучения отдельной модели на каждую группу (магазин, регион, товарная категория):

from sklearn.linear_model import LinearRegression
import pandas as pd
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, LongType

model_output_schema = StructType([
    StructField("store_id", LongType()),
    StructField("feature", StringType()),
    StructField("coefficient", DoubleType()),
    StructField("intercept", DoubleType()),
])

def train_store_model(store_df: pd.DataFrame) -> pd.DataFrame:
    store_id = store_df["store_id"].iloc[0]

    if len(store_df) < 10:  # недостаточно данных для обучения
        return pd.DataFrame(columns=["store_id", "feature", "coefficient", "intercept"])

    X = store_df[["price", "promo_flag", "day_of_week"]].values
    y = store_df["sales"].values

    model = LinearRegression()
    model.fit(X, y)

    features = ["price", "promo_flag", "day_of_week"]
    return pd.DataFrame({
        "store_id": store_id,
        "feature": features,
        "coefficient": model.coef_,
        "intercept": model.intercept_,
    })

# Обучаем отдельную модель для каждого магазина параллельно
coefficients = df.groupBy("store_id").applyInPandas(train_store_model, schema=model_output_schema)

Этот паттерн называется federated learning или local model ensemble. Вместо одной глобальной модели - N маленьких моделей, каждая заточена под свою группу. Spark распараллеливает обучение между executors: каждый executor обучает несколько групп.

Риски GROUPED_MAP: Memory explosion

Главная опасность GROUPED_MAP - неравномерное распределение размеров групп. Если у вас 1 000 магазинов, и 999 из них имеют по 1 000 строк, а один - 50 миллионов строк, то executor, которому достался мегамагазин, попытается загрузить 50 миллионов строк в pandas DataFrame. При ширине схемы в 20 double-колонок это ~8 GB только для данных, плюс pandas overhead ~2×.

# Диагностика: смотрим распределение размеров групп ДО applyInPandas
df.groupBy("store_id").count().orderBy("count", ascending=False).show(10)

# Если есть giant groups - рассмотреть:
# 1. repartition по подгруппе (store_id + week_id) вместо store_id
# 2. Предварительная агрегация перед applyInPandas
# 3. Фильтрация малых/больших групп перед обработкой

Правило: в GROUPED_MAP нужно знать максимальный размер группы. Если максимальная группа > 1–2 GB в памяти - это риск OOM.

GROUPED_AGG: кастомная агрегация

GROUPED_AGG (Series to Scalar) - для случаев, когда нужна агрегация, которую нельзя выразить через встроенные F.sum(), F.avg(), F.percentile_approx() и т.д.

Функция получает pandas.Series всех значений в группе и должна вернуть одно скалярное значение. Spark затем собирает одно значение от каждой группы - это и есть агрегация.

from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType, LongType
import pandas as pd
import numpy as np
from scipy import stats

# Кастомная агрегация: мода (самое частое значение)
@pandas_udf(DoubleType())
def mode_agg(series: pd.Series) -> float:
    mode_result = stats.mode(series.dropna(), keepdims=True)
    return float(mode_result.mode[0]) if len(mode_result.mode) > 0 else None

# Trimmed mean: среднее без 5% outliers с каждой стороны
@pandas_udf(DoubleType())
def trimmed_mean(series: pd.Series) -> float:
    return float(stats.trim_mean(series.dropna(), proportiontocut=0.05))

# Коэффициент вариации (CV = std/mean)
@pandas_udf(DoubleType())
def coeff_variation(series: pd.Series) -> float:
    clean = series.dropna()
    if clean.mean() == 0 or len(clean) < 2:
        return None
    return float(clean.std() / clean.mean())

# Использование как обычная агрегация
df.groupBy("category").agg(
    mode_agg(col("price")).alias("modal_price"),
    trimmed_mean(col("amount")).alias("trimmed_avg"),
    coeff_variation(col("daily_sales")).alias("sales_cv"),
)

GROUPED_AGG vs applyInPandas: когда что выбрать

Задача Метод
Новая метрика для каждой группы (одно число) GROUPED_AGG: groupBy().agg(pandas_udf(...))
Трансформация строк внутри группы GROUPED_MAP: groupBy().applyInPandas()
Обучение модели per-group GROUPED_MAP: groupBy().applyInPandas()
Нормализация / ранжирование внутри группы GROUPED_MAP: groupBy().applyInPandas()
t-test / chi-test между группами GROUPED_AGG на каждую группу, затем join

GROUPED_AGG концептуально проще: он возвращает ровно одну строку на группу, схему описывать не нужно (она вычисляется из returnType декоратора). GROUPED_MAP гибче, но требует явного описания output schema.

Iterator Pandas UDF: memory-efficient батч-обработка

С Spark 3.0 появился новый подход: Iterator of Series и Iterator of DataFrames. Вместо того, чтобы функция получала один батч, она получает итератор по всем батчам. Это позволяет:

  • Загрузить тяжёлый ресурс (ML модель, справочник) один раз для всей партиции
  • Обрабатывать батчи один за другим, не загружая всю партицию в память
from typing import Iterator
import pandas as pd
import numpy as np
import joblib
from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType

@pandas_udf(DoubleType())
def predict_batch_iter(iterator: Iterator[pd.Series]) -> Iterator[pd.Series]:
    # Модель загружается ОДИН РАЗ для всей партиции
    model = joblib.load("/shared/models/scorer_v2.pkl")

    for amount_batch in iterator:
        # Обрабатываем каждый батч по очереди
        X = amount_batch.fillna(0).values.reshape(-1, 1)
        scores = model.predict(X)
        yield pd.Series(scores)

df.withColumn("score", predict_batch_iter(col("amount")))

Iterator UDF с несколькими входными колонками:

from typing import Iterator, Tuple

@pandas_udf(DoubleType())
def predict_multi_iter(
    iterator: Iterator[Tuple[pd.Series, pd.Series, pd.Series]]
) -> Iterator[pd.Series]:
    model = joblib.load("/shared/models/fraud_v3.pkl")

    for amount, category, hour in iterator:
        X = np.column_stack([
            amount.fillna(0).values,
            category.fillna(-1).values,
            hour.fillna(12).values,
        ])
        yield pd.Series(model.predict_proba(X)[:, 1])

df.withColumn("fraud_score",
    predict_multi_iter(col("amount"), col("merchant_cat"), col("hour"))
)

Разница между обычным SCALAR и Iterator UDF:

  • SCALAR: модель загружается при первом батче и кэшируется в глобальной переменной → один раз на Python worker (жизненный цикл процесса)
  • Iterator: модель загружается явно в начале итерации → один раз на партицию, явный контроль

Iterator UDF - более явный и предсказуемый паттерн. Для production ML inference предпочтительнее Iterator, потому что понятно, когда модель загружается, и нет скрытого состояния через глобальные переменные.

Arrow batch sizing: баланс throughput и memory

Размер Arrow батча (по умолчанию 10 000 строк) - ключевой параметр для настройки производительности.

Как рассчитать безопасный размер батча:

# Оценка размера одной строки
row_size_bytes = df.schema.jsonValue()  # грубо: количество полей × средний размер

# Допустим: 50 колонок × 8 байт = 400 байт/строка
# При 10 000 строк на батч: 4 MB
# pandas overhead ~2×: ~8 MB на батч в Python worker
# Если модель занимает 200 MB, суммарно: ~208 MB на батч-цикл

# Если executor memory = 4 GB и tasks_per_executor = 4:
# Доступно на Python worker: ~400 MB
# Максимальный батч: (400 MB / 2) / 400 bytes ≈ 500 000 строк

# Но учитываем, что pandas.DataFrame обычно занимает 2-5× от сырых данных:
# Безопасный батч: ~50 000 строк при умеренной схеме

На практике: начинайте с дефолтного 10 000, измеряйте время выполнения Stage и utilization памяти в Spark UI, затем увеличивайте батч до тех пор, пока не увидите признаки memory pressure (GC pressure, slow tasks).

Pandas memory model

pandas.DataFrame хранит данные в numpy arrays - один array на колонку. Numpy array занимает в памяти примерно N строк × size(dtype) байт. Для int64 - 8 байт на элемент, для float64 - 8 байт, для строк (object dtype) - 8 байт указатель + размер строки в куче.

Дополнительный overhead: при выполнении операций над Series (series * 2, series.str.lower()) pandas создаёт новый array с результатом - исходный при этом не удаляется, пока на него есть ссылки. В момент вычисления в памяти одновременно существуют: исходный батч + промежуточные результаты + финальный результат. Это объясняет, почему pandas реально потребляет 3–5× от «теоретического» размера данных.

Pandas UDF в explain plan

Pandas UDF видны в физическом плане через explain():

df.withColumn("score", predict_fraud_score(col("amount"), col("hour"))).explain()

# Physical Plan:
# *(2) Project [id#5, amount#6, ...]
# +- ArrowEvalPython [predict_fraud_score(amount#6, hour#7)], [id#5, ..., score#42]
#    +- *(1) FileScan parquet [...] PushedFilters: [], ...

ArrowEvalPython - маркер Pandas UDF в физическом плане. В отличие от BatchEvalPython (Row UDF), здесь данные передаются через Arrow.

Как и Row UDF, Pandas UDF разрывает WholeStageCodeGen pipeline. Фильтры выше по плану не могут быть протолкнуты сквозь ArrowEvalPython. Поэтому правило «применяй фильтры до UDF» работает и для Pandas UDF.

Ограничения Pandas UDF

Pandas UDF значительно быстрее Row UDF, но сохраняет часть его ограничений.

1. Нет полной Catalyst optimization. Catalyst не видит внутренности Python функции. Предикаты после Pandas UDF не pushdown в источник данных. Это общая черта с Row UDF.

2. Arrow serialization cost. Arrow быстрее pickle, но не бесплатен. Конвертация между Tungsten UnsafeRow (строчный формат) и Arrow RecordBatch (колончатый формат) требует CPU и памяти. Для simple operations (сложение двух колонок) встроенные функции всё равно быстрее.

3. pandas overhead. Создание pandas.Series из Arrow данных - дополнительный шаг. Для очень маленьких батчей или простой логики этот overhead заметен.

4. applyInPandas требует shuffle. groupBy().applyInPandas() всегда производит shuffle. Для малых датасетов это может быть дороже, чем обработать данные без группировки.

5. Python GIL для compute-heavy logic. Если логика UDF не пробрасывается в numpy/scipy (а остаётся Python кодом с циклами), GIL снова становится ограничением.

6. Версия pandas на executor. Pandas UDF требует, чтобы версия pandas на Python worker совпадала с ожидаемой. В production кластерах это означает контроль за версиями Python environment.

Практика: benchmark Row UDF vs Pandas UDF vs built-in

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.functions import pandas_udf, udf
from pyspark.sql.types import DoubleType
import pandas as pd
import numpy as np
import time

spark = SparkSession.builder \
    .appName("PandasUDF-Benchmark") \
    .config("spark.sql.execution.arrow.pyspark.enabled", "true") \
    .getOrCreate()

# Генерируем тестовые данные: 10 миллионов строк
df = spark.range(10_000_000).select(
    (F.rand() * 10000).alias("amount"),
    (F.rand() * 100).alias("frequency"),
).cache()
df.count()

# ── Вариант 1: Python Row UDF ──────────────────────────────────────────────────
@udf(DoubleType())
def row_udf_score(amount, frequency):
    if amount is None or frequency is None:
        return None
    return float(np.log1p(amount) * np.sqrt(frequency))

start = time.time()
df.withColumn("score", row_udf_score(F.col("amount"), F.col("frequency"))).count()
print(f"Row UDF:        {time.time() - start:.1f}s")

# ── Вариант 2: Pandas SCALAR UDF ──────────────────────────────────────────────
@pandas_udf(DoubleType())
def pandas_udf_score(amount: pd.Series, frequency: pd.Series) -> pd.Series:
    return np.log1p(amount.fillna(0)) * np.sqrt(frequency.fillna(0))

start = time.time()
df.withColumn("score", pandas_udf_score(F.col("amount"), F.col("frequency"))).count()
print(f"Pandas UDF:     {time.time() - start:.1f}s")

# ── Вариант 3: встроенные функции ─────────────────────────────────────────────
start = time.time()
df.withColumn("score",
    F.log1p(F.col("amount")) * F.sqrt(F.col("frequency"))
).count()
print(f"Built-in:       {time.time() - start:.1f}s")

Ожидаемые результаты:

Row UDF:        52.3s   (pickle overhead, row-by-row)
Pandas UDF:      8.7s   (Arrow batches, numpy vectorization)
Built-in:        2.1s   (WholeStageCodeGen, JVM native)

Pandas UDF в 6× быстрее Row UDF. Built-in functions в 4× быстрее Pandas UDF. Итого разница Row UDF vs built-in - ~25×.

Вывод: Pandas UDF - хороший выбор когда нет встроенного аналога и логика выражается через numpy/scipy. Но если логика есть в pyspark.sql.functions - встроенные функции всегда быстрее.

Практика: SCALAR UDF - ML inference на prod данных

Реалистичный пример: ML scoring через Pandas UDF в production ETL.

import pandas as pd
import numpy as np
import joblib
from typing import Iterator, Tuple
from pyspark.sql.functions import pandas_udf, col
from pyspark.sql.types import DoubleType

# Используем Iterator UDF для явной загрузки модели один раз на партицию
@pandas_udf(DoubleType())
def churn_score(
    iterator: Iterator[Tuple[pd.Series, pd.Series, pd.Series, pd.Series]]
) -> Iterator[pd.Series]:
    # Модель загружается ОДИН РАЗ при первом батче на данной партиции
    model = joblib.load("/shared/models/churn_v5.pkl")
    feature_cols = ["days_since_login", "total_purchases", "avg_order_value", "support_tickets"]

    for days_login, purchases, avg_order, tickets in iterator:
        X = pd.DataFrame({
            "days_since_login":  days_login.fillna(999),
            "total_purchases":   purchases.fillna(0),
            "avg_order_value":   avg_order.fillna(0),
            "support_tickets":   tickets.fillna(0),
        })[feature_cols].values

        probs = model.predict_proba(X)[:, 1]
        yield pd.Series(probs)

# В production pipeline:
users_df = spark.table("dim_users").filter(col("is_active") == True)

scored = users_df.withColumn("churn_probability",
    churn_score(
        col("days_since_login"),
        col("total_purchases"),
        col("avg_order_value"),
        col("support_tickets"),
    )
)

# Сохраняем результат - только пользователи с высоким риском оттока
scored.filter(col("churn_probability") > 0.7) \
      .select("user_id", "churn_probability", "segment") \
      .write.mode("overwrite").saveAsTable("ml.churn_scores")

Практика: GROUPED_MAP - feature engineering по сессиям

Задача: для каждой пользовательской сессии посчитать накопленные фичи (rolling features) - то, что встроенные window functions в Spark делают, но для сложной кастомной логики удобнее в pandas:

import pandas as pd
import numpy as np
from pyspark.sql.types import StructType, StructField, LongType, StringType, DoubleType

session_features_schema = StructType([
    StructField("user_id", LongType()),
    StructField("session_id", StringType()),
    StructField("event_timestamp", LongType()),
    StructField("event_type", StringType()),
    StructField("time_since_session_start", DoubleType()),
    StructField("events_so_far", LongType()),
    StructField("avg_interval_seconds", DoubleType()),
])

def compute_session_features(session_df: pd.DataFrame) -> pd.DataFrame:
    # session_df: все события одной сессии (один user_id + session_id)
    session_df = session_df.sort_values("event_timestamp").copy()

    session_start = session_df["event_timestamp"].min()
    session_df["time_since_session_start"] = (
        session_df["event_timestamp"] - session_start
    ) / 1000.0  # в секундах

    # Количество событий до текущего (включая текущее)
    session_df["events_so_far"] = range(1, len(session_df) + 1)

    # Средний интервал между событиями до текущего
    intervals = session_df["event_timestamp"].diff().fillna(0) / 1000.0
    session_df["avg_interval_seconds"] = intervals.expanding().mean()

    return session_df[[
        "user_id", "session_id", "event_timestamp", "event_type",
        "time_since_session_start", "events_so_far", "avg_interval_seconds"
    ]]

events = spark.table("raw.user_events")
features = events.groupBy("user_id", "session_id") \
                 .applyInPandas(compute_session_features, schema=session_features_schema)

Практика: GROUPED_AGG - кастомные статистические метрики

Задача: посчитать для каждой категории товаров нестандартные метрики (Gini коэффициент продаж, IQR), которых нет в Spark built-in:

import pandas as pd
import numpy as np
from pyspark.sql.functions import pandas_udf, col
from pyspark.sql.types import DoubleType

@pandas_udf(DoubleType())
def gini_coefficient(series: pd.Series) -> float:
    """Коэффициент Джини - мера неравномерности продаж внутри категории."""
    values = series.dropna().values
    if len(values) < 2:
        return None
    values = np.sort(values)
    n = len(values)
    # Формула Gini через cumulative sum
    index = np.arange(1, n + 1)
    return float((2 * np.sum(index * values) / (n * np.sum(values))) - (n + 1) / n)

@pandas_udf(DoubleType())
def interquartile_range(series: pd.Series) -> float:
    """IQR: разница между 75-м и 25-м перцентилями."""
    clean = series.dropna()
    if len(clean) < 4:
        return None
    return float(clean.quantile(0.75) - clean.quantile(0.25))

@pandas_udf(DoubleType())
def sales_entropy(series: pd.Series) -> float:
    """Энтропия распределения продаж (насколько равномерно распределены)."""
    clean = series.dropna().clip(lower=0.001)
    total = clean.sum()
    if total == 0:
        return None
    probs = clean / total
    return float(-np.sum(probs * np.log(probs)))

# Применяем как обычные агрегации
df.groupBy("category_id").agg(
    gini_coefficient(col("daily_sales")).alias("sales_gini"),
    interquartile_range(col("daily_sales")).alias("sales_iqr"),
    sales_entropy(col("daily_sales")).alias("sales_entropy"),
)

Anti-patterns Pandas UDF

1. GROUPED_MAP с гигантскими группами без проверки

# ОПАСНО: если одна group = 50M строк → OOM
df.groupBy("country_code").applyInPandas(heavy_func, schema)

# СНАЧАЛА: проверить размеры групп
df.groupBy("country_code").count().orderBy("count", ascending=False).show(5)
# Если топ-группа > 10M строк - нужен другой подход

2. Python цикл внутри Pandas UDF

# ПЛОХО: цикл по строкам - теряем преимущество векторизации
@pandas_udf(DoubleType())
def bad_score(series: pd.Series) -> pd.Series:
    result = []
    for val in series:  # Python цикл - как Row UDF!
        result.append(np.log1p(val) if val else 0.0)
    return pd.Series(result)

# ХОРОШО: numpy операция над всей Series
@pandas_udf(DoubleType())
def good_score(series: pd.Series) -> pd.Series:
    return np.log1p(series.fillna(0))  # векторизованно

3. Использование вместо встроенных функций Spark

# ПЛОХО: Pandas UDF для того, что есть в F.*
@pandas_udf(StringType())
def upper_pandas(s: pd.Series) -> pd.Series:
    return s.str.upper()

# ХОРОШО: встроенная функция
F.upper(col("name"))

# ПЛОХО: Pandas UDF для date extraction
@pandas_udf(IntegerType())
def extract_year(s: pd.Series) -> pd.Series:
    return pd.to_datetime(s).dt.year

# ХОРОШО: встроенная функция
F.year(F.col("created_at"))

4. Конвертация Spark DataFrame в pandas внутри UDF

# ПЛОХО: попытка использовать spark внутри UDF - это невозможно
@pandas_udf(DoubleType())
def lookup_from_spark(ids: pd.Series) -> pd.Series:
    ref = spark.table("reference").toPandas()  # ОШИБКА: spark недоступен в worker
    ...

# ХОРОШО: broadcast dictionary или join до UDF
ref_dict = spark.table("reference").rdd.collectAsMap()
ref_bv = spark.sparkContext.broadcast(ref_dict)

@pandas_udf(DoubleType())
def lookup_broadcast(ids: pd.Series) -> pd.Series:
    ref = ref_bv.value
    return ids.map(lambda x: ref.get(x, 0.0))

5. Загрузка модели на каждый батч (не на партицию)

# ПЛОХО: модель загружается на каждый батч (тысячи раз)
@pandas_udf(DoubleType())
def score_bad(series: pd.Series) -> pd.Series:
    model = joblib.load("/path/model.pkl")  # загружается на каждый батч!
    return pd.Series(model.predict(series.values.reshape(-1, 1)))

# ХОРОШО: Iterator UDF - модель загружается один раз на партицию
@pandas_udf(DoubleType())
def score_good(iterator: Iterator[pd.Series]) -> Iterator[pd.Series]:
    model = joblib.load("/path/model.pkl")  # загружается один раз
    for batch in iterator:
        yield pd.Series(model.predict(batch.values.reshape(-1, 1)))

Checklist: когда использовать Pandas UDF

Задача Рекомендация
Математика, строки, даты Встроенные функции F.*
HOF над массивами F.transform, F.aggregate
NumPy/scipy вычисления поверх колонки SCALAR Pandas UDF
ML inference (sklearn, xgboost, lightgbm) Iterator Pandas UDF (явная загрузка модели)
Кастомная агрегация (Gini, entropy, IQR) GROUPED_AGG Pandas UDF
Трансформация строк внутри группы GROUPED_MAP (applyInPandas)
Обучение per-group моделей GROUPED_MAP (applyInPandas)
Обращение к внешним ресурсам (DB, HTTP) mapPartitions (не Pandas UDF)
Группы > 5M строк в GROUPED_MAP Сначала проверить размеры, возможно repartition
Медленный Pandas UDF Убрать Python циклы, использовать numpy операции
Повторная загрузка модели Перейти на Iterator UDF