Pandas API on Spark (pyspark.pandas): когда и как
pyspark.pandas (бывш. Koalas) даёт pandas-совместимый интерфейс поверх Spark. Разбираем, где это ускоряет миграцию, где замедляет работу, и как API соотносится с нативным PySpark.
Что такое 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.pandas → to_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 медленнее нативного PySparkpyspark.pandasсdistributedиндексом в 1.5x медленнее нативного PySpark- Для большинства задач разница между
distributedpyspark.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")