RDD, DataFrame и Dataset: эволюция абстракций Spark
Три уровня API Spark изнутри: UnsafeRow и off-heap в DataFrame, Pickle vs Arrow в Python RDD, Encoders в Dataset - и почему в PySpark нет Dataset.
В предыдущих уроках мы работали с DataFrame как с данностью. Пора разобраться, что за ним стоит: почему он быстрее RDD в 10–50 раз, почему Dataset недоступен в Python, как Spark хранит миллиарды строк эффективнее Python-объектов, и что происходит на уровне памяти и процессора при каждом вызове filter() или groupBy().
Урок идёт от истории к архитектуре: RDD и его проблемы → катастрофа Python-сериализации → DataFrame и UnsafeRow → Catalyst и Tungsten → Dataset и Encoders → почему в PySpark нет Dataset → практика и домашнее задание.
Зачем три разных API? История и мотивация¶
Spark развивался итерационно: каждая следующая абстракция решала конкретные проблемы предыдущей. Понять почему API устроен именно так невозможно без понимания истории.
| API | Появился | Схема | Catalyst | Типизация | Доступен в Python |
|---|---|---|---|---|---|
| RDD | Spark 1.0 | Нет | Нет | Объекты JVM / Python | Да (медленно) |
| DataFrame | Spark 1.3 | Да | Да | Нет (runtime errors) | Да (рекомендуется) |
| Dataset | Spark 1.6 | Да | Да | Да (compile-time) | Нет |
Короткий ответ на вопрос «что использовать в PySpark» - DataFrame. Длинный ответ - читайте дальше: понимание того, почему именно DataFrame, поможет вам писать более быстрый код и правильно диагностировать проблемы производительности.
RDD¶
RDD (Resilient Distributed Dataset) - это фундаментальная структура данных, которая является отказоустойчивой, неизменяемой и представляет собой распределенную коллекцию объектов. RDD неизменяемы, то есть их нельзя изменить после создания. Любое преобразование RDD приводит к созданию нового RDD. Каждый набор данных в RDD разделен на логические разделы, которые могут быть вычислены на разных узлах кластера.
DataFrame¶
DataFrame - это распределенный набор данных, состоящий из данных, расположенных в строках и столбцах с именованными атрибутами. Он имеет сходство с таблицами реляционных баз данных или фреймами данных R/Python, но включает в себя сложные оптимизации.
Если вы знакомы с Python, я предполагаю, что вы уже знаете, что такое Pandas DataFrame. PySpark DataFrame во многом похож на Pandas DataFrame, за исключением того, что DataFrame распределены по кластеру (это означает, что данные в DataFrame хранятся на разных машинах в кластере), и любые операции в PySpark выполняются параллельно на всех машинах, тогда как Pandas DataFrame хранит и обрабатывает данные на одной машине.
Пример работы с DataFrame в PySpark:
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, avg
# 1. Создаем сессию Spark
spark = SparkSession.builder \
.appName("PySparkDataFrameExample") \
.getOrCreate()
# 2. Подготовим тестовые данные (список кортежей и названия колонок)
data = [
("Иван", "IT", 85000, 28),
("Мария", "IT", 95000, 32),
("Анна", "HR", 70000, 25),
("Петр", "HR", 65000, 40),
("Елена", "Финансы", 120000, 35)
]
columns = ["Имя", "Департамент", "Зарплата", "Возраст"]
# 3. Создаем DataFrame
df = spark.createDataFrame(data, schema=columns)
# 4. Выводим схему и первые строки в консоль
print("Исходный DataFrame:")
df.printSchema()
df.show()
# 5. Трансформации (Фильтрация и выбор колонок)
# Выбираем сотрудников из IT, которые старше 30 лет
it_experienced_df = df.filter((col("Департамент") == "IT") & (col("Возраст") > 30)) \
.select("Имя", "Зарплата")
print("Сотрудники IT старше 30:")
it_experienced_df.show()
# 6. Агрегация (Группировка)
# Считаем среднюю зарплату по каждому департаменту
avg_salary_df = df.groupBy("Департамент") \
.agg(avg("Зарплата").alias("Средняя_Зарплата"))
print("Средняя зарплата по департаментам:")
avg_salary_df.show()
# 7. Сохранение результата (например, в формат CSV)
# coalesce(1) нужен, чтобы объединить данные в один файл вместо нескольких партиций
avg_salary_df.coalesce(1).write.mode("overwrite").csv("path/to/output_report", header=True)
# 8. Останавливаем сессию
spark.stop()
Dataset¶
Dataset в Apache Spark - это строго типизированная, объектно-ориентированная распределенная коллекция данных, которая объединяет в себе лучшие свойства низкоуровневого RDD API (высокую безопасность типов во время компиляции и возможность работы со стандартными классами языков программирования) и высокоуровневого DataFrame API (высокую производительность и встроенные оптимизации механизма Catalyst).
Эта структура данных, доступная преимущественно в Scala и Java, представляет собой таблицу с жестко заданной схемой, где каждая строка является не просто абстрактным объектом общего назначения, а конкретным экземпляром пользовательского класса (например, объекта в формате Case Class в Scala или Java Bean).
Основное преимущество Dataset заключается в том, что он позволяет разработчикам использовать привычные лямбда-выражения и функции высшего порядка для трансформации данных с полной проверкой корректности кода на этапе сборки приложения, но при этом под капотом Spark не выполняет дорогостоящую десериализацию объектов в стандартную память Java/Scala, а оперирует данными в оптимизированном бинарном формате Tungsten, обеспечивая минимальное потребление оперативной памяти и максимальную скорость вычислений.
1. RDD: фундамент и проклятие JVM-объектов¶
Что такое RDD¶
RDD (Resilient Distributed Dataset) - базовая абстракция Spark 1.x, появившаяся в 2012 году в статье Матея Захарии «Resilient Distributed Datasets: A Fault-Tolerant Abstraction for In-Memory Cluster Computing». Это иммутабельная распределённая коллекция объектов без фиксированной схемы. DataFrame и Dataset построены поверх RDD - это важно понимать: RDD никуда не делся, он просто спрятан под слоем абстракций.
Каждый RDD определяется пятью характеристиками, которые Spark хранит как метаданные и использует для планирования вычислений:
| Свойство | Что это | Пример |
|---|---|---|
| Partitions | Список разделов, из которых состоит RDD | rdd.getNumPartitions() |
| Dependencies | Ссылки на родительские RDD (Lineage) | Narrow vs Wide зависимости |
| Compute function | Функция вычисления данных партиции | lambda partition: ... |
| Partitioner | Опциональная функция разбиения по ключу | HashPartitioner(100) |
| Preferred locations | Подсказки где запускать Task (Data Locality) | Адреса HDFS-блоков |
# Создать RDD напрямую из Python-коллекции
rdd = spark.sparkContext.parallelize(range(1_000_000), numSlices=8)
print(rdd.getNumPartitions()) # 8
# Из файла - партиции соответствуют HDFS-блокам (~128 MB каждый)
rdd_file = spark.sparkContext.textFile("s3a://data/logs/*.gz")
# Трансформации возвращают новый RDD (иммутабельность)
rdd_filtered = rdd.filter(lambda x: x % 2 == 0)
rdd_mapped = rdd_filtered.map(lambda x: (x, x * x))
# Ничего не вычисляется до Action - Lazy Evaluation
Иммутабельность - ключевое свойство. Нельзя «изменить» существующий RDD: каждая трансформация создаёт новый RDD с ссылкой на родительский. Это делает возможным lineage-граф для отказоустойчивости.
Оверхед JVM-объектов: почему Java расточительна с памятью¶
Первая фундаментальная проблема RDD в Spark 1.x - расход памяти JVM-объектов. Каждая строка данных хранилась как полноценный Java-объект на JVM Heap. Это звучит нормально для привычного ООП, но катастрофично для Big Data.
Рассмотрим что происходит когда Java хранит простую строку с тремя полями: (userId: Long, name: String, age: Int).
Почему Java-объект такой раздутый:
- Object Header (16 байт): каждый объект в Java имеет заголовок из двух слов - mark word (hashcode, GC-метки, монитор синхронизации) и class pointer (ссылка на метаданные класса). Это 16 байт накладных расходов на каждый объект - даже на
new Integer(42). - Указатели вместо значений: для ссылочных типов (String, массивы) JVM хранит 8-байтный указатель на другой объект на куче. Один логический «пользователь» в памяти - это на самом деле граф из трёх объектов:
User → String → char[]. - Padding для выравнивания: JVM выравнивает размер объектов до кратного 8 байтам. Если данные занимают 12 байт - объект займёт 16 байт (4 байта мусора).
- UTF-16 кодировка String: Java хранит строки в UTF-16 (2 байта на символ). Слово «hello» → 10 байт данных + 32 байта заголовков двух объектов. В UTF-8 то же слово - 5 байт.
Итог: для хранения 1 млрд строк с тремя простыми полями в виде Java-объектов нужно ~80 GB памяти. В UnsafeRow тот же массив данных - ~29 GB. Разница в 2,5–3×, и это без учёта коллекций внутри объектов.
GC паузы: Stop-the-World катастрофа на Big Data¶
Вторая фундаментальная проблема RDD - сборщик мусора (Garbage Collector) Java. JVM автоматически освобождает память от недостижимых объектов, но делает это в «паузах сборки мусора» (GC pauses).
Для обычных серверных приложений GC - незаметный фоновый процесс. Для кластера Spark с миллиардами объектов это катастрофа.
Механика проблемы на Spark-кластере:
- Spark Executor обрабатывает партицию из 10 млн строк - в JVM heap 10 млн объектов.
- После завершения Stage объекты становятся мусором, GC запускается.
- Во время GC все потоки Executor полностью остановлены - Stop-the-World.
- При heap 16–32 GB и миллиардах мелких объектов пауза GC может занимать секунды или минуты.
- Driver не получает heartbeat от Executor → считает его мёртвым → задача пересчитывается → снова GC → дедлок.
Именно поэтому в Spark 1.x с RDD-кластеры периодически «замерзали» на многие минуты без видимой причины - это был GC. Проект Tungsten, появившийся в Spark 1.4, был создан именно для решения этой проблемы.
Lineage: отказоустойчивость без репликации данных¶
Несмотря на проблемы с памятью, RDD принёс революционную идею - Lineage. «Resilient» в названии означает устойчивость к отказам без физической репликации данных, как делают HDFS или Kafka.
Spark хранит не сами данные, а граф зависимостей - последовательность трансформаций, которые привели к данному RDD. Если Executor упал и потерял партицию, Spark пересчитывает её заново, повторив цепочку трансформаций от источника.
Это работает потому что трансформации детерминированы и чистые (pure functions): один и тот же входной RDD с одними и теми же трансформациями всегда даёт одинаковый результат. Если функция имеет побочные эффекты или недетерминирована (например, считывает из нестабильного источника) - lineage ломается.
Ограничение lineage проявляется при длинных цепочках трансформаций: пересчёт от источника занимает много времени. Именно для этого существует persist()/cache() - сохранить промежуточный RDD в память, чтобы при отказе пересчитывать не от самого начала, а от ближайшей материализованной точки.
# Сохранить промежуточный результат после дорогой трансформации
expensive_rdd = raw_rdd.map(complex_parse).filter(validate).reduceByKey(merge)
expensive_rdd.persist() # materialize in memory
# Теперь при отказе Spark пересчитывает только с этой точки
result = expensive_rdd.groupByKey().mapValues(aggregate)
2. Python RDD: Pickle и катастрофа производительности¶
Анатомия Python Worker: два процесса, один сокет¶
Когда вы пишете rdd.map(lambda x: x * 2) в PySpark, происходит нечто гораздо более сложное, чем кажется. Spark написан на Scala и работает в JVM. Ваш Python-код - отдельный процесс.
Каждая строка данных проходит следующий путь:
- JVM читает строку из UnsafeRow-буфера или с диска
- JVM сериализует строку через Python Pickle в байтовый поток
- Байты передаются через Unix Domain Socket в Python Worker (отдельный процесс)
- Python Worker десериализует байты обратно в Python-объект
- Python выполняет вашу лямбду над объектом
- Python сериализует результат снова через Pickle
- Байты передаются обратно в JVM через сокет
- JVM десериализует результат
Это происходит для каждой строки отдельно. При 100 млн строк - 200 млн операций сериализации/десериализации + 200 млн пересечений границы процессов.
Pickle: почему это катастрофа производительности¶
Pickle - стандартный Python-протокол сериализации объектов. Он универсален, но катастрофически медленен для Big Data по нескольким причинам:
1. Построчная обработка. Каждый объект сериализуется и передаётся отдельно. Нет пакетной передачи, нет векторизации. При 10 млн строк - 10 млн отдельных вызовов pickle.dumps().
2. Медленная сериализация. Pickle должен пройти всю иерархию объекта (поля, вложенные объекты), записать типы, значения, метаданные. Это в 10–20 раз медленнее бинарной сериализации Tungsten.
3. Раздутые байты. Pickle хранит не просто данные, но и информацию о типах Python. Один целочисленный объект 42 в pickle - это ~14 байт вместо 4 байт сырого int.
4. Давление на GC с обеих сторон. JVM создаёт байтовые массивы для каждой строки - мусор для JVM GC. Python создаёт объекты для каждой строки - мусор для Python GC (циклический сборщик). Оба GC работают одновременно.
# Бенчмарк: Python RDD vs DataFrame
# Медленно: каждая строка → Pickle → сокет → Python → сокет → JVM
rdd = spark.sparkContext.parallelize(range(10_000_000))
result_rdd = rdd.map(lambda x: x * 2).filter(lambda x: x > 1_000_000).count()
# ~30–60 секунд на типичном кластере
# Быстро: всё в JVM, Catalyst + Tungsten, без Python overhead
from pyspark.sql.functions import col
df = spark.range(10_000_000)
result_df = df.select(col("id") * 2).filter(col("id") * 2 > 1_000_000).count()
# ~1–3 секунды на том же кластере
# Разница: 10–50× в пользу DataFrame
Именно поэтому Python RDD в production Big Data - антипаттерн. Он работает, но тратит 90% CPU и времени на сериализацию, а не на реальные вычисления.
Apache Arrow: как Python UDF стал быстрым¶
Осознав катастрофу Pickle, команда Spark интегрировала Apache Arrow - открытый стандарт колоночного представления данных в памяти. Arrow появился в PySpark 2.3 (2018) и радикально изменил производительность Python UDF.
Принципиальные отличия Arrow от Pickle:
- Батчевая передача: вместо строки за строкой Arrow передаёт батчи (RecordBatch) из 1024 строк. Одна операция IPC вместо 1024.
- Zero-Copy: данные не копируются. JVM и Python Worker разделяют один буфер памяти через shared memory или file descriptor. Нет сериализации - есть только передача указателя.
- Колоночный формат: колонки хранятся как непрерывные массивы примитивов. Pandas, numpy и SIMD-инструкции процессора работают с этим форматом нативно.
- Без накладных расходов типов: Arrow хранит только значения, без Python-метаданных типов.
# Pandas UDF (Vectorized UDF) - использует Arrow
from pyspark.sql.functions import pandas_udf
import pandas as pd
@pandas_udf("double")
def normalize_udf(series: pd.Series) -> pd.Series:
"""Весь батч приходит как Pandas Series - vectorized!"""
return (series - series.mean()) / series.std()
# Spark передаёт данные в Python как Arrow RecordBatch → Pandas Series
# Не одна строка, а тысячи строк за один IPC-вызов
df.withColumn("amount_norm", normalize_udf(col("amount")))
Pandas UDF с Arrow работает в 10–100× быстрее обычного Python UDF с Pickle. Для операций, которые действительно нужно делать в Python (ML-инференс, сложная бизнес-логика), Pandas UDF - правильный выбор.
3. DataFrame: схема + Catalyst + Tungsten¶
Что изменилось с появлением DataFrame¶
DataFrame появился в Spark 1.3 (2014) как результат осознания, что RDD - слишком низкоуровневый инструмент для большинства задач аналитики. Ключевое изменение - у DataFrame появилась схема.
Когда у данных есть схема, Spark перестаёт быть слепым исполнителем и становится умным оптимизатором:
- Он знает, что колонка
amountимеет типDouble→ может выбрать правильный алгоритм join - Он знает, что колонка
countryиспользуется вfilter→ может прочитать только нужные колонки из Parquet (Column Pruning) - Он знает структуру запроса целиком → может переставить операции в оптимальном порядке (Predicate Pushdown, Join Reordering)
# DataFrame несёт схему - Spark "видит" структуру данных
df = spark.read.parquet("s3a://data/events/")
df.printSchema()
# root
# |-- user_id: long (nullable = false)
# |-- event_type: string (nullable = true)
# |-- amount: double (nullable = true)
# |-- event_date: date (nullable = true)
# RDD - чёрный ящик. Spark не знает что внутри
rdd = df.rdd # RDD[Row] - структура скрыта за абстракцией Row
UnsafeRow: революция в хранении данных¶
Project Tungsten (2015, Spark 1.4) полностью переосмыслил хранение данных в Spark. Центральное изобретение - UnsafeRow: бинарный формат хранения строк DataFrame в памяти.
Идея проста и радикальна: вместо хранения данных как Java-объектов на JVM Heap, хранить их как компактные бинарные буферы в стиле C. Spark берёт управление памятью на себя, обходя JVM GC полностью.
Детальная анатомия UnsafeRow:
Null Bitmap (битовая маска null-значений). Первые ceil(numFields / 8) байт строки - битовая маска: бит i равен 1 если поле i равно null. Это позволяет проверить null без распаковки поля - просто проверить один бит.
Фиксированная часть. Следующие 8 * numFields байт - фиксированные данные каждого поля. Для примитивных типов (Long, Int, Double, Boolean) здесь хранится само значение. Для переменной длины (String, Array) здесь хранится пара (смещение, длина) - 4+4 байта, указывающая на место в переменной части.
Переменная часть. После фиксированной части идут данные переменной длины: строки в UTF-8, сериализованные массивы, вложенные структуры. Каждый элемент выровнен по 8 байт.
Почему это быстрее Java-объектов:
- Нет Object Header: сэкономлено 16 байт на каждой записи
- Нет указателей на подобъекты: строки хранятся встроенно, нет разыменования
- Смежное расположение в памяти: последовательный доступ → L1/L2 кэш процессора работает эффективно (cache hit ratio > 95%)
- Off-heap: данные не попадают в JVM GC, нет паузы Stop-the-World
# Получить доступ к UnsafeRow напрямую (для диагностики)
df = spark.createDataFrame([(1, "Alice", 30.5)], ["id", "name", "score"])
internal_row = df.rdd.take(1)[0] # Row (не UnsafeRow - RDD возвращает Row)
# UnsafeRow виден в explain() как внутренний формат
df.filter(col("id") > 0).explain(mode="formatted")
# В выводе видны операторы работающие с UnsafeRow напрямую
Off-Heap память: Spark выходит за пределы JVM¶
Off-Heap - техника выделения памяти в обход JVM Heap через нативный API операционной системы (sun.misc.Unsafe.allocateMemory()). Память выделяется напрямую в RAM ОС, минуя JVM-кучу и, следовательно, Garbage Collector.
Когда Off-Heap включён (spark.memory.offHeap.enabled=true, spark.memory.offHeap.size=...), все данные DataFrame - строки UnsafeRow, shuffle-буферы, кэшированные партиции - хранятся вне JVM Heap. GC просто не знает об их существовании и не тратит время на их сканирование.
# Включить Off-Heap в SparkSession
spark = SparkSession.builder \
.config("spark.memory.offHeap.enabled", "true") \
.config("spark.memory.offHeap.size", "10g") \
.config("spark.executor.memory", "4g") \
.getOrCreate()
# Теперь данные DataFrame уходят в Off-Heap
# JVM Heap 4 GB используется только для Spark-метаданных
# Off-Heap 10 GB хранит UnsafeRow данных
# Проверить использование Off-Heap (в Spark UI → Executors)
# колонки "Off-Heap Memory Used"
Риски Off-Heap: Spark сам управляет памятью на байтовом уровне через sun.misc.Unsafe. Memory leak в Spark или неправильная конфигурация может привести к OOM вне JVM - такой OOM не виден стандартными JVM-инструментами и тяжело диагностируется.
Catalyst: четыре фазы оптимизации запроса¶
Catalyst - оптимизатор запросов Spark SQL. Именно Catalyst даёт DataFrame преимущество перед RDD: он «понимает» запрос как математическое выражение и может переписать его в более эффективную форму.
Разберём каждую фазу подробно:
1. Unresolved Logical Plan. Spark строит дерево выражений из вашего кода, не проверяя корректность. На этом этапе col("country") - просто узел UnresolvedAttribute("country"), без проверки существует ли такая колонка. Это чистая синтаксическая структура.
2. Analyzed Logical Plan (Analysis Phase). Catalyst обращается к Catalog - реестру таблиц, схем и метаданных. Каждый UnresolvedAttribute разрешается в конкретную колонку с типом. Если колонка не найдена - AnalysisException именно здесь. Типы проверяются и при несовместимости генерируется ошибка. Null-правила применяются согласно схеме.
3. Optimized Logical Plan. Catalyst применяет набор правил оптимизации (Rule-Based Optimization):
- Predicate Pushdown:
filter(country == "RU")перемещается как можно ближе к источнику - перед join, перед агрегацией. Это уменьшает объём данных на каждом следующем шаге. - Column Pruning: если после
select("user_id", "amount")вы делаетеgroupBy("user_id").sum("amount")- Catalyst уберёт из плана все остальные колонки. Parquet не будет читать ненужные данные вообще. - Constant Folding: выражения со статическими значениями вычисляются сразу.
col("amount") * 1.18→ числовое умножение,lit(2) + lit(3)→lit(5). - Boolean Simplification:
NOT (a AND NOT b)→NOT a OR b(законы Де Моргана).
4. Physical Planning. Оптимизированный логический план конвертируется в несколько физических планов - конкретных алгоритмов выполнения. Для join это может быть Sort-Merge Join или Broadcast Hash Join. Для агрегации - HashAggregate или SortAggregate. Cost-Based Optimizer (CBO) выбирает лучший план на основе статистики таблиц (размер, кардинальность колонок).
# Увидеть каждую фазу Catalyst
df_result = (
df_events
.filter(col("country") == "RU")
.join(df_regions, "region_id")
.groupBy("region_name").agg(sum("amount").alias("total"))
)
# Полный детальный план с форматированием
df_result.explain(mode="formatted")
# == Parsed Logical Plan ==
# == Analyzed Logical Plan ==
# == Optimized Logical Plan ==
# == Physical Plan ==
# Только физический план (краткий)
df_result.explain()
# Режим cost: покажет статистику CBO
df_result.explain(mode="cost")
WholeStageCodegen: компилятор вместо интерпретатора¶
WholeStageCodegen - самая сложная оптимизация Tungsten, появившаяся в Spark 2.0. Чтобы понять её, нужно понять проблему, которую она решает.
Традиционная модель выполнения запросов - Volcano Model (или Iterator Model): каждый оператор - отдельный класс с методом next(), который возвращает одну строку. Оператор filter вызывает next() у оператора scan, проверяет условие, и возвращает строку вверх по стеку к оператору project. Очень элегантно, но катастрофически медленно:
- Один вызов
next()- это виртуальный вызов через интерфейс. Процессор не может предсказать адрес следующей инструкции (branch misprediction). - Данные проходят через цепочку операторов по одной строке - плохая утилизация CPU cache.
- Каждый оператор - отдельный уровень абстракции - накладные расходы на каждом слое.
WholeStageCodegen устраняет всё это: Catalyst динамически генерирует Java-bytecode, который сливает все операторы в один тугой цикл без виртуальных вызовов.
Сгенерированный код специализирован под конкретный запрос: вместо обобщённого row.get(columnIndex) генерируется row.getLong(0) - прямой доступ к байту в UnsafeRow без полиморфизма. JIT-компилятор Java оптимизирует такой код до нативного ассемблера, часто применяя SIMD-инструкции для параллельной обработки нескольких строк.
# WholeStageCodegen виден в explain() - ищите * перед операторами
df.filter(col("id") > 100).groupBy("country").agg(sum("amount")).explain()
# == Physical Plan ==
# *(2) HashAggregate(...) ← * означает WholeStageCodegen
# +- Exchange hashpartitioning(...)
# +- *(1) HashAggregate(...) ← ещё один WholeStageCodegen stage
# +- *(1) Filter (id > 100)
# +- *(1) FileScan parquet ...
# Все операторы с одним номером (*1) или (*2) сгенерированы в один Java-метод
Эффект WholeStageCodegen: ускорение в 2–5× для CPU-bound операций по сравнению с Volcano Model. Именно поэтому DataFrame в некоторых бенчмарках обходит ручной Java-код на Hadoop MapReduce.
4. Dataset: типизация для JVM-разработчиков¶
Dataset доступен только в Scala и Java. Если вы пишете исключительно на Python - этот раздел полезен для понимания архитектуры Spark и ответа на вопрос «почему в PySpark нет Dataset», но не применяется в daily работе.
Что такое Dataset¶
Dataset - типизированный DataFrame. Вместо безымянных Row каждая запись - конкретный Scala case class или Java Bean, и это проверяется компилятором на этапе написания кода, а не в runtime.
// Scala: Dataset[User] vs DataFrame (= Dataset[Row])
case class User(userId: Long, name: String, age: Int)
val ds: Dataset[User] = spark.read.json("...").as[User]
val df: DataFrame = spark.read.json("...") // = Dataset[Row]
// Dataset: ошибка типов → подчёркивается IDE сразу при написании
ds.map(u => u.nme) // error: value nme is not a member of User
ds.filter(u => u.age > "string") // type error: comparing Int with String
// DataFrame: ошибка обнаруживается только при выполнении Action
df.select("nme") // AnalysisException при df.show() или count()
Compile-time safety - главная ценность Dataset для больших команд на Scala/Java: ошибки в именах колонок и типах не добираются до production, IDE подсказывает правильные поля через автодополнение, рефакторинг работает корректно (переименование поля компилятор найдёт везде).
Encoders: как Dataset сохраняет скорость Tungsten¶
Казалось бы, если Dataset хранит объекты (User, Transaction), он должен страдать от тех же проблем JVM-объектов, что и RDD. Но нет - Encoders решают эту проблему.
Encoder - сгенерированный Catalyst байткод для конвертации между UnsafeRow и вашим классом. Spark не хранит объекты User в памяти постоянно. Внутренний формат хранения - всё тот же UnsafeRow. Encoder вступает в дело только когда вы применяете лямбду, требующую объект User.
Encoder не использует обобщённую Java-рефлексию или стандартную сериализацию. Catalyst генерирует специализированный байткод для конкретного класса: он знает, что userId - Long в позиции 0, name - String в позиции 1, age - Int в позиции 2, и генерирует прямые операции чтения байт из UnsafeRow.
// Scala: операции с column-syntax работают с UnsafeRow напрямую (нет overhead)
ds.filter($"age" > 18).groupBy($"name").count()
// Одинаково быстро с:
df.filter($"age" > 18).groupBy($"name").count()
// Операции с lambda-syntax требуют Encoder roundtrip (небольшой overhead)
ds.filter(u => u.age > 18) // UnsafeRow → User → boolean → UnsafeRow
Когда Dataset медленнее DataFrame: Encoder overhead¶
Парадокс Dataset: в некоторых сценариях он медленнее DataFrame. Это происходит когда вы используете lambda-синтаксис (ds.map(u => ...)) вместо column-синтаксиса (ds.filter($"age" > 18)).
// Одинаково быстро - Catalyst работает с UnsafeRow напрямую
ds.filter($"age" > 18) // column syntax: нет Encoder roundtrip
df.filter($"age" > 18) // то же самое под капотом
// Dataset медленнее - каждая строка: UnsafeRow → User → UnsafeRow
ds.filter(u => u.age > 18) // lambda syntax: Encoder overhead
// Худший случай: map с лямбдой создаёт новый объект на каждую строку
ds.map(u => User2(u.userId, u.name.toUpperCase, u.age * 2))
// Миллиард строк = миллиард создаваемых User2-объектов → GC давление
Практическое правило для Scala/Java: используйте column-синтаксис ($"col") для фильтрации и агрегации, lambda-синтаксис (u => u.field) только там, где это необходимо для сложной бизнес-логики.
5. Почему Dataset нет в PySpark¶
Это один из самых частых вопросов Python-разработчиков, переходящих от Scala к PySpark. Ответ лежит в фундаментальных различиях между Python и JVM-языками.
Динамическая типизация Python vs compile-time generics¶
Dataset работает через параметрический полиморфизм - Dataset[User] это обобщённый тип (generic), который на этапе компиляции Scala/Java проверяет, что тип T в Dataset[T] соответствует данным. Compile-time type safety.
Python - динамически типизированный язык. Концепция «тип переменной известен на этапе компиляции» фундаментально противоречит Python:
- Нет фазы компиляции в смысле JVM: Python компилируется в байткод только на уровне синтаксиса, без проверки типов
- Duck typing: Python не требует явного объявления типов переменных
- Рефлексия во всём: любой объект можно изменить в runtime, добавить поля, сменить класс
- Нет generics в смысле Java:
List[int]в Python - только подсказка для статических анализаторов (mypy), не для интерпретатора
Следовательно, Encoder для Dataset[User] в Python написать невозможно: нет compile-time информации о типе User, которую Catalyst мог бы использовать для генерации байткода. Python-объекты в runtime могут иметь произвольную структуру.
PySpark DataFrame = Scala Dataset[Row] под капотом¶
Важное понимание: когда вы пишете df.filter(col("amount") > 100) в Python - ваша команда через Py4J (Python-Java bridge) вызывает Scala-метод. Объект df в Python - это тонкая обёртка над Dataset[Row] в JVM.
Это означает несколько важных вещей:
- Скорость нативная.
df.filter(),df.groupBy(),df.join()- все эти операции выполняются в JVM с полными оптимизациями Catalyst и Tungsten. Python-overhead существует только для Py4J-вызовов при построении плана, не при выполнении. - Python не видит строк. До тех пор пока вы не вызываете
df.collect(),df.toPandas()или Python UDF - ни одна строка данных не покидает JVM. Все данные - в UnsafeRow-буферах. - Python UDF - единственный overhead. Если вы применяете обычный Python UDF (
@udf), Spark вынужден сериализовать строку через Arrow/Pickle, передать в Python Worker, получить результат обратно. Всё остальное - чистый JVM.
Почему это хорошо для Python-разработчика¶
Отсутствие Dataset в PySpark - не ограничение, а правильное архитектурное решение. Вот почему:
1. Вы получаете скорость Dataset без его сложности. PySpark DataFrame выполняется с той же эффективностью Tungsten, что и Scala DataFrame (= Dataset[Row]). Вам не нужен Dataset чтобы получить UnsafeRow и WholeStageCodegen.
2. Ошибки типов ловятся на Analysis Phase. Catalyst проверяет схему при построении плана. Да, это runtime, а не compile-time - но на практике ошибка AnalysisException: column not found обнаруживается при первом explain() или show(), задолго до production.
3. mypy + pyspark-stubs. Для статической проверки типов в Python-коде можно использовать mypy с stub-файлами PySpark. Это не даёт 100% гарантии как Scala Dataset, но покрывает большинство ошибок.
# Python: нет Dataset, но есть статическая проверка через mypy
from pyspark.sql import DataFrame
from pyspark.sql.functions import col
def process_events(df: DataFrame) -> DataFrame:
"""mypy проверит типы аргументов и возвращаемого значения."""
return (
df.filter(col("amount") > 0)
.groupBy("user_id")
.agg({"amount": "sum"})
)
# Явная схема - максимальная защита от ошибок
from pyspark.sql.types import StructType, StructField, LongType, DoubleType
schema = StructType([
StructField("user_id", LongType(), False),
StructField("amount", DoubleType(), True),
])
df = spark.read.schema(schema).parquet("s3a://events/")
# AnalysisException будет немедленно при любой опечатке в имени колонки
6. Сравнение API: производительность, использование, ограничения¶
Сводная таблица производительности¶
Для одного и того же вычисления (filter + group by + aggregate, 1 млрд строк):
| API | Время | Память | Причина |
|---|---|---|---|
| DataFrame | 1× (baseline) | 1× | UnsafeRow + Catalyst + WholeStageCodegen |
| Dataset (column syntax) | ~1× | ~1× | Идентично DataFrame - Catalyst работает с UnsafeRow |
| Dataset (lambda syntax) | 1.5–2× | 1.5× | Encoder roundtrip на каждой строке |
| Scala RDD | 3–5× | 2–3× | Kryo/Java сериализация, нет Catalyst |
| Python RDD | 20–50× | 5–10× | Pickle + Python Worker overhead |
| Python UDF на DataFrame | 3–10× | 2× | Arrow батчи быстрее Pickle, но всё равно IPC |
| Pandas UDF на DataFrame | 1.5–3× | 1.5× | Arrow + векторизованный Python |
Dataset с lambda-синтаксисом медленнее чистого DataFrame - контринтуитивно, но объясняется тем что каждая строка десериализуется в JVM-объект и обратно.
Дерево решений: что выбрать¶
7. Практика: бенчмарк, конвертации и инспекция планов¶
Конвертации между API в production-коде¶
В production часто встречается legacy-код с RDD, который нужно интегрировать с новым DataFrame-кодом. Безопасные паттерны конвертации:
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, LongType, StringType, DoubleType
from pyspark.sql.functions import col
spark = SparkSession.builder.master("local[*]").appName("APIDemo").getOrCreate()
# ── DataFrame → RDD (редко нужно, дорого!) ──
df = spark.range(1000).toDF("id")
rdd_rows = df.rdd # RDD[Row] - UnsafeRow → Row materialization
rdd_ids = df.rdd.map(lambda r: r["id"]) # RDD[long]
# ВНИМАНИЕ: df.rdd материализует UnsafeRow в обычные Row-объекты
# Это дорогая операция - избегайте в hot path
# ── RDD → DataFrame (типичный сценарий при работе с legacy) ──
raw_rdd = spark.sparkContext.parallelize([
(1, "Alice", 30.5),
(2, "Bob", 25.0),
(3, "Carol", 35.8),
])
schema = StructType([
StructField("id", LongType(), nullable=False),
StructField("name", StringType(), nullable=True),
StructField("score", DoubleType(), nullable=True),
])
df_from_rdd = spark.createDataFrame(raw_rdd, schema)
# toDF() - быстрый способ без явной схемы (типы выводятся автоматически)
df2 = raw_rdd.toDF(["id", "name", "score"])
df_from_rdd.printSchema()
# root
# |-- id: long (nullable = false)
# |-- name: string (nullable = true)
# |-- score: double (nullable = true)
# ── Правильная явная схема для CSV/JSON ──
# Без явной схемы Spark читает весь файл для inference - дорого
schema_csv = StructType([
StructField("user_id", LongType(), False),
StructField("event_type", StringType(), True),
StructField("amount", DoubleType(), True),
])
df_csv = spark.read.schema(schema_csv).csv("s3a://data/events.csv")
# Сразу читает с правильными типами, без лишнего прохода по данным
Инспекция планов Catalyst: видим оптимизации в действие¶
# Создать датасет для демонстрации
df_events = spark.createDataFrame([
(1, "RU", "purchase", 100.0),
(2, "US", "view", 5.0),
(3, "RU", "purchase", 200.0),
], ["user_id", "country", "event_type", "amount"])
df_regions = spark.createDataFrame([
(1, "Москва", "RU"),
(2, "Питер", "RU"),
], ["region_id", "region_name", "country"])
# Сложный запрос
result = (
df_events
.filter(col("country") == "RU") # predicate - будет pushed down
.filter(col("event_type") == "purchase") # ещё один predicate
.join(df_regions, "country")
.groupBy("region_name")
.agg({"amount": "sum"})
)
# Полный план с форматированием - читаем снизу вверх
result.explain(mode="formatted")
# == Physical Plan ==
# *(2) HashAggregate(keys=[region_name], functions=[sum(amount)])
# +- Exchange hashpartitioning(region_name, 200)
# +- *(1) HashAggregate(keys=[region_name], functions=[partial_sum(amount)])
# +- *(1) BroadcastHashJoin [country], [country], Inner, BuildRight
# :- *(1) Filter ((country = RU) AND (event_type = purchase)) ← оба filter слиты!
# : +- *(1) Scan ExistingRDD[user_id,country,event_type,amount]
# +- BroadcastExchange HashedRelationBroadcastMode
# +- *(1) Scan ExistingRDD[region_id,region_name,country]
# Что видим:
# 1. *(1) - WholeStageCodegen, операции в одном сгенерированном методе
# 2. Filter слил два условия в одно (constant folding + predicate merge)
# 3. BroadcastHashJoin - Catalyst автоматически выбрал broadcast для маленькой таблицы
# 4. partial_sum - HashAggregate выполнен в две фазы (partial → shuffle → final)
Бенчмарк: RDD vs DataFrame в реальных условиях¶
import time
# Генерировать данные
spark.conf.set("spark.sql.shuffle.partitions", "8")
N = 5_000_000
# ── Тест 1: Python RDD ──
rdd = spark.sparkContext.parallelize(range(N), 8)
start = time.time()
result_rdd = (
rdd
.map(lambda x: (x % 1000, x)) # (key, value)
.filter(lambda kv: kv[1] > N // 2) # отфильтровать половину
.reduceByKey(lambda a, b: a + b) # сумма по ключу
.count()
)
print(f"Python RDD: {time.time() - start:.1f}s")
# ── Тест 2: DataFrame ──
from pyspark.sql.functions import col, sum as spark_sum
df = spark.range(N).withColumn("key", (col("id") % 1000).cast("long"))
start = time.time()
result_df = (
df
.filter(col("id") > N // 2)
.groupBy("key")
.agg(spark_sum("id").alias("total"))
.count()
)
print(f"DataFrame: {time.time() - start:.1f}s")
# Типичные результаты на кластере с 4 Executor × 4 CPU:
# Python RDD: 45–90 секунд
# DataFrame: 2–5 секунд
# Разница: 15–40×
Когда Python UDF неизбежен: делаем его правильно¶
from pyspark.sql.functions import udf, pandas_udf, col
from pyspark.sql.types import DoubleType
import pandas as pd
import re
# ── МЕДЛЕННО: обычный Python UDF (Pickle, строка за строкой) ──
@udf(returnType=DoubleType())
def parse_amount_slow(raw_str: str) -> float:
"""Извлечь число из строки вида '1,234.56 USD'."""
if raw_str is None:
return None
return float(re.sub(r"[^\d.]", "", raw_str.split()[0]))
df.withColumn("amount", parse_amount_slow(col("raw_amount")))
# Каждая строка: JVM → Pickle → сокет → Python → сокет → JVM
# ── БЫСТРО: Pandas UDF (Arrow, батч за раз) ──
@pandas_udf(DoubleType())
def parse_amount_fast(series: pd.Series) -> pd.Series:
"""То же самое, но для целого батча строк."""
return series.str.replace(r"[^\d.]", "", regex=True) \
.str.split().str[0] \
.astype(float)
df.withColumn("amount", parse_amount_fast(col("raw_amount")))
# Батч из 1024 строк: JVM → Arrow RecordBatch → Python → Arrow → JVM
# В 10–100× быстрее обычного UDF
# ── ЛУЧШЕ ВСЕГО: нативные SQL-функции (нет Python overhead вообще) ──
from pyspark.sql.functions import regexp_replace, split
df.withColumn("amount",
regexp_replace(col("raw_amount"), r"[^\d.]", "")
.cast(DoubleType())
)
# Всё выполняется в JVM WholeStageCodegen, Python не задействован
Золотое правило: сначала проверьте, есть ли нативная Spark-функция. Функций в pyspark.sql.functions - более 200. Если нет - используйте Pandas UDF с Arrow. Обычный @udf - только в крайнем случае.
8. Итог: куда движется Spark¶
Тренд последних лет - уход от низкоуровневых API к декларативным. RDD остаётся фундаментом, но прямое использование уходит в нишевые сценарии.
| Период | Основной стиль в PySpark |
|---|---|
| 2012–2014 | RDD - единственный API |
| 2015–2018 | DataFrame вытеснил RDD для ETL и аналитики |
| 2019–2021 | Pandas UDF с Arrow сделал Python UDF приемлемыми |
| 2022–2024 | Pandas API on Spark (pyspark.pandas) для миграции pandas-кода |
| 2025+ | ~95% production PySpark-кода - DataFrame/Spark SQL без RDD |
RDD незаменим когда:
- данные неструктурированы (бинарные файлы, нестандартные форматы)
- нужен кастомный
Partitioner(например, для geo-sharding) - нужны операции без аналога в DataFrame API (
aggregate(),treeReduce()) - интеграция с legacy-системами, возвращающими RDD
Во всех остальных случаях - DataFrame.
Best Practices для PySpark Data Engineer¶
- В PySpark использовать DataFrame, не RDD
- Схему задавать явно для CSV и JSON (
schema=StructType(...)) df.rdd- дорогая операция, избегать в production hot path- Для Python-логики: нативные SQL-функции → Pandas UDF → обычный @udf (в порядке убывания скорости)
explain(mode="formatted")после сложных трансформаций - убедиться, что Catalyst нашёл BroadcastHashJoin и операторы WholeStageCodegen (*)- Off-Heap включить для большинства production-кластеров:
spark.memory.offHeap.enabled=true - Смотреть на Spill в Spark UI: если есть Spill (Disk) → GC давление или мало памяти
9. Домашнее задание¶
Дан legacy-скрипт на PySpark, написанный в стиле Python RDD:
# Legacy код - Python RDD стиль (2014 год)
sc = spark.sparkContext
# Читаем транзакции
raw = sc.textFile("s3a://data/transactions.csv")
# Убираем заголовок
header = raw.first()
data = raw.filter(lambda line: line != header)
# Парсим CSV
parsed = data.map(lambda line: line.split(",")) \
.filter(lambda f: len(f) == 4) \
.map(lambda f: {
"user_id": int(f[0]),
"country": f[1].strip(),
"amount": float(f[2]),
"status": f[3].strip(),
})
# Фильтруем только успешные транзакции из RU
ru_success = parsed.filter(lambda r: r["country"] == "RU" and r["status"] == "success")
# Сумма по пользователям
user_totals = ru_success.map(lambda r: (r["user_id"], r["amount"])) \
.groupByKey() \
.mapValues(sum)
# Только пользователи с суммой > 10000
top_users = user_totals.filter(lambda kv: kv[1] > 10000)
print(top_users.count())
Задание 1 - Аудит. Для каждой строки с трансформацией опишите физические накладные расходы:
- Где происходит Pickle-сериализация?
- Где данные уходят в Python Worker и возвращаются в JVM?
- Почему
groupByKey()хужеreduceByKey()(подсказка: раздел 2 урока 6)? - Сколько раз каждая строка пересекает границу JVM → Python → JVM?
Задание 2 - Переписать. Полностью переписать скрипт на DataFrame API:
- Использовать явную схему (
StructType) при чтении CSV - Заменить все Python UDF нативными Spark-функциями
- Заменить
groupByKey().mapValues(sum)наgroupBy().agg(sum(...)) - Убедиться что нет ни одного
.rddи ни одногоlambdaв финальном коде
Задание 3 - Инспекция плана. Вызвать explain(mode="formatted") на финальном DataFrame и найти:
- Узлы
*(N)- это WholeStageCodegen, перечислить их - Оператор который выполняет filter - убедиться что он применяется до join/groupBy (Predicate Pushdown)
- Тип агрегации -
HashAggregate (partial)+HashAggregate (final)должны быть оба видны - Есть ли
BroadcastExchange? Если нет, почему?
В следующем уроке разберём трансформации и actions подробно: какие операции вызывают shuffle, как читать план из Spark UI и почему порядок filter → join важнее, чем кажется.