Глоссарий PySpark

Справочник ~100 ключевых терминов: от RDD до Schema Evolution. Используй как шпаргалку на протяжении всего курса.

core reference

Термины сгруппированы по слоям - от модели выполнения к инфраструктуре. Разделы идут от низкого уровня к высокому.


I. Модель выполнения

Action

Операция, которая запускает реальное выполнение вычислений и возвращает результат в Driver или записывает данные во внешнее хранилище. До вызова action Spark только строит граф DAG, ничего не считает.

Примеры: collect(), count(), show(), write.parquet(), foreach().

Transformation

Операция над DataFrame/RDD, которая возвращает новый DataFrame/RDD без запуска вычислений. Spark записывает трансформацию в DAG и откладывает выполнение до action.

Примеры: filter(), select(), groupBy(), join(), withColumn().

Narrow Transformation

Трансформация, где каждая output-партиция зависит только от одной input-партиции. Не вызывает shuffle. Spark склеивает несколько narrow-трансформаций в один проход по данным (pipeline).

Примеры: map(), filter(), withColumn(), select(), flatMap(), union().

Wide Transformation

Трансформация, где output-партиция зависит от нескольких input-партиций из разных Executor-ов. Вызывает shuffle - данные перемещаются по сети. Граница между Stage-ами в DAG.

Примеры: groupBy(), join(), repartition(), distinct(), orderBy().

Разница в зависимостях партиций - главный критерий:

Lazy Evaluation (ленивые вычисления)

Принцип, по которому Spark не выполняет трансформации немедленно, а накапливает их в DAG. Только action форсирует исполнение. Это позволяет Catalyst оптимизировать весь план целиком до первого shuffle.

Immutability (неизменяемость)

DataFrame нельзя изменить на месте. Любая операция создаёт новый DataFrame, старый остаётся нетронутым. Это основа отказоустойчивости: если узел упал, Spark восстанавливает партицию, повторяя шаги по Lineage.

DAG (Directed Acyclic Graph)

Направленный ациклический граф зависимостей между операциями. Spark строит DAG при каждом action: узлы - операции, рёбра - потоки данных. DAG разбивается на Stage-ы по границам wide-трансформаций.

Job

Единица работы, запускаемая одним action. Один вызов df.count() - один Job. В Spark UI видны все Jobs с временем выполнения и количеством Stage-ов.

Stage

Часть Job, в которой все операции narrow (без shuffle). Граница между Stage-ами - wide-трансформация. Задачи внутри Stage выполняются параллельно без сетевого обмена.

Task

Минимальная единица работы: обработка одной партиции на одном Executor. Если Stage работает с 200 партициями - в нём 200 Task-ов, выполняющихся параллельно.

Иерархия выполнения от action до Task:

Shuffle

Процесс перераспределения данных между Executor-ами по одной из wide-трансформаций (groupBy, join, distinct). Три основных источника накладных расходов:

  • Disk IO - каждый Executor сначала записывает свои данные на локальный диск (Shuffle Write), затем читает нужные блоки других Executor-ов (Shuffle Read)
  • Network IO - блоки данных передаются между узлами кластера по сети; при большом объёме это может стать узким местом всего кластера
  • Сериализация/десериализация - каждая запись упаковывается в байты при записи и распаковывается при чтении

Shuffle - главная причина падения производительности; сокращение его объёма - центральный приём оптимизации Spark-приложений.

Механизм shuffle между Stage-ами:

Lineage (родословная)

«Рецепт» данных - история того, из какого источника и какими шагами получен каждый DataFrame. Позволяет Spark восстановить утраченную партицию без хранения промежуточных резервных копий.


II. Архитектура кластера

Application

Программа пользователя, которую запускает Spark. Состоит из Driver-процесса и набора Executor-ов. Одно приложение может последовательно выполнить сотни Job-ов.

Driver

JVM-процесс с кодом пользователя. Строит DAG, планирует Task-и, отслеживает прогресс. В режиме client запускается на машине пользователя, в режиме cluster - на узле кластера. Главный риск: collect() переносит все данные с кластера в память Driver - при большом объёме он «умирает».

Executor

JVM-процесс на рабочем узле (Worker Node). Выполняет Task-и, кэширует данные в памяти и на диске. Каждый Executor держит N параллельных слотов (= ядер). При сбое Executor Driver переназначает его Task-и.

Worker Node

Физическая или виртуальная машина в кластере, где запускаются Executor-ы. Один Worker может содержать несколько Executor-ов.

Cluster Manager

Внешняя система, у которой Spark запрашивает ресурсы для Executor-ов.

Manager Когда использовать
YARN Hadoop-кластеры, Amazon EMR
Kubernetes Современные self-hosted и облачные деплои
Standalone Встроенная очередь Spark, простые сценарии
local / local[N] Локальная разработка

DAG Scheduler

Внутренний компонент Driver-а: разбивает DAG на Stage-ы по границам shuffle и определяет их зависимости. Передаёт Stage-ы Task Scheduler-у.

Task Scheduler

Внутренний компонент Driver-а: принимает Stage от DAG Scheduler, разбивает его на Task-и по числу партиций и распределяет их по свободным Executor-ам с учётом Data Locality.

SparkContext

Точка входа в Spark до версии 2.0. Управляет соединением с кластером. В современном коде доступен через spark.sparkContext.

SparkSession

Единая точка входа с версии 2.0. Объединяет SparkContext, SQLContext и HiveContext.

spark = SparkSession.builder \
    .appName("MyApp") \
    .master("yarn") \
    .config("spark.sql.shuffle.partitions", "200") \
    .getOrCreate()

Slot (Core)

Ресурс внутри Executor для выполнения одной Task. Один слот = одно ядро = одна параллельная задача. Если у Executor 4 ядра - он одновременно обрабатывает 4 партиции.

Data Locality

Принцип «код едет к данным, а не данные к коду». Spark предпочитает запускать Task там, где физически лежит нужная партиция, чтобы избежать сетевой передачи. Настройка: spark.locality.wait.

