Pandas API on Spark (pyspark.pandas): когда и как

pyspark.pandas (бывш. Koalas) даёт pandas-совместимый интерфейс поверх Spark. Разбираем, где это ускоряет миграцию, где замедляет работу, и как API соотносится с нативным PySpark.

core

Что такое pyspark.pandas

Pandas API on Spark (pyspark.pandas, ранее - проект Koalas от Databricks) - слой совместимости, который позволяет запускать pandas-код на Spark без переписывания.

Идея выглядит привлекательно: Data Scientist написал notebook на pandas - 500 строк хорошо выверенного кода для очистки, агрегации и подготовки признаков. Данных стало в 100 раз больше, pandas падает по памяти. Что делать? Переписывать на нативный PySpark - это недели работы. Именно для этой ситуации появился pyspark.pandas.

# Обычный pandas - выполняется на одной машине
import pandas as pd
df = pd.read_csv("data.csv")
df["revenue"].sum()

# pyspark.pandas - тот же синтаксис, выполняется на кластере
import pyspark.pandas as ps
df = ps.read_csv("s3://bucket/data.csv")
df["revenue"].sum()

Два блока кода выглядят идентично - строка import pyspark.pandas as ps вместо import pandas as pd, и всё. Но за этим синтаксическим сахаром скрывается принципиально другой механизм выполнения: вместо операций над in-memory массивом на одной машине запускается распределённый Spark job на кластере из N нод.

Доступен начиная с PySpark 3.2 как встроенная библиотека (ранее устанавливался отдельно: pip install koalas). До PySpark 3.2 нужно было: pip install koalas + import databricks.koalas as ks.


От Koalas к официальному ядру: история проекта

Понимание истории pyspark.pandas помогает объяснить некоторые архитектурные решения и ограничения, которые встречаются при работе с ним сегодня.

2019 - рождение Koalas. Databricks выпускает open-source проект Koalas - pandas-совместимый DataFrame API поверх Apache Spark. Проблема была конкретной: в крупных компаниях (Uber, Airbnb, Netflix) работали тысячи аналитиков с pandas, но данных стало так много, что однонодовая обработка превратилась в узкое место. Переучить всех на нативный PySpark - непосильная задача. Koalas предложил постепенную миграцию.

2021 - передача в Apache. Databricks передаёт Koalas в upstream Apache Spark. API перестаёт быть проприетарным дополнением и становится частью экосистемы.

2022 - PySpark 3.2. Koalas интегрируется напрямую в дистрибутив Apache Spark как pyspark.pandas. Отдельная установка больше не нужна. API стабилизируется, покрытие pandas-методов расширяется до ~80%.

Сегодня. pyspark.pandas активно развивается. Однако фундаментальные ограничения - прежде всего проблема индексов и неявных коллектов - никуда не делись, потому что они обусловлены самой природой распределённых вычислений, а не недоработками реализации.


Смена парадигмы: от однонодового к распределённому

Чтобы понять ограничения pyspark.pandas, нужно осознать принципиальную разницу между вычислительными моделями pandas и Spark.

Pandas - eager, imperative, single-node. Каждая операция над pandas DataFrame выполняется немедленно (eager execution), изменяет состояние объекта in-place или создаёт новый объект. Индекс - это фундамент: каждая строка имеет уникальный номер, по которому можно обратиться мгновенно. Порядок строк гарантирован. Всё в памяти одной машины.

Spark - lazy, functional, distributed. Каждая трансформация строит логический план (DAG), не выполняя вычислений. Данные распределены по партициям на разных нодах. Понятия «строка №5» не существует - партиция может лежать на любой ноде, строки внутри неё никак не упорядочены относительно других партиций. Выполнение начинается только при вызове Action.

Задача pyspark.pandas - транслировать привычный императивный синтаксис pandas в ленивые распределённые операции Catalyst. Это работает, но не без потерь:

  • In-place мутации (df["col"] = value) транслируются в создание нового DataFrame
  • Операции, требующие глобального порядка (sort_values, iloc), вызывают дорогостоящий Shuffle
  • Методы, требующие полного обхода данных (to_dict, iterrows), тайно собирают всё на Driver
# В pandas: мгновенная мутация, O(1) по памяти
import pandas as pd
df = pd.read_csv("data.csv")
df["tax"] = df["price"] * 0.2     # мутация на месте
df.iloc[5]                         # мгновенный доступ по индексу

# В pyspark.pandas: тот же синтаксис, но под капотом - новый DataFrame
import pyspark.pandas as ps
df = ps.read_csv("s3://bucket/data.csv")
df["tax"] = df["price"] * 0.2     # создаётся НОВЫЙ DataFrame, перестраивается план
df.iloc[5]                         # вызывает полный Shuffle для построения индекса!

Код выглядит одинаково, но поведение принципиально разное. Понимание этой разницы - ключ к правильному использованию pyspark.pandas.


Архитектура: Arrow как мост

Под капотом pyspark.pandas конвертирует pandas-вызовы в Spark DataFrame операции, используя Apache Arrow для эффективной передачи данных между Python и JVM:

Apache Arrow - это columnar in-memory формат, разработанный специально для эффективного обмена данными между системами. В pyspark.pandas Arrow выполняет роль универсального переходника:

  • JVM → Python без сериализации. Данные из Spark JVM передаются в Python без конвертации через pickle. Arrow batches читаются напрямую как zero-copy memory-mapped буферы.
  • Векторизованная обработка. Python получает данные не строками, а колонками (батчами), что позволяет NumPy-совместимым операциям работать эффективно.
  • Двунаправленный обмен. При df.to_pandas() - из Spark в Python через Arrow; при ps.from_pandas() - обратно.

