Pandas UDF: Arrow-based векторизация
Проблема, которую решает Pandas UDF
Как мы выяснили в предыдущем уроке, Python UDF выполняет сериализацию на каждую строку. При 1 миллиарде строк это 1 миллиард циклов serialize -> socket -> deserialize.
Pandas UDF решают эту проблему через батчевую обработку: вместо строк данные передаются пакетами (batches) по тысячам строк за раз, используя Apache Arrow columnar format.
Результат: Pandas UDF в 5-20x быстрее обычного Python UDF.
Apache Arrow: ключ к производительности
Apache Arrow — это in-memory columnar format, который позволяет:
- Columnar layout: данные одного столбца хранятся непрерывно в памяти, что оптимально для vectorized операций
- Zero-copy reads: Python может читать Arrow-данные без копирования — JVM и Python разделяют один участок памяти
- Batch transfer: тысячи строк передаются одним блоком вместо per-row сериализации
Типы 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.
Выбор типа 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 быстрее по двум причинам:
- Arrow батчи вместо per-row сериализации (сокращает transfer overhead)
- NumPy vectorized операции вместо Python for-loop (сокращает compute overhead)
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-уровне.
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 UDF | Simpler, лучше оптимизация |
| 1 row -> N rows фиксированной схемы | UDTF (фиксированная) | explode() работает только для arrays |
| 1 row -> N rows динамической схемы | Polymorphic UDTF | Невозможно через explode/Pandas UDF |
| Группа -> новый DataFrame | applyInPandas (Grouped Map) | UDTF работает per-row, не per-group |
| Парсинг JSON в variant-колонки | Polymorphic UDTF | Схема зависит от schema-inferring аргумента |
Polymorphic UDTF лучше всего подходят для случаев, где схема выхода известна на planning-фазе (зависит от literal-аргументов), но не известна на write-time UDTF-кода. Классические применения: schema-aware JSON-парсеры, конфигурируемые pivot-функции, динамическое разбиение строк по конфигу разделителей.
Что дальше?
В следующем уроке мы разберём Scala UDF — как UDF на JVM-языке обходит проблему сериализации, обеспечивая производительность ~2x от встроенных функций.