Speculative Execution

Если задача на одном узле выполняется подозрительно долго, Spark запускает её копию на другом узле. Первая завершившаяся версия принимается, вторая отменяется. Защищает от «зависших» узлов. Включается: spark.speculation = true.

Dynamic Allocation

Автоматическое масштабирование числа Executor-ов в зависимости от нагрузки. При простое Spark отдаёт лишние Executor-ы другим приложениям, при росте очереди Task-ов - запрашивает новые. Работает вместе с External Shuffle Service.

Базовая формула параллелизма: executors × executor-cores = максимальное число одновременных Task. Например, 10 Executor-ов × 4 ядра = 40 параллельных задач. При Dynamic Allocation это число растёт и сокращается вслед за фактической нагрузкой.


III. Структуры данных и API

RDD (Resilient Distributed Dataset)

Низкоуровневая абстракция Spark: неизменяемая коллекция объектов, распределённая по кластеру. Не имеет схемы - просто Python/Java-объекты. Не проходит через Catalyst. В современном коде используется редко - только для задач, которые нельзя выразить через DataFrame API.

DataFrame

Таблица с именованными колонками и схемой. Аналог таблицы в SQL или pandas.DataFrame, но распределённая по кластеру. Проходит через оптимизатор Catalyst. Основной API в современном PySpark.

Dataset

Типизированный DataFrame (строгая типизация на уровне JVM). В PySpark не используется - Python не даёт compile-time type-safety. Актуален только в Scala/Java Spark.

Schema

Структура DataFrame: список колонок с именами, типами и флагом nullable.

from pyspark.sql.types import StructType, StructField, StringType, LongType

schema = StructType([
    StructField("user_id",    LongType(),   nullable=False),
    StructField("event_type", StringType(), nullable=True),
])

StructType / StructField

StructType - класс для описания схемы DataFrame, содержит список StructField. StructField - описание одного поля: имя, тип, nullable. Задавать схему вручную лучше, чем полагаться на inference - это быстрее и надёжнее.

Row

Объект одной строки DataFrame. При df.collect() возвращается список Row. Доступ: row.column_name или row["column_name"].

Column (Expression)

Объект, представляющий колонку или вычисление над ней. Создаётся через col("name"), df["name"] или литерал lit(42). Используется в select(), filter(), withColumn().

Partition

Горизонтальный срез DataFrame, хранящийся на одном Executor. Количество партиций = уровень параллелизма. После shuffle: spark.sql.shuffle.partitions (по умолчанию 200). Правило: 128–256 МБ данных на партицию.

Catalog

Системный реестр метаданных: таблицы, базы данных, представления, функции. В PySpark доступен через spark.catalog. Поддерживает Hive Metastore и REST Catalog (Apache Iceberg).

Spark Core

Базовый движок Apache Spark: RDD API, управление памятью, планирование задач (DAG Scheduler, Task Scheduler) и взаимодействие с кластерными менеджерами. Все остальные компоненты - Spark SQL, Structured Streaming, MLlib, GraphX - построены поверх Spark Core и используют его механизм распределённого выполнения.

MLlib

Встроенная библиотека машинного обучения Spark. Поддерживает алгоритмы классификации, регрессии, кластеризации (k-means), рекомендательных систем (ALS) и инструменты feature engineering: VectorAssembler, StandardScaler, StringIndexer. Современный API работает через DataFrame (pyspark.ml); старый RDD-based пакет pyspark.mllib считается устаревшим.

GraphX

Библиотека обработки графов на Spark (доступна только в Scala/Java). Представляет граф как пару RDD: вершины и рёбра. В PySpark официального аналога нет - для графовых задач используется сторонняя библиотека GraphFrames, работающая поверх DataFrame API.


IV. Оптимизатор и рантайм

Catalyst Optimizer

Правило-базированный и стоимостной оптимизатор запросов Spark SQL. Преобразует logical plan в physical plan: проталкивает фильтры вниз, устраняет лишние колонки, выбирает тип join.

Оптимизация и планы выполнения (Catalyst & Plans)

  1. Catalyst Optimizer - встроенный в Spark декларативный оптимизатор SQL-запросов и DataFrame-трансформаций. Автоматически трансформирует граф вычислений, применяя эвристики и оптимизации (например, Predicate Pushdown).
  2. Logical Plan (Логический план) - абстрактное представление запроса, отражающее что именно должно быть сделано с данными без привязки к физической реализации. Делится на Unresolved (до проверки по каталогу), Analyzed (после валидации схем) и Optimized.
  3. Physical Plan (Физический план) - финальный план выполнения, описывающий как конкретно Spark будет обрабатывать данные на кластере (какие алгоритмы джойнов, чтения и агрегаций будут запущены). Виден при вызове метода .explain().
  4. Adaptive Query Execution (AQE) - механизм динамической оптимизации планов выполнения в процессе работы (Runtime), доступный начиная со Spark 3.x. Пересчитывает логику джойнов, склеивает мелкие партиции и борется с перекосом данных на лету.
  5. Cost-Based Optimizer (CBO) - компонент оптимизатора Catalyst, выбирающий наиболее эффективный физический план (например, порядок джойнов нескольких таблиц) на основе собранной статистики данных (команда ANALYZE TABLE).
  6. Whole-Stage Code Generation (WholeStageCodegen) - фича физического движка Spark (Tungsten), которая компилирует цепочку независимых Spark-операций в единый плоский Java-байткод, минимизируя виртуальные вызовы функций и работу процессора.
  7. Dynamic Partition Pruning (DPP) - оптимизация чтения, при которой Spark на этапе выполнения использует отфильтрованные данные одной таблицы (например, справочника), чтобы пропустить чтение ненужных физических партиций другой (большой таблицы фактов).