Включить Arrow (по умолчанию включён в новых версиях, но лучше проверить явно):

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

Без Arrow каждый элемент сериализовался бы через pickle - в 3–7 раз медленнее. С Arrow колонка из 1M числовых значений передаётся как один непрерывный буфер памяти размером ~8 МБ (8 байт × 1M).


Три модели выполнения: pandas, pyspark.pandas, нативный PySpark

Чтобы принимать осознанные архитектурные решения, нужно чётко понимать, чем эти три модели отличаются:

Ключевые наблюдения из диаграммы:

pandas - вся работа происходит на одной ноде в памяти. Данных больше RAM - OutOfMemoryError. Преимущество: полный контроль над состоянием, индексы работают мгновенно, порядок строк гарантирован.

pyspark.pandas - добавляет промежуточный слой трансляции между pandas-синтаксисом и Spark. Этот слой несёт накладные расходы: каждый вызов метода должен быть проанализирован, транслирован и передан в Spark. Однако преимущество Catalyst-оптимизации сохраняется - транслированный план оптимизируется так же, как нативный PySpark план.

Нативный PySpark - самый прямой путь к движку. Нет промежуточного слоя, нет накладных расходов на трансляцию. Catalyst оптимизирует план максимально эффективно. Недостаток: непривычный синтаксис для тех, кто пришёл из pandas.


Lazy evaluation под pandas-синтаксисом

Одна из самых непривычных особенностей pyspark.pandas - сохранение lazy evaluation несмотря на то, что pandas-код выглядит eagerly.

В обычном pandas каждая строка выполняется немедленно и возвращает результат. В pyspark.pandas - нет:

import pyspark.pandas as ps

df = ps.read_parquet("s3://bucket/events/")  # Action - Spark читает метаданные

# Это всё трансформации - план строится, вычислений нет
filtered  = df[df["amount"] > 100]
enriched  = filtered.copy()
enriched["category"]   = enriched["type"].str.upper()
enriched["amount_tax"] = enriched["amount"] * 1.2

# Action - только здесь Spark реально обрабатывает данные
result = enriched.groupby("category")["amount_tax"].sum()
print(result)   # ещё один Action (вывод требует данных)

Это отличается от поведения обычного pandas, где каждая операция на строчке filtered = df[...] немедленно создаёт новый DataFrame в памяти. В pyspark.pandas она лишь добавляет узел в логический план.

Практическое следствие: промежуточные результаты не кешируются автоматически. Если использовать enriched в двух разных вычислениях, план выполнится дважды:

# Плохо: enriched вычисляется дважды
total = enriched["amount"].sum()   # Action 1: полный scan
count = enriched["amount"].count() # Action 2: ещё один полный scan

# Хорошо: пересчитать один раз через нативный Spark
spark_df  = enriched.to_spark().cache()
ps_cached = spark_df.pandas_api()
total     = ps_cached["amount"].sum()   # Action 1: из кеша
count     = ps_cached["amount"].count() # Action 2: из кеша

Проклятие индексов - The Index Problem

Это самая серьёзная ловушка pyspark.pandas, которая неожиданно обваливает производительность. Чтобы понять проблему, нужно разобраться, зачем индекс нужен pandas и почему Spark не может его поддержать так же эффективно.

Зачем индекс в pandas? Индекс pandas - это структура данных, которая хранит уникальную метку для каждой строки. По умолчанию - монотонный целочисленный RangeIndex: 0, 1, 2, 3, ... Это позволяет за O(1) обращаться к любой строке по позиции (df.iloc[5]), выравнивать два DataFrame при операциях (df1 + df2 выровняет по индексу), делать label-based slicing (df.loc[10:20]).

Почему Spark не может это поддержать? В Spark данные распределены по партициям на разных нодах. Партиция №3 может лежать на ноде 5, а партиция №1 - на ноде 2. Строки не имеют глобального порядка. Понятия "строка №1000 от начала датасета" не существует без глобального сканирования.

Когда pyspark.pandas пытается построить последовательный монотонный индекс, Spark вынужден выполнить глобальный sort + shuffle - одну из самых дорогих операций в распределённых вычислениях:

Три типа индекса - управляются через compute.default_index_type:

Тип Поведение Стоимость Когда использовать
sequence Глобальный монотонный 0,1,2... Очень дорого: полный Shuffle Никогда в production
distributed-sequence Монотонный по партициям, уникальный глобально Дорого: частичный sort Только если порядок реально нужен
distributed monotonically_increasing_id() Почти бесплатно По умолчанию для production
import pyspark.pandas as ps

# Рекомендуется явно установить в начале сессии:
ps.set_option("compute.default_index_type", "distributed")

# Теперь создание DF не будет триггерить дорогой Shuffle
df = ps.read_parquet("s3://bucket/events/")

# Если индекс не нужен - явно отбросить при операциях
result = df.groupby("user_id")["amount"].sum().reset_index(drop=True)

Практический урок: большинство Data Engineering пайплайнов вообще не нуждаются в индексе. Фильтрация, группировка, агрегации, join - всё это работает без позиционного доступа к строкам. Установите distributed и забудьте о проблеме индексов.

