RDD, DataFrame и Dataset: эволюция абстракций Spark

Три уровня API Spark изнутри: UnsafeRow и off-heap в DataFrame, Pickle vs Arrow в Python RDD, Encoders в Dataset - и почему в PySpark нет Dataset.

core internals

В предыдущих уроках мы работали с 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-код - отдельный процесс.

Каждая строка данных проходит следующий путь:

  1. JVM читает строку из UnsafeRow-буфера или с диска
  2. JVM сериализует строку через Python Pickle в байтовый поток
  3. Байты передаются через Unix Domain Socket в Python Worker (отдельный процесс)
  4. Python Worker десериализует байты обратно в Python-объект
  5. Python выполняет вашу лямбду над объектом
  6. Python сериализует результат снова через Pickle
  7. Байты передаются обратно в JVM через сокет
  8. 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) 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× 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 важнее, чем кажется.