Управление памятью и производительностью

  1. Project Tungsten - инициатива по перепроектированию физического движка Spark. Переводит работу с памятью в бинарный вид в обход стандартных Java-объектов (Off-Heap память) и оптимизирует работу кэша процессора (L1/L2/L3).
  2. Data Skew (Перекос данных) - аномалия распределения, при которой один или несколько ключей содержат подавляющий объем данных (например, user_id = 'guest'). Ведет к тому, что один Task на кластере выполняется часами, пока остальные воркеры простаивают.
  3. Salting (Соление ключей) - инженерный метод борьбы с Data Skew. К перекошенному ключу джойна искусственно подмешивается случайный префикс (соль), чтобы распределить данные равномерно по кластеру.
  4. Spill to Disk (Пролив на диск) - защитный, но крайне медленный процесс, когда Spark сбрасывает промежуточные данные (из агрегаций, окон или сортировок) из оперативной памяти Executor на локальный диск воркера из-за нехватки RAM.
  5. Storage Level - конфигурация кэширования через .persist(), определяющая, где и как сохранять DataFrame (в памяти, на диске, в сериализованном бинарном виде _SER или с репликацией на несколько нод _2).
  6. Unpersist - метод явного освобождения ресурсов оперативной памяти (.unpersist(blocking=True)), удаляющий DataFrame из кэша кластера после завершения работы с ним.
  7. Shuffle Partitions - параметр (spark.sql.shuffle.partitions), определяющий количество партиций, создаваемых после широких трансформаций (Join, groupBy). По умолчанию равен 200, требует ручной настройки под объем данных.

Продвинутая семантика и типы (Wrangling)

  1. Higher-Order Functions - встроенные в Spark функции высшего порядка (transform, filter, exists, aggregate), позволяющие обрабатывать массивы (ArrayType) прямо внутри ячеек строк на уровне байт-кода без использования медленных Python UDF.
  2. Left Anti Join / Left Semi Join - фильтрующие типы джойнов. Semi возвращает только те строки левой таблицы, для которых нашлось совпадение в правой. Anti возвращает строки левой, для которых совпадений в правой нет.
  3. Dynamic Pivot - операция разворачивания строк в колонки без явного указания списка уникальных значений. Является антипаттерном, так как заставляет Spark делать скрытый Full Scan данных для генерации схемы.
  4. Posexplode - функция взрыва массива, которая, в отличие от стандартного explode(), возвращает две колонки: порядковый индекс элемента (позицию от 0) и сам элемент. Базовый шаблон для трансформаций Wide-to-Long.
  5. EqualNullSafe (<=>) - NULL-безопасный оператор сравнения в Spark SQL. В отличие от стандартного =, выражение NULL <=> NULL вернет True, что критично при джойнах и фильтрации данных со множественными пустотами.
  6. Snapshot Diffing - паттерн вычисления дельты (изменений) между полными историческими срезами данных (снапшотами) на определенные даты с использованием операторов множеств (exceptAll, intersect).
  7. Predicate Pushdown - оптимизация, при которой фильтрация данных (where, filter) спускается на уровень источника данных или файлового движка (например, Parquet футера), отсекая ненужные строки до их считывания в сеть и память кластера.
  8. Columnar Pruning - оптимизация, позволяющая считывать с диска только те колонки, которые явно указаны в .select(), полностью игнорируя остальные бинарные данные в колоночных форматах хранения.

Этапы плана Catalyst

Этап Что происходит
Unresolved Logical Plan SQL/DataFrame → AST, без проверки имён
Analyzed Logical Plan Разрешение имён колонок и типов через Catalog
Optimized Logical Plan Применение правил (pushdown, constant folding)
Physical Plans Несколько вариантов физического плана
Selected Physical Plan CBO выбирает лучший вариант
Executed Plan Codegen, Tungsten

Пайплайн оптимизатора от SQL-запроса до байткода:

Cost-Based Optimizer (CBO)

Модуль Catalyst, использующий статистику таблиц (количество строк, cardinality колонок) для выбора оптимального физического плана. Например: если после фильтра таблица стала маленькой, CBO заменит Sort-Merge Join на Broadcast Join. Включить сбор статистики: ANALYZE TABLE t COMPUTE STATISTICS.

AQE (Adaptive Query Execution)

Адаптивная оптимизация во время выполнения (Spark 3.0+). Работает по 4-шаговому циклу после каждого shuffle/exchange:

  1. Выполнить текущий Stage
  2. Shuffle/Broadcast + сбор реальной статистики партиций
  3. Re-оптимизация невыполненной части плана на основе реальных данных
  4. Запуск следующего Stage с обновлённым планом

Три ключевые оптимизации AQE:

  • Coalescing Shuffle Partitions - объединяет мелкие shuffle-партиции (было 200 → стало 5) на основе реального размера
  • Switching Join Strategies - переключает Sort-Merge Join на Broadcast Join, если реальный размер таблицы оказался меньше порога
  • Optimizing Skew Joins - автоматически дробит горячую партицию (A0 → A0-0 + A0-1) и дублирует соответствующую часть второй таблицы

Включение и настройка:

spark.conf.set("spark.sql.adaptive.enabled", "true")  # по умолчанию с Spark 3.2
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")

Наиболее эффективен для batch-обработки с join-операциями и неточной статистикой.

DPP (Dynamic Partition Pruning)

При join-е таблицы фактов с маленькой таблицей измерений Spark вычисляет список нужных партиций из dimension-таблицы и передаёт его как subquery прямо в scan таблицы фактов - ненужные партиции не читаются вообще.

Классический сценарий: star schema - FACT_SALES партиционирована по date_id, DIM_DATE содержит фильтр year = 2023. Без DPP читается вся FACT; с DPP - только партиции с нужными date_id. На TPC-DS 1 ТБ некоторые запросы ускоряются в 10–20× (q25: 390 с → 20 с).

# DPP включён по умолчанию в Spark 3.0+
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")

Работает с Parquet, Iceberg, Hive-partitioned таблицами. Требует, чтобы таблица фактов была партиционирована по ключу join.

Tungsten

Движок физического выполнения: управляет памятью напрямую (off-heap), минуя GC Java. Работает прозрачно, не требует настройки.