Когда индекс реально нужен? Только в специфических аналитических задачах: time-series resample с выравниванием по времени, merge по индексу между двумя DF. В таких случаях смиритесь со стоимостью или переписывайте на нативный PySpark.


Неявный сбор данных на Driver (Implicit Collect)

Вторая серьёзная ловушка - операции, которые тайно собирают все данные на Driver. Это происходит, когда Python-метод не может быть транслирован в распределённую операцию Spark и вынужден подтянуть всё к себе.

Опасные операции и почему они небезопасны:

import pyspark.pandas as ps
df = ps.read_parquet("s3://bucket/events/")  # 50 GB данных

# ОПАСНО: весь датасет едет в память Driver
records = df.to_dict()           # OutOfMemoryError на 50 GB
values  = df.values              # то же - возвращает numpy array
rows    = list(df.iterrows())    # то же - итерирует построчно

# ОПАСНО: iloc с отрицательным индексом требует знания общей длины
last_row = df.iloc[-1]           # Spark вынужден собрать всё для нумерации

# БЕЗОПАСНО: остаётся распределённым
filtered   = df[df["amount"] > 1000]                          # Spark Filter
aggregated = df.groupby("user_id")["amount"].sum()            # Spark GroupBy
spark_df   = df.to_spark()                                    # без перемещения данных

Особый случай - apply(): в pandas df.apply(func) применяет Python-функцию к каждой строке или колонке. В pyspark.pandas поведение зависит от контекста - Spark попытается использовать pandas_udf, но может неявно собрать данные:

# Потенциально опасно - поведение зависит от функции
df.apply(lambda row: custom_python_func(row), axis=1)

# Безопаснее и явно быстрее - нативный UDF
from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType

@udf(returnType=DoubleType())
def custom_spark_udf(value):
    return custom_python_func(value)

df.to_spark().withColumn("result", custom_spark_udf("amount"))

In-place операции и перестройка плана

В pandas df["new_col"] = ... - это настоящая мутация: объект df изменяется на месте, в памяти создаётся новый массив для колонки, план не меняется.

В pyspark.pandas всё иначе. Spark DataFrames иммутабельны - их нельзя изменить. Каждая «мутация» через df["new_col"] = ... создаёт новый DataFrame с расширенным планом выполнения:

import pyspark.pandas as ps
df = ps.read_parquet("s3://bucket/events/")

# Каждое присваивание создаёт новый DataFrame и перестраивает план
df["tax"]      = df["amount"] * 0.2        # DataFrame #1, план глубиной 1
df["net"]      = df["amount"] - df["tax"]  # DataFrame #2, план глубиной 2
df["category"] = df["type"].str.upper()    # DataFrame #3, план глубиной 3

# После многих мутаций план становится очень глубоким
# При explain() увидите длинную цепочку Project→Filter→Project→...
print(df.to_spark().explain())

Это особенно критично в циклах:

# ПЛОХО: 100 итераций = 100 вложенных планов = потенциальный StackOverflowError
for feature_name in feature_list:
    df[feature_name] = compute_feature(df, feature_name)

# ХОРОШО: накопить все вычисления через нативный Spark за один раз
from pyspark.sql import functions as F

spark_df = df.to_spark()
for feature_name, feature_expr in feature_expressions.items():
    spark_df = spark_df.withColumn(feature_name, feature_expr)
# Catalyst оптимизирует всё в один физический план

Где pyspark.pandas ускоряет миграцию

Несмотря на ловушки, pyspark.pandas действительно полезен в ряде сценариев:

1. Zero-Rewrite Migration. Команда аналитиков написала 2000 строк pandas-кода для feature engineering. Данных стало в 50 раз больше. Замена import pandas as pd на import pyspark.pandas as ps + исправление нескольких несовместимых мест (обычно 5–10% кода) - это неделя работы против 3–4 месяцев полного переписывания.

2. Распределённый EDA. Исследовательский анализ данных на терабайтах с привычными методами:

import pyspark.pandas as ps
ps.set_option("compute.default_index_type", "distributed")

df = ps.read_parquet("s3://datalake/events/year=2024/")

# Знакомые методы работают на терабайтах
print(df.describe())                 # статистика по всем колонкам
print(df["revenue"].value_counts())  # топ значений
print(df.isnull().sum())             # количество пропусков по колонкам

Без pyspark.pandas такой EDA требовал либо полного переписывания на Spark SQL/PySpark, либо работы на sample данных.

3. Интеграция с ML-экосистемой. Feature engineering перед обучением моделей:

import pyspark.pandas as ps
from pyspark.ml.feature import VectorAssembler

df = ps.read_parquet("s3://datalake/features/")
df["log_amount"] = ps.np.log1p(df["amount"])
df["amount_z"]   = (df["amount"] - df["amount"].mean()) / df["amount"].std()
df["day_of_week"] = df["date"].dt.dayofweek

# Конвертация в нативный Spark для обучения модели
spark_df  = df.to_spark()
assembler = VectorAssembler(
    inputCols=["log_amount", "amount_z", "day_of_week"],
    outputCol="features"
)
ml_df = assembler.transform(spark_df)

4. Прототипирование с горизонтом переписывания. Быстро запустить код на большом кластере, проверить бизнес-логику, потом оптимизировать критические части на нативном PySpark. "Сначала работает, потом быстро."


Основные операции

import pyspark.pandas as ps

# Чтение
df = ps.read_parquet("s3://bucket/events/")
df = ps.read_csv("data.csv")
df = ps.from_pandas(pd_df)       # из обычного pandas DF (данные идут на кластер)

# DataFrame операции (pandas-синтаксис)
df_filtered = df[df["amount"] > 100]
df["category"] = df["type"].str.upper()
df["rank"] = df.groupby("user_id")["amount"].rank(method="dense")

# Конвертация
pd_df    = df.to_pandas()        # собрать на Driver - осторожно с объёмом!
spark_df = df.to_spark()         # вернуться в нативный Spark DataFrame (без перемещения)
ps_df    = spark_df.pandas_api() # Spark DF в pyspark.pandas DF (без перемещения)

Разница между to_pandas() и to_spark() принципиальна:

  • to_pandas() - это Action. Все данные едут на Driver через Arrow. Опасно для больших DF.
  • to_spark() - это трансформация. Никаких перемещений данных. Просто меняется тип обёртки с pyspark.pandas.DataFrame на pyspark.sql.DataFrame. Бесплатно.
  • pandas_api() - обратная операция к to_spark(). Тоже бесплатна.
# Безопасный паттерн: работать в ps, при необходимости переключаться в Spark
ps_df    = ps.read_parquet("s3://bucket/events/")
spark_df = ps_df.to_spark()             # бесплатно: просто меняем обёртку
cached   = spark_df.cache()             # кешируем нативными средствами Spark
ps_cached = cached.pandas_api()         # возвращаемся в ps для дальнейших операций

total_per_user = ps_cached.groupby("user_id")["amount"].sum()

Обход несовместимостей

pyspark.pandas не на 100% совместим с pandas - часть операций ведёт себя иначе или не поддерживается:

# Индекс по умолчанию - распределённый (нет гарантии порядка)
# pandas: целочисленный монотонный индекс
# pyspark.pandas: по умолчанию sequence (глобальный sort) или distributed (без порядка)

# Порядок строк не гарантирован
df.sort_values("date")  # OK, но не меняет физический порядок партиций

# Некоторые методы требуют сбора на Driver
df.to_dict()   # Опасно! Весь DF уходит в память Driver
df.values      # То же

# Явно включить режим совместимости (замедляет - разрешает join разных DF)
ps.set_option("compute.ops_on_diff_frames", True)

ops_on_diff_frames - особый случай. В pandas можно выполнять операции между двумя DF с разными индексами:

# В pandas это работает: автовыравнивание по индексу
pd_df1 = pd.DataFrame({"a": [1, 2, 3]})
pd_df2 = pd.DataFrame({"b": [4, 5, 6]})
result = pd_df1["a"] + pd_df2["b"]   # [5, 7, 9] - выравнивание по индексу

# В pyspark.pandas по умолчанию это ошибка - нет общего индекса
ps_df1 = ps.DataFrame({"a": [1, 2, 3]})
ps_df2 = ps.DataFrame({"b": [4, 5, 6]})
# result = ps_df1["a"] + ps_df2["b"]   # ValueError!

# Включить (медленно - требует join по индексу = Shuffle)
ps.set_option("compute.ops_on_diff_frames", True)
result = ps_df1["a"] + ps_df2["b"]   # работает, но дорого

# Лучше: явный join на нативном Spark
spark_result = ps_df1.to_spark().join(ps_df2.to_spark(), on="key")

Сравнение API: нативный PySpark vs Pandas API on Spark

Понимание синтаксических различий помогает правильно конвертировать код в обоих направлениях.

Фильтрация:

# Нативный PySpark
from pyspark.sql.functions import col
spark_df.filter(col("amount") > 100)
spark_df.where("amount > 100")

# pyspark.pandas
ps_df[ps_df["amount"] > 100]
ps_df.query("amount > 100")   # тоже работает, транслируется в WHERE

Добавление новых колонок:

# Нативный PySpark
from pyspark.sql.functions import upper
spark_df.withColumn("tax", col("amount") * 0.2)
spark_df.withColumn("category", upper(col("type")))

# pyspark.pandas (создаёт новый DataFrame под капотом)
ps_df["tax"]      = ps_df["amount"] * 0.2
ps_df["category"] = ps_df["type"].str.upper()

Агрегации:

# Нативный PySpark
from pyspark.sql.functions import sum as spark_sum, avg, count, max as spark_max

spark_df.groupBy("user_id").agg(
    spark_sum("amount").alias("total"),
    avg("amount").alias("avg_amount"),
    count("*").alias("cnt"),
    spark_max("amount").alias("max_amount")
)

# pyspark.pandas - лаконичнее для простых случаев
ps_df.groupby("user_id")["amount"].agg(["sum", "mean", "count", "max"])
# Или словарь:
ps_df.groupby("user_id").agg({"amount": ["sum", "mean"], "events": "count"})

Join:

# Нативный PySpark
orders.join(users, orders["user_id"] == users["id"], "left")

# pyspark.pandas
orders_ps.merge(users_ps, left_on="user_id", right_on="id", how="left")

Ключевой совет по выбору синтаксиса: для простых агрегаций и фильтраций pyspark.pandas синтаксис удобнее. Для сложных window functions, lateral join, оконных функций с несколькими partitionBy - переключайтесь на нативный PySpark через to_spark().


GroupBy и apply: скрытые ловушки

groupby().apply() - один из самых мощных инструментов pandas, но в pyspark.pandas он имеет существенные ограничения.