Code Generation (Codegen / Whole-Stage Codegen)

Catalyst генерирует «сырой» Java-байткод прямо под ваш конкретный запрос - несколько операций объединяются в один цикл без виртуальных вызовов. Результат: скорость, близкая к C. Видно в df.explain(mode="codegen").

Логические оптимизации Catalyst (Rule-Based)

Catalyst применяет набор правил к Optimized Logical Plan до физического планирования:

Правило Суть
Predicate Pushdown Проталкивает WHERE/filter к источнику данных - меньше данных читается
Constant Folding Вычисляет константные выражения при планировании: amount * 1.0amount, 100 + 0100
Column Pruning Убирает неиспользуемые колонки из scan: SELECT id FROM t читает только id
Empty Relation Propagation Если источник пуст (WHERE 1=0) - пропускает join/aggregation
Boolean Expression Simplification a AND TRUE → a, a OR FALSE → a и подобные упрощения

Sort Aggregate vs Hash Aggregate

Два алгоритма groupBy:

  • Hash Aggregate - строит хэш-таблицу в памяти. Быстрый, используется по умолчанию.
  • Sort Aggregate - предварительно сортирует данные. Используется при нехватке памяти для хэш-таблицы.

V. Joins и Shuffle

Broadcast Join (BHJ)

Broadcast Hash Join (BHJ) - самая быстрая стратегия объединения таблиц. Маленькая таблица целиком рассылается на все Executor-ноды по сети, исключая тяжелый этап Shuffle для большой таблицы. Контролируется параметром autoBroadcastJoinThreshold.

Маленькая таблица целиком рассылается каждому Executor - shuffle большой таблицы не нужен. Порог: spark.sql.autoBroadcastJoinThreshold (по умолчанию 10 МБ). Принудительно: broadcast(df).

Sort-Merge Join (SMJ)

Sort-Merge Join (SMJ) - стандартная стратегия соединения двух больших таблиц. Состоит из этапов Shuffle (перенос строк с одинаковыми ключами на одну ноду), Sort (сортировка внутри ноды) и Merge (итеративное слияние списков).

Стандартный join для двух больших таблиц:

  1. Обе стороны shuffle-ятся по ключу
  2. На каждом Executor данные сортируются
  3. Выполняется merge-сортировка

Дорого, но хорошо масштабируется.

Сравнение двух основных стратегий join:

Shuffle Hash Join

Альтернатива SMJ: одна сторона шаффлится и загружается в хэш-таблицу в памяти Executor-а, вторая - сканируется по ней. Быстрее SMJ, но требует больше памяти. Catalyst выбирает его автоматически, если данных мало.

Broadcast Nested Loop Join

Broadcast Nested Loop Join (BNLJ) - низкопроизводительная стратегия джойна, запускаемая при отсутствии условий равенства (non-equi joins). Рассылает одну таблицу и выполняет вложенный цикл строк на воркерах.

Cartesian Product (Cross Join)

Соединение «каждый с каждым» - каждая строка левой таблицы с каждой строкой правой. Размер результата = N × M. Без явного crossJoin() Spark выдаст ошибку, чтобы защитить от случайных декартовых произведений. Потенциальный «убийца» памяти кластера.

Anti Join

Возвращает строки из левой таблицы, которых нет в правой. SQL: LEFT ANTI JOIN. Эффективная альтернатива NOT IN (subquery).

Semi Join

Возвращает строки из левой таблицы, которые есть в правой, но без колонок правой таблицы. SQL: LEFT SEMI JOIN. Аналог EXISTS (subquery).

Data Skew (перекос данных)

Ситуация, когда одна или несколько партиций содержат несоразмерно большой объём данных (например, city = 'Москва' - 80% строк). Один Executor перегружен, остальные простаивают. Spark UI: на вкладке Tasks одна задача занимает в 10× больше времени, чем остальные.

Join Hint

Join Hint - явное указание разработчика для оптимизатора Catalyst (например, /*+ BROADCAST(df) */ или df.hint("broadcast")), форсирующее выбор определенной физической стратегии джойна вопреки автоматическим расчетам.

Salting (соление)

Техника борьбы со skew: к «горячему» ключу добавляется случайный суффикс (Москва_1, Москва_2, ...), что разбивает одну огромную группу на несколько. После агрегации результаты объединяются повторно.

Как salting устраняет перекос - одна горячая партиция разбивается на несколько равных:

Bucketing (бакетинг)

Физическое сохранение данных на диск, заранее разбитых по хэшу ключа. При последующем join двух забакетированных таблиц по тому же ключу Spark пропускает shuffle - данные уже лежат правильно. Настройка: .bucketBy(N, "key").saveAsTable("t").


VI. Память и производительность

Unified Memory Manager

Единый пул памяти Executor-а (Spark 1.6+): Execution и Storage делят один бюджет и занимают друг у друга при необходимости.

Область Описание
Execution Memory Shuffle, sort, join-буферы
Storage Memory Кэш RDD/DataFrame (persist())
User Memory Пользовательские структуры данных
Reserved Memory 300 МБ для системных нужд Spark

Структура памяти Executor-а - Execution и Storage делят один бюджет и могут занимать друг у друга:

Off-heap Memory

Память вне управления Garbage Collector Java. Spark через Tungsten управляет ею напрямую, избегая пауз GC. Включить: spark.memory.offHeap.enabled = true.

Persist / Cache

Сохранение DataFrame в памяти или на диске для повторного использования. Без persist() Spark пересчитывает DataFrame с нуля при каждом action.

from pyspark import StorageLevel

df.persist(StorageLevel.MEMORY_AND_DISK)
df.cache()        # алиас для MEMORY_AND_DISK_DESER
# ... несколько action-ов ...
df.unpersist()    # явно освободить

StorageLevel

Уровень кэширования, определяющий где и как хранятся данные:

Уровень Где хранится Сериализован
MEMORY_ONLY RAM Нет
MEMORY_AND_DISK RAM, переполнение → диск Нет
DISK_ONLY Только диск Да
MEMORY_AND_DISK_SER RAM сжато, переполнение → диск Да

Spill (сброс на диск)

Запись shuffle-буферов или кэша на диск при нехватке Execution Memory. Замедляет работу в 10–100×, но не приводит к ошибке. Частый spill - сигнал увеличить spark.executor.memory или уменьшить spark.sql.shuffle.partitions.

Checkpointing (для длинных DAG)

Жёсткое сохранение данных на HDFS/S3 с обрывом Lineage. В отличие от кэша, Checkpoint гарантирует восстановление после перезапуска. Применяется в итеративных алгоритмах (MLlib, GraphX), где накопленный DAG может обрушить Driver.

spark.sparkContext.setCheckpointDir("s3a://bucket/checkpoints/")
df.checkpoint()

mapPartitions (паттерн переиспользования ресурсов)

Трансформация RDD/DataFrame, передающая итератор всей партиции в функцию - вместо вызова функции на каждую строку отдельно. Позволяет инициализировать дорогие ресурсы (DB-соединение, ML-модель) один раз на партицию вместо N раз на N строк.

def write_partition(rows):
    conn = create_db_connection()  # Один раз на партицию
    for row in rows:
        conn.insert(row)
    conn.close()

df.rdd.mapPartitions(write_partition)  # RDD API
df.foreachPartition(write_partition)   # DataFrame API

Разница с map: map(f) - f вызывается для каждой строки. mapPartitions(f) - f получает итератор всей партиции. Особенно важно для Non-serializable ресурсов (JDBC-соединение нельзя передать через closure - его надо создавать на Executor'е).

Coalesce vs Repartition

coalesce(N) repartition(N)
Shuffle Нет (склеивает соседние) Да (полный перемешывания)
Увеличить N Нет Да
Когда Перед записью (уменьшить мелкие файлы) Для выравнивания данных

partitionBy (дисковое партиционирование)

Физическое разбиение данных на диске по значению колонки при записи. Создаёт иерархию папок вида date=2024-01-01/country=RU/. При чтении Spark использует Partition Pruning и читает только нужные папки.

df.write \
    .partitionBy("date", "country") \
    .parquet("s3a://bucket/events/")
# Результат: events/date=2024-01-01/country=RU/part-0000.parquet

Не путать с repartition()/coalesce() - те управляют числом партиций в памяти во время выполнения, не структурой файлов на диске.

Broadcast Variable

Переменная «только для чтения», которая рассылается на все Executor-ы один раз и живёт там всё время жизни Application. Эффективно для справочников и небольших моделей ML.

lookup = spark.sparkContext.broadcast({"Москва": "RU", "Берлин": "DE"})

@udf("string")
def get_country(city):
    return lookup.value.get(city)

Accumulator

Переменная «только для записи» на Executor-ах - Driver читает итог. Используется для счётчиков (битые строки, пропущенные значения) без collect().

bad_rows = spark.sparkContext.accumulator(0)

def process(row):
    if row["value"] is None:
        bad_rows.add(1)

df.foreach(process)
print(f"Битых строк: {bad_rows.value}")

Garbage Collection (GC) Overhead

Ситуация, когда JVM тратит больше времени на очистку памяти, чем на полезную работу. В Spark UI отображается красными полосками в таймлайне Task-ов. Решение: G1GC (-XX:+UseG1GC), уменьшить spark.storage.memoryFraction, перейти на off-heap.


VII. Форматы и хранилище

Parquet

Колоночный формат: данные одной колонки хранятся вместе. Поддерживает predicate pushdown, column pruning и статистику (min/max) на уровне row-group. Стандарт для Spark в lakehouse-стеке.

Колоночное хранилище читает только нужные колонки - для аналитики в разы быстрее строкового:

SELECT age FROM t читает только age-блок - id и name не трогаются вообще.

Avro

Строковый (row-based) формат. Хранит схему вместе с данными. Хорош для Kafka (схема в каждом сообщении), CDC и систем с частыми изменениями схемы. Медленнее Parquet для аналитических запросов.

ORC

Колоночный формат, популярный в экосистеме Hive. Хорошо интегрирован с YARN/Hive Metastore. По производительности сопоставим с Parquet; в новых проектах Parquet предпочтительнее.

Delta Lake

Надстройка над Parquet с ACID-транзакциями, версионностью и Schema Enforcement:

  • ACID - гарантирует корректность при параллельной записи
  • Time Travel - df = spark.read.format("delta").option("versionAsOf", 5).load(path)
  • OPTIMIZE - компактация мелких файлов
  • Z-Order - кластеризация по нескольким колонкам

Lakehouse

Архитектура, объединяющая Data Lake (дешёвое хранилище, S3/HDFS) и Data Warehouse (ACID, схемы, запросы). Реализуется через table formats: Apache Iceberg, Delta Lake, Apache Hudi.

Predicate Pushdown

Оптимизация: Spark передаёт условие WHERE вниз к хранилищу (Parquet, Iceberg), чтобы не читать ненужные строки. Работает автоматически.

Partition Pruning

При чтении Spark пропускает папки, не удовлетворяющие фильтру по partition-колонке. WHERE date = '2024-01-01' при партиционировании по date читает одну папку вместо всех.

Column Pruning

SELECT id, amount из 100-колоночного Parquet-файла читает физически только 2 колонки. Catalyst автоматически передаёт список нужных колонок в reader.

Z-Ordering

Метод кластеризации данных внутри файлов Delta Lake: строки, близкие по значению указанных колонок, группируются вместе физически. Ускоряет фильтрацию по нескольким колонкам (city + category). Команда: OPTIMIZE t ZORDER BY (city, category).

Small Files Problem

Ситуация, когда на S3/HDFS лежат миллионы файлов по несколько КБ. Spark тратит больше времени на листинг и открытие файлов, чем на их чтение. Возникает при стриминге или неправильном repartition. Лечится: coalesce() перед записью, OPTIMIZE в Delta Lake.

Snappy / Zstd

Алгоритмы сжатия Parquet:

  • Snappy - быстрый, умеренное сжатие, стандарт по умолчанию
  • Zstd - лучшее сжатие при сопоставимой скорости, рекомендуется для cold storage

VIII. PySpark Bridge (Python-слой)

Py4J

Библиотека-«мост»: позволяет Python-коду вызывать методы JVM-объектов Spark. Когда вы пишете df.show(), Py4J отправляет команду в JVM через localhost-сокет. Работает только в Driver-процессе.

Python Worker

При выполнении обычного Python UDF Spark запускает отдельный Python-процесс на каждом Executor. Данные передаются из JVM в этот воркер через сокет с сериализацией Pickle - это главный тормоз row-level UDF.

Разница в передаче данных между JVM и Python:

Pickle Serialization

Стандартный способ сериализации Python-объектов. В Spark используется для передачи данных в/из Python UDF. Медленный и неэффективный по сравнению с Apache Arrow. По возможности заменяйте Pandas UDF.

Apache Arrow

Колоночный формат в памяти с нулевой сериализацией между JVM и Python. Основа Pandas UDF: данные передаются батчами как pd.Series, без построчной упаковки. Также ускоряет df.toPandas() и spark.createDataFrame(pandas_df).

Включить: spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true").

UDF (User Defined Function)

Пользовательская функция Python, зарегистрированная в Spark SQL. Работает построчно через Python Worker - медленнее встроенных функций. Используй только если нет эквивалента в pyspark.sql.functions.

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

@udf(returnType=StringType())
def clean_phone(phone: str) -> str:
    return phone.replace("-", "").strip() if phone else None

Pandas UDF (Vectorized UDF)

UDF, получающий pd.Series вместо скалярного значения. Передача через Arrow - в 10–100 раз быстрее row-level UDF. Рекомендуется для любых Python-вычислений над колонками.

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

@pandas_udf("double")
def fahrenheit_to_celsius(temp: pd.Series) -> pd.Series:
    return (temp - 32) * 5 / 9

Kryo Serializer

Альтернативный JVM-сериализатор: компактнее и быстрее стандартного Java Serializer. Используется для передачи объектов между Driver и Executor-ами (не для данных в DataFrame - там Tungsten). Включить: spark.serializer = org.apache.spark.serializer.KryoSerializer.


IX. Развёртывание и мониторинг

spark-submit

CLI-инструмент для запуска Spark-приложений на кластере.

spark-submit \
  --master yarn \
  --deploy-mode cluster \
  --num-executors 10 \
  --executor-memory 4g \
  --executor-cores 2 \
  my_job.py --date 2024-01-01

Deploy Mode

Режим Где Driver Когда
client На машине пользователя Разработка, отладка; именно в этом режиме работает JupyterHub
cluster На узле кластера (в YARN - внутри Application Master) Production

local[*] / local[N]

Специальный master URL для запуска на одной машине без кластера. local[*] - все ядра CPU, local[4] - 4 ядра. Удобно для тестирования.

Executor Memory Overhead

Память сверх JVM heap: native-память, stack, Python Worker-ы. По умолчанию max(384 МБ, 10% executor.memory). В PySpark критично - именно здесь живут Python Worker-ы. Настройка: spark.executor.memoryOverhead.

External Shuffle Service

Отдельный процесс на каждом узле кластера, хранящий файлы shuffle. Позволяет Executor-ам завершиться (при Dynamic Allocation), не теряя результатов shuffle. Обязателен при Dynamic Allocation на YARN.

Fair Scheduler / FIFO Scheduler

Планировщики задач внутри одного Spark-приложения:

  • FIFO (по умолчанию): первый пришёл - первый обслуживается, один тяжёлый Job блокирует остальных
  • Fair: ресурсы делятся между Pool-ами, короткие запросы не ждут длинных

Spark UI

Веб-интерфейс для мониторинга работающего приложения: Jobs, Stages, Tasks, SQL-планы, использование памяти. По умолчанию на http://driver:4040.

Spark History Server

Веб-интерфейс для просмотра завершённых приложений по сохранённым Event Log-ам. Настройка: spark.eventLog.enabled = true, spark.eventLog.dir = s3a://bucket/logs/.

S3A

Современный протокол для работы Spark с S3-совместимым хранилищем (MinIO, Ceph RGW, AWS S3). Поддерживает multipart upload, список с пагинацией и Magic Committer для атомарной записи. Префикс пути: s3a://bucket/path.

SparkMeasure

Open-source библиотека для программного сбора Stage и Task метрик из кода Spark-приложения. Реализует кастомный Spark Listener с двумя компонентами: StageInfoRecorder и TaskInfoRecorder.

from sparkmeasure import StageMetrics

stage_metrics = StageMetrics(spark)
stage_metrics.begin()
df.groupBy("key").count().collect()  # Ваш код
stage_metrics.end()
stage_metrics.print_report()
# → executorRunTime, shuffleReadBytes, jvmGCTime, peakExecutionMemory ...

FlightRecorder-режим позволяет писать метрики в фоне в файл, InfluxDB или Kafka без явного begin/end.

Ключевые метрики: peakExecutionMemory, inputMetrics.recordsRead, shuffleBytesWritten, jvmGCTime, memoryBytesSpilled.

spark.metrics.namespace

Namespace для метрик Spark MetricsSystem. По умолчанию равен spark.app.id - уникальному ID каждого запуска. Это значит, что каждый запуск создаёт новый namespace в Graphite/InfluxDB и история одного приложения разбивается на тысячи разных серий.

Правильная настройка - использовать имя приложения:

spark.conf.set("spark.metrics.namespace", "${spark.app.name}")
# Или в spark-defaults.conf:
# spark.metrics.namespace=${spark.app.name}

Это позволяет видеть тренды одного приложения (daily_etl_sales) в Grafana через несколько запусков.

Block Manager

Внутреннее хранилище данных на каждом узле: управляет тем, какие блоки (партиции) лежат в RAM, а какие на диске. Каждый Block Manager регистрируется у Driver-а и участвует в broadcast и shuffle.


X. Structured Streaming

Micro-Batch

Основной режим Structured Streaming: поток делится на маленькие батчи с фиксированным интервалом и обрабатывается как обычный Spark job. Latency - секунды. Exactly-once с правильным sink.

Continuous Processing

Экспериментальный режим «настоящего» realtime: задачи на Executor-ах запущены постоянно, latency - миллисекунды. Не поддерживает stateful-операции в полном объёме.

Trigger

Настройка частоты запуска micro-batch:

df.writeStream \
  .trigger(processingTime="30 seconds")  # каждые 30 сек
  # .trigger(once=True)                 # один раз и стоп
  # .trigger(availableNow=True)         # все накопленные данные
  .start()

Watermark

Граница времени, до которой Spark принимает опоздавшие события. Позволяет очищать state stateful-операций. withWatermark("event_time", "10 minutes") - данные старше 10 минут игнорируются.

Output Modes

Режим Что выводится
append Только новые строки (без агрегаций)
update Только изменившиеся строки
complete Вся таблица целиком пересчитывается

Stateful Transformation

Операция, которая «помнит» данные из прошлых батчей. Например, running total или скользящее окно. Spark хранит state в памяти Executor-а и персистирует его в checkpoint location.

Checkpoint Location

Директория (HDFS/S3), куда стрим записывает прогресс (offsets Kafka) и state агрегаций. При перезапуске приложение читает checkpoint и продолжает с места остановки.


XI. Data Engineering

Idempotency (идемпотентность)

Свойство пайплайна: сколько бы раз ни запустить расчёт за один период, результат в хранилище одинаков - без дублей. Достигается через overwrite или MERGE INTO с уникальным ключом.

Schema Evolution

Способность адаптироваться к изменениям схемы источника без падения кода. Delta Lake и Iceberg поддерживают добавление колонок, переименование, изменение типов с контролируемой миграцией.

Exactly-once Semantics

Гарантия: каждое событие обработано ровно один раз - ни потерь, ни дублей. В Structured Streaming достигается сочетанием Kafka offset-ов (at-least-once чтение) + идемпотентного sink (MERGE/overwrite).

Data Lineage (lineage данных)

Карта происхождения данных от источника до конечного отчёта. Позволяет отследить ошибку на любом этапе ETL. Инструменты: OpenLineage, Apache Atlas, DataHub.



XII. Apache Kafka

Message Broker

Промежуточное ПО (middleware), выступающее посредником между Producer и Consumer сообщений. Обеспечивает decoupling (Producer и Consumer не знают друг о друге), надёжное хранение сообщений и гарантии доставки. Два типа: Point-to-Point (FIFO, одно сообщение → один Consumer) и Publish/Subscribe (одно сообщение → все подписчики Topic'а).

Apache Kafka (брокер)

Распределённый брокер сообщений с моделью Publish/Subscribe, созданный LinkedIn в 2011 году. Назван в честь Франца Кафки («оптимизирован для записи»). Основные характеристики: append-only log (сообщения только добавляются), высокая пропускная способность (sequential disk writes), горизонтальное масштабирование через партиции.

Topic (Kafka)

Логическая очередь сообщений в Kafka. Append-only - нельзя удалить отдельное сообщение. Состоит из одной или нескольких физических Partition. Consumer'ы подписываются на Topic по имени.

Partition (Kafka)

Физическая единица хранения внутри Topic. Упорядоченная неизменяемая последовательность сообщений. Каждая партиция имеет один активный Leader (обрабатывает reads/writes) и набор Follower-реплик. Роутинг по ключу: hash(key) % num_partitions; без ключа - round-robin. Уровень параллелизма Consumer Group = число партиций.

Segment (Kafka)

Физический файл данных внутри Partition. Один активный сегмент принимает новые записи; при достижении log.segment.bytes (default 1 GB) закрывается и создаётся новый. Закрытые сегменты подлежат retention/compaction политике.

Broker (Kafka)

Один Kafka-сервер. Хранит лидерские и реплицированные партиции. При падении Leader-брокера один из Follower'ов становится новым Leader автоматически.

Offset (Kafka)

Монотонно возрастающий порядковый номер сообщения в Partition. Уникален в рамках одной партиции. Consumer Group хранит committed offset для каждой партиции - это «закладка», с которой продолжится чтение после рестарта.

Consumer Group (Kafka)

Группа Consumer'ов, совместно читающих Topic. Одна партиция назначается ровно одному Consumer в группе - это обеспечивает параллелизм без дублей. Разные Consumer Groups читают Topic независимо и полностью (каждая группа получает все сообщения).

Log Retention (Kafka)

Политика хранения данных в Kafka. Два типа: по времени (log.retention.hours=168, приоритет: ms > minutes > hours) и по размеру (log.retention.bytes). По умолчанию: 7 дней. Срабатывает то условие, которое выполнилось первым.

Cleanup Policy (Kafka)

Политика очистки закрытых сегментов:

Policy Суть
delete (default) Физически удаляет сегменты старше retention
compact Оставляет только последнее сообщение для каждого ключа
delete,compact Сначала compact, потом delete по retention

Log Compaction (Kafka)

Механизм очистки compact policy. Cleaner читает Log Head, собирает последний offset для каждого ключа, копирует в новый сегмент только актуальные сообщения, атомарно заменяет старые. В Log Tail (compacted части) возможны «дыры» в offset'ах. Применяется для changelog-таблиц и key-value state.

Окружение Data Lakehouse и Ingestion

  1. Schema Evolution - возможность транзакционных табличных форматов (Delta Lake, Iceberg) или файлов (Avro, Parquet) автоматически расширять и изменять свою структуру (метаданные) при записи новых колонок без перезаписи исторических данных.
  2. Schema Enforcement - защитный механизм Lakehouse-таблиц, блокирующий запись датафрейма, если его схема (типы данных или состав колонок) не совпадает со схемой целевой таблицы. Предотвращает data corruption.
  3. Read Mode (PERMISSIVE / FAILFAST / DROPMALFORMED) - конфигурации парсинга грязных текстовых данных (CSV, JSON). Определяют поведение Spark при встрече с битой строкой: замена на NULL (PERMISSIVE), падение (FAILFAST) или игнорирование (DROPMALFORMED).
  4. _corrupt_record - специальная системная колонка, создаваемая Spark при чтении в режиме PERMISSIVE, куда автоматически складывается исходный сырой текст испорченной строки, не прошедшей валидацию по схеме.
  5. badRecordsPath - параметр чтения, указывающий Spark путь в Data Lake (S3/HDFS), куда необходимо автоматически складировать сырые логи и файлы ошибок парсинга (Dead Letter Queue паттерн).
  6. Testable Transforms - архитектурный паттерн написания Spark-приложений, требующий отделения бизнес-логики (чистых функций трансформации DataFrame-в-DataFrame) от кода ввода-вывода (I/O) для обеспечения возможности юнит-тестирования.
  7. Chispa - специализированная библиотека для тестирования PySpark-кода в связке с Pytest. Предоставляет производительные методы детального сравнения схем и контента датафреймов (assert_df_equality).
  8. Z-Ordering - алгоритм многомерной кластеризации и физической сортировки данных на диске в Delta Lake / Iceberg. Группирует похожую информацию в одних файлах, обеспечивая максимальный эффект для Data Skipping.
  9. Manifest Files - технические файлы метаданных, содержащие точные списки путей к физическим файлам данных (Parquet), актуальным для конкретного среза времени (используются в Apache Iceberg / Delta Lake).
  10. Time Travel - возможность транзакционных движков Data Lakehouse читать состояние таблицы на определенную точку в прошлом (по Timestamp или по ID версии коммита), обращаясь к историческим файлам через Transaction Log.

Быстрая шпаргалка

Термин Одной строкой
Action Запускает выполнение, возвращает результат
Transformation Описывает операцию, не выполняет её
Lazy Evaluation Вычисления откладываются до action
Immutability DataFrame неизменяем, операция создаёт новый
DAG Граф операций; разбивается на Stage по shuffle
Job Одна action = один Job
Stage Набор narrow-трансформаций без shuffle
Task Обработка одной партиции
Shuffle Disk IO + Network IO + сериализация между Executor-ами
Lineage История создания данных для восстановления
Driver Процесс с кодом пользователя и планировщиком
Executor Процесс на рабочем узле, выполняет Task-и
Slot Одно ядро Executor-а = одна параллельная Task
Data Locality Код едет к данным, а не данные к коду
Speculative Execution Дублирует медленную Task на другом узле
Dynamic Allocation Авто-масштабирование числа Executor-ов
Partition Горизонтальный срез DataFrame на одном Executor
Catalyst Оптимизатор запросов Spark SQL
CBO Выбирает план по статистике таблиц
AQE Адаптивная оптимизация во время выполнения
DPP Динамическое отсечение ненужных партиций при join
Codegen Генерация байткода под конкретный запрос
Tungsten Управление памятью вне GC, почти как C
Broadcast Join Join без shuffle для маленьких таблиц
Sort-Merge Join Стандартный join двух больших таблиц с shuffle
Data Skew Перекос: одна партиция несоразмерно больше
Salting Добавление суффикса к ключу для борьбы со skew
Bucketing Физическое хранение, уже разбитое по ключу
Spill Запись на диск при нехватке памяти
Persist / Cache Сохранение DataFrame для повторного использования
Coalesce Уменьшение партиций без shuffle
Repartition Изменение партиций с shuffle
partitionBy Физические папки на диске при записи по колонке
Broadcast Variable Read-only переменная на всех Executor-ах
Accumulator Write-only счётчик на Executor-ах
Predicate Pushdown Фильтр передаётся в хранилище, не в Spark
Partition Pruning Чтение только нужных папок
Column Pruning Чтение только нужных колонок
Parquet Колоночный формат, стандарт для аналитики
Delta Lake Parquet + ACID + Time Travel
Z-Ordering Кластеризация данных по нескольким колонкам
Small Files Problem Миллионы мелких файлов = тормоз листинга
Py4J Мост Python ↔ JVM в Driver-е
Python Worker Отдельный Python-процесс на Executor для UDF
Pickle Медленная сериализация для row-level UDF
Pandas UDF Быстрый Python-UDF через Arrow-батчи
Exactly-once Каждое событие обработано ровно один раз
Idempotency Повторный запуск за период даёт тот же результат
Schema Evolution Адаптация к изменениям схемы источника
Watermark Граница ожидания опоздавших событий в стриминге
Micro-Batch Стриминг как серия мини-batch заданий
Stateful Transformation Операция, хранящая состояние между батчами
Spark Core Базовый движок: RDD, DAG, планировщик, Cluster Manager
MLlib Встроенная ML-библиотека (DataFrame API: pyspark.ml)
GraphX Библиотека графов Spark (только Scala/Java)
Логические оптимизации Catalyst Predicate Pushdown, Constant Folding, Column Pruning, Empty Relation, Boolean Simplification
mapPartitions Паттерн: один ресурс (DB-коннект) на всю партицию, не на каждую строку
SparkMeasure Программный сбор Stage/Task метрик через кастомный Listener
spark.metrics.namespace Namespace метрик; ставить = ${spark.app.name}, а не app.id
Message Broker Middleware: decoupling Producer и Consumer через очередь
Kafka Topic Логическая append-only очередь; состоит из Partition
Kafka Partition Физический лог; Leader + Followers; hash(key) % N
Kafka Offset Монотонный порядковый номер сообщения в партиции
Consumer Group Группа Consumer'ов; одна партиция = один Consumer в группе
Log Retention Хранение по времени (default 7 дней) или по размеру
Log Compaction Compact policy: оставляет только последнее сообщение на ключ
DPP Star schema: фильтр DIM → отсечение партиций FACT при scan
AQE 4-step cycle Execute Stage → collect stats → re-optimize → next Stage