Как работает groupby().apply() в pandas: принимает функцию f(group: pd.DataFrame) -> pd.DataFrame. Применяет к каждой группе. Быстро, потому что всё в памяти одной машины.

Как это работает в Spark: каждая группа сериализуется, отправляется на Python-воркер, выполняется Python-функцией, результат сериализуется обратно. Это дорого из-за сериализации и Python overhead.

# pandas-style groupby apply - работает, но не оптимально
def normalize_group(group):
    group["normalized"] = (group["amount"] - group["amount"].mean()) / group["amount"].std()
    return group

# Транслируется в pandas_udf под капотом (если можно)
result = ps_df.groupby("category").apply(normalize_group)

Проблемы с apply():

  • Схема результата должна совпадать со схемой входного DF или быть явно указана
  • Python GIL: каждый Python-воркер выполняет функцию последовательно для своей партиции
  • Крупные группы: если группа не помещается в памяти воркера - OOM
  • Нет pushdown оптимизаций: Catalyst не может оптимизировать внутри Python-функции

Когда предпочесть нативный Spark:

# Вместо groupby().apply() с Python-функцией:
ps_df.groupby("category").apply(normalize_group)  # медленно

# Используйте Window functions через нативный Spark:
from pyspark.sql import functions as F
from pyspark.sql.window import Window

spark_df = ps_df.to_spark()
w = Window.partitionBy("category")
spark_df.withColumn(
    "normalized",
    (F.col("amount") - F.avg("amount").over(w)) / F.stddev("amount").over(w)
)
# Catalyst оптимизирует, выполняется в JVM без Python overhead

Pandas UDF vs Pandas API on Spark

Это две разные технологии с разными применениями, хотя обе используют Apache Arrow и обе работают с pandas объектами в Python:

Pandas UDF (@pandas_udf) - это способ добавить кастомную Python-логику в нативный PySpark pipeline. Принимает pd.Series батч, возвращает pd.Series. Используется для вычисления новой колонки по существующей.

from pyspark.sql.functions import pandas_udf
from pyspark.sql.types import DoubleType
import pandas as pd

# Pandas UDF: функция над батчем значений одной колонки
@pandas_udf(DoubleType())
def risk_score(amounts: pd.Series) -> pd.Series:
    return amounts.apply(lambda x: x * 0.1 if x > 1000 else x * 0.05)

# Применяется в нативном PySpark
spark_df.withColumn("risk", risk_score("amount"))

Pandas API on Spark - полноценный DataFrame интерфейс для работы с целым датасетом в pandas-стиле. Не ограничен одной колонкой, поддерживает groupby, merge, join и другие межколоночные операции.

import pyspark.pandas as ps

ps_df = ps.read_parquet("s3://bucket/events/")
# Работает со всем DataFrame, не только с одной колонкой
result = ps_df.groupby("user_id").agg({
    "amount": ["sum", "mean", "std"],
    "date":   ["min", "max"]
})

Когда что использовать:

Сценарий Что выбрать
Кастомная Python-логика для одной колонки @pandas_udf
Сложная бизнес-логика с numpy/scipy @pandas_udf
Полноценный EDA над большим датасетом pyspark.pandas
Миграция pandas notebook на кластер pyspark.pandas
Feature engineering с groupby/window pyspark.pandasto_spark() для сложных случаев
Production ETL pipeline Нативный PySpark

NULL и NaN: разные семантики

Это тонкая, но важная разница, которая может приводить к трудноуловимым багам при миграции.

В pandas есть два понятия "отсутствующего значения":

  • NaN (Not a Number) - float-значение из IEEE 754. Может жить в числовых и object колонках.
  • None - Python объект. Для object-колонок.
  • pandas часто смешивает их: pd.isna(None) → True, pd.isna(float('nan')) → True.

В Spark есть только NULL - это SQL NULL семантика. Для float есть NaN как отдельное значение, которое isNull() не поймает:

from pyspark.sql import functions as F

# Spark: NULL и NaN - разные вещи!
spark_df = spark.createDataFrame([
    (1, None),           # NULL
    (2, float("nan")),   # NaN - не NULL!
    (3, 42.0)
], ["id", "value"])

spark_df.filter(F.col("value").isNull()).show()
# Результат: только строка с None (id=1)
# NaN (id=2) НЕ считается NULL!

spark_df.filter(F.isnan("value")).show()
# Результат: только строка с NaN (id=2)

# Правильная обработка: поймать и то и другое
spark_df.filter(
    F.col("value").isNull() | F.isnan("value")
).show()

В pyspark.pandas поведение приближено к pandas - isnull() ловит и NULL и NaN. Но при конвертации через to_spark() исходная семантика Spark сохраняется:

import pyspark.pandas as ps
import numpy as np

ps_df = ps.DataFrame({
    "value": [1.0, None, np.nan, 4.0]
})

# pyspark.pandas: pandas-семантика
print(ps_df["value"].isnull().sum())   # 2 (None и nan оба считаются)

# Нормализовать NaN в NULL перед to_spark():
ps_df["value"] = ps_df["value"].fillna(None)  # все nan → None → NULL в Spark

# Нормализовать NaN в нативном Spark:
spark_df = spark_df.withColumn(
    "value",
    F.when(F.isnan("value"), None).otherwise(F.col("value"))
)

Оконные функции в Pandas API

Window functions - аналитические операции, которые вычисляют значения относительно группы строк. В pandas они реализованы через groupby().transform(), rolling(), expanding(). В pyspark.pandas часть из них поддерживается:

import pyspark.pandas as ps

df = ps.read_parquet("s3://bucket/orders/")

# Накопительная сумма по пользователю (работает)
df["cumulative_amount"] = df.groupby("user_id")["amount"].cumsum()

# Ранжирование (работает)
df["rank"] = df.groupby("user_id")["amount"].rank(method="dense", ascending=False)

# Процент от суммы группы (работает)
group_total = df.groupby("category")["amount"].transform("sum")
df["pct_of_category"] = df["amount"] / group_total

Однако для сложных window frames (ROWS BETWEEN ... AND ...) нативный Spark Window API даёт больше возможностей:

from pyspark.sql.window import Window
from pyspark.sql import functions as F

# Нативный Spark: lag/lead с гибким окном
spark_df = df.to_spark()
w = Window.partitionBy("user_id").orderBy("date")
spark_df = spark_df.withColumn("prev_amount", F.lag("amount", 1).over(w))
spark_df = spark_df.withColumn("next_amount", F.lead("amount", 1).over(w))
spark_df = spark_df.withColumn(
    "rolling_sum_3",
    F.sum("amount").over(w.rowsBetween(-2, 0))  # скользящее окно в 3 строки
)

# Вернуться в pyspark.pandas при желании
result_ps = spark_df.pandas_api()

Анализ планов выполнения (explain)

Один из ключевых инструментов для понимания проблем производительности - план выполнения Spark. pyspark.pandas позволяет получить его через to_spark().explain():

import pyspark.pandas as ps

# Намеренно медленный тип - для демонстрации
ps.set_option("compute.default_index_type", "sequence")

df = ps.read_parquet("s3://bucket/events/")
result = df.groupby("user_id")["amount"].sum()

# Получить план выполнения
result.to_spark().explain(mode="extended")

Пример вывода с sequence индексом - обратите внимание на лишние Exchange и Sort:

== Physical Plan ==
*(4) Sort [__index_level_0__ ASC NULLS FIRST], true, 0
+- Exchange rangepartitioning(__index_level_0__ ASC NULLS FIRST, 200)
   +- *(3) HashAggregate(keys=[user_id], functions=[sum(amount)])
      +- Exchange hashpartitioning(user_id, 200)
         +- *(2) HashAggregate(keys=[user_id], functions=[partial_sum(amount)])
            +- *(1) ColumnarToRow
               +- Scan parquet s3://bucket/events/

Видны два Exchange (Shuffle) - один для GroupBy, второй для сортировки индекса. При distributed индексе второй Exchange исчезает:

== Physical Plan ==
*(3) HashAggregate(keys=[user_id], functions=[sum(amount)])
+- Exchange hashpartitioning(user_id, 200)
   +- *(2) HashAggregate(keys=[user_id], functions=[partial_sum(amount)])
      +- *(1) ColumnarToRow
         +- Scan parquet s3://bucket/events/

Только один Shuffle - для группировки. Это и есть то, что даёт distributed индекс.

Что искать в explain():

  • Exchange - это Shuffle. Каждый Shuffle = сетевые расходы. Минимизировать.
  • Sort - дорогая сортировка. Нужна ли она реально?
  • BroadcastExchange - малая таблица разослана всем исполнителям. Хорошо.
  • Filter перед Scan - pushdown работает. Хорошо.
  • Filter после Scan - pushdown не применился, весь файл читается. Плохо для Parquet.

Антипаттерны производительности

Собранные в одном месте самые распространённые ошибки при работе с pyspark.pandas:

1. Повторное вычисление без кеширования:

# ПЛОХО: df_enriched сканируется дважды
total = df_enriched["amount"].sum()
count = df_enriched["amount"].count()

# ХОРОШО:
df_cached = df_enriched.to_spark().cache().pandas_api()
total = df_cached["amount"].sum()
count = df_cached["amount"].count()
df_cached.to_spark().unpersist()  # освободить кеш

2. to_pandas() на большом датасете:

# ПЛОХО: 50 GB → Driver RAM → OOM
pd_result = large_ps_df.to_pandas()

# ХОРОШО: агрегировать сначала, потом конвертировать
summary  = large_ps_df.groupby("category")["amount"].sum()
pd_summary = summary.to_pandas()   # маленький агрегат - безопасно

3. sequence index в production:

# ПЛОХО: каждое создание DF с sequence index = Shuffle
ps.set_option("compute.default_index_type", "sequence")

# ХОРОШО:
ps.set_option("compute.default_index_type", "distributed")

4. apply() с тяжёлой Python-функцией без pandas_udf:

# ПЛОХО: нет векторизации, медленная сериализация строк
df["score"] = df["text"].apply(lambda x: ml_model.predict(x))

# ХОРОШО: pandas_udf для векторизованной обработки батчами
@pandas_udf(DoubleType())
def batch_predict(texts: pd.Series) -> pd.Series:
    return pd.Series(ml_model.predict(texts.tolist()))

spark_df.withColumn("score", batch_predict("text"))

5. Мутации в цикле:

# ПЛОХО: глубокий план, потенциальный StackOverflow при 100+ итерациях
for col in many_columns:
    ps_df[col] = compute(ps_df, col)

# ХОРОШО: выйти в нативный Spark
spark_df = ps_df.to_spark()
for col_name, expr in column_expressions.items():
    spark_df = spark_df.withColumn(col_name, expr)
result_ps = spark_df.pandas_api()

6. ops_on_diff_frames без необходимости включено глобально:

# ПЛОХО: включено глобально, все операции замедляются
ps.set_option("compute.ops_on_diff_frames", True)

# ХОРОШО: только там где нужно, или через нативный join
ps.set_option("compute.ops_on_diff_frames", False)  # по умолчанию
joined = ps_df1.merge(ps_df2, on="key")  # явный join

pyspark.pandas vs нативный PySpark API

pyspark.pandas Нативный PySpark DataFrame
Целевая аудитория Миграция pandas-кода Новая разработка под Spark
Синтаксис pandas-совместимый Spark-специфичный
Производительность Хуже (дополнительный слой) Лучше
Поддержка операций ~80% pandas API 100% Spark API
Отладка Сложнее Проще (Spark UI, explain())
Index semantics Эмулируется (дорого) Нет (данные без порядка)
Рекомендация Миграция легаси, EDA, прототипы Все новые пайплайны

Benchmark: pandas vs pyspark.pandas vs нативный PySpark

Реальные числа сильно зависят от кластера, данных и операций, но порядок величин полезно понимать. Пример для датасета 100M строк, операция: groupBy + agg (sum, mean, count) по 5 колонкам:

Реализация Время Память Driver Примечание
pandas (single node) OOM crash >100 GB Не помещается в RAM
pyspark.pandas, sequence ~185s <1 GB 2 лишних Shuffle на индекс
pyspark.pandas, distributed ~48s <1 GB 1 Shuffle (только для groupBy)
Нативный PySpark ~31s <1 GB Оптимально: нет overhead на трансляцию

Выводы:

  • pandas вообще не справляется с такими объёмами
  • pyspark.pandas с sequence индексом в 6x медленнее нативного PySpark
  • pyspark.pandas с distributed индексом в 1.5x медленнее нативного PySpark
  • Для большинства задач разница между distributed pyspark.pandas и нативным PySpark приемлема

Небольшие датасеты (1M строк):

Реализация Время Когда использовать
pandas ~0.1s Данные помещаются в RAM
pyspark.pandas ~3–6s Overhead на Spark инициализацию
Нативный PySpark ~1–3s Чуть быстрее pyspark.pandas для малых данных

Для малых датасетов pandas значительно быстрее обоих вариантов Spark. Spark - это инструмент для больших данных, накладные расходы на распределение не окупаются на малых объёмах.


Pandas API и Lakehouse-стек

Современные data platform строятся на Iceberg/Delta Lake таблицах в S3. pyspark.pandas корректно работает с ними через обычные механизмы Spark:

import pyspark.pandas as ps

# Чтение из Iceberg таблицы через Spark catalog
spark_df = spark.read.table("iceberg.events.orders")
ps_df    = spark_df.pandas_api()

# Аналитика в pandas-стиле
monthly_summary = (
    ps_df
    .assign(month=ps_df["created_at"].dt.to_period("M"))
    .groupby(["month", "category"])["revenue"]
    .agg(["sum", "count", "mean"])
    .reset_index()
)

# Запись обратно в Iceberg
monthly_summary.to_spark().writeTo("iceberg.analytics.monthly_summary") \
    .tableProperty("write.format.default", "parquet") \
    .createOrReplace()

Важный момент про схему. Iceberg и Delta Lake обеспечивают schema enforcement на уровне таблицы. Если pyspark.pandas попытается записать DataFrame с другой схемой - получит ошибку:

# Iceberg отклонит запись если типы не совпадают
ps_df["amount"] = ps_df["amount"].astype(str)   # конвертировали в string
# ps_df.to_spark().write.saveAsTable("iceberg.events.orders")
# AnalysisException: Cannot write incompatible data to table

Это поведение правильное - Iceberg защищает целостность данных. При записи с партиционированием лучше использовать нативный Spark API:

# Запись с контролем партиционирования - через нативный Spark
ps_df.to_spark() \
    .write \
    .partitionBy("year", "month") \
    .mode("overwrite") \
    .parquet("s3://bucket/output/")

Когда использовать

Хорошо:

  • Перенос существующего pandas-кода на кластер с минимальными изменениями
  • Data Science notebooks с привычным pandas-синтаксисом
  • Прототипирование с последующим переводом на нативный API
  • EDA (Exploratory Data Analysis) на больших датасетах
  • Аналитические запросы в интерактивных сессиях (Jupyter, Databricks)
  • Feature engineering для ML, если команда состоит в основном из DS, а не DE

Плохо:

  • Production ETL-пайплайны - накладные расходы слоя трансляции и риск скрытых проблем
  • Операции, требующие гарантированного порядка строк
  • Код, активно использующий pandas-индексы (в Spark нет аналога)
  • Пайплайны, написанные с нуля - нет смысла добавлять слой совместимости
  • Критичные по производительности пути: сложные join, оконные функции с большим окном

Золотое правило: если код уже написан на pandas - pyspark.pandas поможет перенести его на Spark с минимальными усилиями. Если пишется с нуля - сразу нативный PySpark.


Типичный путь миграции

Пошаговый процесс безопасного переноса pandas-кода на Spark:

import pyspark.pandas as ps

# Стартовая конфигурация для миграции
ps.set_option("compute.default_index_type", "distributed")
ps.set_option("compute.ops_on_diff_frames", False)
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")

# Было (pandas):
# import pandas as pd
# df = pd.read_csv("data.csv")

# Стало (pyspark.pandas):
df = ps.read_csv("s3://bucket/data.csv")  # или ps.read_parquet(...)

# Весь остальной pandas-код остаётся почти без изменений
df_clean    = df.dropna(subset=["amount", "user_id"])
df_filtered = df_clean[df_clean["amount"] > 0]
result      = df_filtered.groupby("category")["amount"].agg(["sum", "count", "mean"])

# Узкое место - переписать на нативный Spark
complex_result = df_filtered.to_spark() \
    .join(reference_df, "product_id", "left") \
    .groupBy("category", "subcategory") \
    .agg(F.sum("amount").alias("total"))

Практика: рефакторинг ML-пайплайна

Исходный код дата-сайентиста (имитирует реальный notebook с типичными проблемами):

import pyspark.pandas as ps

# Проблема 1: sequence индекс по умолчанию - лишние Shuffle
df = ps.read_parquet("s3://datalake/user_events/")

# Проблема 2: in-place мутации в цикле создают глубокий план
feature_cols = ["page_views", "session_time", "clicks", "purchases", "cart_adds"]
for col in feature_cols:
    df[col + "_log"] = ps.np.log1p(df[col])  # каждая итерация = новый DataFrame

# Проблема 3: z-score через apply = Python overhead на каждую строку
for col in feature_cols:
    mean = df[col].mean()
    std  = df[col].std()
    df[col + "_z"] = df[col].apply(lambda x: (x - mean) / std if std > 0 else 0)

# Проблема 4: to_pandas() на большом DF перед записью
pd_result = df.to_pandas()       # OOM риск
pd_result.to_parquet("output.parquet")

Оптимизированная версия:

import pyspark.pandas as ps
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# Исправление 1: distributed индекс
ps.set_option("compute.default_index_type", "distributed")
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")

df = ps.read_parquet("s3://datalake/user_events/")

feature_cols = ["page_views", "session_time", "clicks", "purchases", "cart_adds"]
spark_df = df.to_spark()

# Исправление 2: log-трансформации за один выход в нативный Spark
for col in feature_cols:
    spark_df = spark_df.withColumn(f"{col}_log", F.log1p(F.col(col)))

# Исправление 3: z-score через Window (JVM, без Python overhead)
for col in feature_cols:
    w = Window.rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing)
    mean_val = F.avg(col).over(w)
    std_val  = F.stddev(col).over(w)
    spark_df = spark_df.withColumn(
        f"{col}_z",
        F.when(std_val > 0, (F.col(col) - mean_val) / std_val).otherwise(F.lit(0.0))
    )

# Исправление 4: записать через нативный Spark без to_pandas()
spark_df.write \
    .mode("overwrite") \
    .parquet("s3://datalake/user_features/")

Ожидаемое улучшение на 100M строк:

  • До оптимизации: ~12 минут (sequence index + Python apply в цикле + to_pandas)
  • После оптимизации: ~2.5 минуты (distributed index + JVM Window + нативная запись)

Стратегический чек-лист архитектора

Перед тем как использовать pyspark.pandas в production, проверьте:

Конфигурация:

  • compute.default_index_type установлен в distributed (не sequence)
  • spark.sql.execution.arrow.pyspark.enabled = true включён
  • compute.ops_on_diff_frames = False (включать только явно там где нужно)

Код-ревью:

  • Нет to_dict(), df.values, iterrows() на больших DF
  • Нет iloc[-1] или iloc с отрицательными индексами на больших DF
  • Нет мутаций в цикле с 100+ итерациями без выхода в нативный Spark
  • Нет apply() с тяжёлой логикой - заменить на pandas_udf
  • Нет to_pandas() без предварительной агрегации

Производительность:

  • Проверить план через to_spark().explain() - нет ли лишних Exchange (Shuffle)
  • Проверить Spark UI: Stage времена, Shuffle Write/Read байты
  • Добавить cache() если ps_df используется в нескольких Action

Архитектурные принципы:

  • pyspark.pandas - только для миграции и EDA. Новые пайплайны - нативный PySpark.
  • Критичные для производительности операции (сложные join, оконные функции) - всегда на нативном Spark.
  • Используйте to_spark() как "аварийный выход" для любой операции, которая работает медленно.

Pandas API on Spark упрощает жизнь разработчику, но усложняет её оптимизатору Catalyst. Каждое удобство имеет цену - знайте эту цену до того, как заплатите её в production.


Конфигурация

import pyspark.pandas as ps

# Показать все опции
ps.get_option("display.max_rows")
ps.set_option("display.max_rows", 100)

# Разрешить операции на DF из разных источников (замедляет - добавляет join по индексу)
ps.set_option("compute.ops_on_diff_frames", True)

# Дефолтный тип индекса (ключевая настройка производительности)
ps.set_option("compute.default_index_type", "distributed")          # рекомендуется
# ps.set_option("compute.default_index_type", "distributed-sequence")  # монотонный, без глобальной сортировки
# ps.set_option("compute.default_index_type", "sequence")            # медленно, только для отладки

# Количество строк для repr
ps.set_option("display.max_rows", 1000)
ps.set_option("display.max_columns", 50)

Настройки Spark, влияющие на pyspark.pandas:

# Arrow для передачи данных (критично для производительности)
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "10000")  # размер батча

# Количество партиций после Shuffle
spark.conf.set("spark.sql.shuffle.partitions", "200")   # дефолт для production

# Для тестов - уменьшить чтобы не гонять 196 пустых партиций
spark.conf.set("spark.sql.shuffle.partitions", "4")