1 · Fundamentals & Architecture1 · Основы и архитектура
The mental model in one screenМентальная модель на одном экране
| Term | What it is |
|---|---|
| Driver | Runs main, builds plan, schedules, collects results. One per app. Don't collect() big data here. |
| Executor | JVM on a worker; runs tasks, holds cache + shuffle files. Many per app. |
| Job | Triggered by one action. |
| Stage | Set of tasks with no shuffle between them. Boundary = wide dependency. |
| Task | Smallest unit; one task per partition of a stage. |
| Термин | Что это |
|---|---|
| Driver | Исполняет main, строит план, планирует, собирает результаты. Один на приложение. Не делай collect() больших данных сюда. |
| Executor | JVM на воркере; выполняет задачи, хранит кеш + shuffle-файлы. Много на приложение. |
| Job | Запускается одним action. |
| Stage | Набор задач без shuffle между ними. Граница = широкая зависимость. |
| Task | Наименьшая единица; одна задача на партицию стадии. |
2 · RDD vs DataFrame vs Dataset2 · RDD vs DataFrame vs Dataset
| RDD | DataFrame | Dataset | |
|---|---|---|---|
| Schema | No | Yes (Row) | Yes (typed) |
| Type safety | Compile-time | Runtime only | Compile-time |
| Catalyst/Tungsten | No | Yes | Yes |
| Languages | All | All | Scala/Java only |
| Use when | Unstructured, custom partitioner, full control | Default ETL/analytics | Typed pipelines, testability |
| RDD | DataFrame | Dataset | |
|---|---|---|---|
| Схема | Нет | Да (Row) | Да (типизированная) |
| Типобезопасность | Compile-time | Только runtime | Compile-time |
| Catalyst/Tungsten | Нет | Да | Да |
| Языки | Все | Все | Только Scala/Java |
| Когда использовать | Неструктурированные данные, кастомный партиционер, полный контроль | Стандартный ETL/аналитика | Типизированные пайплайны, тестируемость |
One-liner: "Prefer DataFrame/Dataset so Catalyst can optimize; drop to RDD only when you need control the optimizer can't give you."
В одну строку: «Предпочитай DataFrame/Dataset, чтобы Catalyst мог оптимизировать; спускайся к RDD только когда нужен контроль, который оптимизатор дать не может.»
3 · Transformations, Actions, Lazy Eval & the DAG3 · Трансформации, actions, ленивое выполнение и DAG
Transformations (lazy)Трансформации (ленивые)
Build the DAG, return a new RDD/DF, execute nothing.
Строят DAG, возвращают новый RDD/DF, ничего не выполняют.
map filter select join groupBy withColumn repartition
Actions (eager)Actions (нетерпеливые)
Trigger execution; return to driver or write out.
Запускают выполнение; возвращают драйверу или пишут наружу.
count collect take show write foreach reduce first
What happens on an action (the end-to-end walk)Что происходит при action (полный путь)
cache(). count() then write() on an uncached DF = DAG runs twice.cache(). count(), а потом write() на некешированном DF = DAG выполняется дважды.4 · Narrow vs Wide Dependencies & Shuffles4 · Узкие vs широкие зависимости и shuffle
| Narrow | Wide (shuffle) | |
|---|---|---|
| Network | None | Disk write + network fetch |
| Examples | map, filter, union, mapPartitions | groupByKey, reduceByKey, join*, distinct, repartition |
| Cost | Cheap, pipelined | Dominant cost + skew surface |
| Recovery | Recompute 1 partition | May recompute upstream |
| Узкая | Широкая (shuffle) | |
|---|---|---|
| Сеть | Нет | Запись на диск + fetch по сети |
| Примеры | map, filter, union, mapPartitions | groupByKey, reduceByKey, join*, distinct, repartition |
| Цена | Дёшево, конвейерно | Доминирующая стоимость + площадка для перекоса |
| Восстановление | Перевычислить 1 партицию | Может перевычислить upstream |
groupByKey shuffles everything → OOM on hot keys.groupByKey шафлит всё → OOM на горячих ключах.reduceByKey over groupByKey · bucket on join keys · prune columns · right-size partitions.
reduceByKey вместо groupByKey · bucket по ключам джойна · обрезать колонки · правильно подбирать партиции.
5 · Join Strategies5 · Стратегии джойнов
| Strategy | When | Mechanism |
|---|---|---|
| Broadcast Hash (BHJ) | One side < ~10MB (autoBroadcastJoinThreshold) | Ship small side to all executors, build hash map. No big-side shuffle. Fastest. |
| Sort-Merge (SMJ) | Default large–large equi-join | Shuffle both by key, sort, merge. Robust, pays full shuffle+sort. |
| Shuffle Hash (SHJ) | One side fits per-partition; sort too costly | Hash one side per partition, no sort. |
| Broadcast Nested Loop / Cartesian | Non-equi / cross | Fallback; expensive. |
| Стратегия | Когда | Механизм |
|---|---|---|
| Broadcast Hash (BHJ) | Одна сторона < ~10MB (autoBroadcastJoinThreshold) | Отправить маленькую сторону всем executors, построить хеш-карту. Без shuffle большой стороны. Быстрейший. |
| Sort-Merge (SMJ) | По умолчанию для больших equi-join | Shuffle обеих по ключу, сортировка, слияние. Надёжный, платит полный shuffle+сортировка. |
| Shuffle Hash (SHJ) | Одна сторона влезает на партицию; сортировка слишком дорогая | Хеш одной стороны на партицию, без сортировки. |
| Broadcast Nested Loop / Cartesian | Non-equi / cross | Резервный; дорогой. |
# force broadcast (PySpark) from pyspark.sql.functions import broadcast big.join(broadcast(small), "key") -- SQL hint SELECT /*+ BROADCAST(dim) */ ... FROM fact JOIN dim ...
# принудительный broadcast (PySpark) from pyspark.sql.functions import broadcast big.join(broadcast(small), "key") -- SQL hint SELECT /*+ BROADCAST(dim) */ ... FROM fact JOIN dim ...
-1.-1.6 · Catalyst Optimizer & Tungsten6 · Оптимизатор Catalyst и Tungsten
Catalyst = planning
Catalyst = планирование
Extensible rule-based + cost-based optimizer over plan trees.
Расширяемый rule-based + cost-based оптимизатор над деревьями планов.
- Phases: analysis → logical opt → physical planning → codegen
- Rules: predicate pushdown, column pruning, constant folding, filter reorder, partition pruning, join reordering (CBO)
- CBO needs stats:
ANALYZE TABLE … COMPUTE STATISTICS
- Фазы: анализ → логическая опт. → физическое планирование → codegen
- Правила: predicate pushdown, обрезка колонок, сворачивание констант, переупорядочивание фильтров, partition pruning, переупорядочивание джойнов (CBO)
- CBO нужна статистика:
ANALYZE TABLE … COMPUTE STATISTICS
Tungsten = execution
Tungsten = выполнение
- Off-heap binary memory (no JVM object overhead/GC)
- Whole-stage codegen: fuse operator chain into one generated Java function
- Cache-aware + vectorized/columnar Parquet reads
- Off-heap бинарная память (без оверхеда JVM-объектов/GC)
- Whole-stage codegen: сливает цепочку операторов в одну сгенерированную Java-функцию
- Cache-aware + векторизованное/колоночное чтение Parquet
One-liner: "Catalyst decides what plan; Tungsten makes its execution fast and memory-lean."
В одну строку: «Catalyst решает, какой план; Tungsten делает его выполнение быстрым и экономным по памяти.»
df.explain(True) / explain("formatted") → look for Exchange (=shuffle), BroadcastHashJoin, *(n) codegen stages, pushed filters.df.explain(True) / explain("formatted") → искать Exchange (=shuffle), BroadcastHashJoin, *(n) стадии codegen, pushed-фильтры.7 · Caching & Persistence7 · Кеширование и персистентность
cache()=persist(MEMORY_AND_DISK)for DataFrames (RDD cache = MEMORY_ONLY).- Levels:
MEMORY_ONLY·MEMORY_AND_DISK·*_SER·DISK_ONLY·*_2(replicated) · off-heap. - Cache when a DF is reused across multiple actions / iterations (ML loops). Don't cache a once-through pipeline.
cache()=persist(MEMORY_AND_DISK)для DataFrames (кеш RDD = MEMORY_ONLY).- Уровни:
MEMORY_ONLY·MEMORY_AND_DISK·*_SER·DISK_ONLY·*_2(реплицированный) · off-heap. - Кешировать, когда DF переиспользуется через несколько actions / итераций (ML-циклы). Не кешировать одноразовый пайплайн.
unpersist() leaks memory.unpersist() течёт память.8 · Partitioning & Pruning8 · Партиционирование и обрезка
| Concept | Meaning |
|---|---|
| Input partitions | From source (file/block size, maxPartitionBytes) |
| Shuffle partitions | spark.sql.shuffle.partitions (default 200) — AQE coalesces now |
| repartition(n[,col]) | Full shuffle; increase/even out; can hash by col |
| coalesce(n) | No shuffle; only merges (reduce); can cause skew |
| write.partitionBy(col) | On-disk directory layout for later pruning (different from repartition) |
| Концепция | Смысл |
|---|---|
| Входные партиции | Из источника (размер файла/блока, maxPartitionBytes) |
| Shuffle-партиции | spark.sql.shuffle.partitions (по умолчанию 200) — AQE сливает теперь |
| repartition(n[,col]) | Полный shuffle; увеличить/выровнять; можно хешировать по col |
| coalesce(n) | Без shuffle; только слияние (уменьшение); может вызвать перекос |
| write.partitionBy(col) | На-диске структура папок для последующего pruning (отличается от repartition) |
Partition pruning: physically partitioned table (by
ds/date) → read only matching dirs.Dynamic Partition Pruning: prune the big fact table at runtime using keys from the small broadcast dim side — star-schema win.
Partition pruning: физически партиционированная таблица (по
ds/дате) → читать только подходящие папки.Dynamic Partition Pruning: обрезать большую fact-таблицу в runtime, используя ключи из маленькой broadcast dim-стороны — выигрыш для star-схемы.
ds) tables, reading only relevant partitions = daily bread; same idea as Hive/Iceberg pruning.ds), чтение только релевантных партиций = ежедневный хлеб; та же идея, что Hive/Iceberg pruning.9 · Skew Handling & AQE9 · Работа с перекосом и AQE
AQE — Adaptive Query Execution (on by default 3.2+)
AQE — Adaptive Query Execution (по умолчанию с 3.2+)
Re-optimizes the physical plan at runtime from real shuffle stats:
Переоптимизирует физический план в runtime из реальной статистики shuffle:
- Coalesce shuffle partitions — collapse tiny post-shuffle partitions (kills the "200 default, mostly empty" problem)
- Switch join strategy — SMJ → BHJ when a side is small after filters
- Optimize skew joins — split skewed partitions into sub-partitions
- Слияние shuffle-партиций — сворачивает крошечные партиции после shuffle (убивает проблему «200 по умолчанию, в основном пустых»)
- Переключение стратегии джойна — SMJ → BHJ, когда сторона маленькая после фильтров
- Оптимизация перекоса в джойнах — разбивает перекошенные партиции на подпартиции
Skew toolkit (concrete)
Инструментарий для перекоса (конкретно)
- AQE skew join — first line of defense.
- Salt the hot key: suffix
0..Non both sides (replicate small side across salts), join on salted key, aggregate back. - Broadcast the small side → remove the shuffle entirely.
- Isolate hot keys: process top-N separately (often broadcast), union the rest.
- Pre-aggregate (map-side combine) before the shuffle.
- Handle nulls: null join keys all hash to one partition — filter/coalesce them.
- AQE skew join — первая линия защиты.
- Соление (salt) горячего ключа: добавить суффикс
0..Nс обеих сторон (реплицировать маленькую сторону по солям), джойнить по солёному ключу, агрегировать обратно. - Broadcast маленькой стороны → убрать shuffle полностью.
- Изолировать горячие ключи: обрабатывать топ-N отдельно (часто broadcast), объединить остальное.
- Предагрегация (комбинация на стороне map) перед shuffle.
- Обработать null: null join-ключи все хешируются в одну партицию — фильтровать/сливать их.
10 · Executor Memory Model10 · Модель памяти Executor
- Execution can evict Storage beyond a floor; the reverse is bounded (asymmetric).
- Spill: execution memory exhausted → sort/agg buffers go to disk (slower, avoids OOM).
- OOM causes: huge broadcast · skewed partition too big ·
collect()to driver · too few partitions.
- Execution может вытеснить Storage сверх минимума; обратное ограничено (асимметрично).
- Spill: execution-память исчерпана → буферы сортировки/агрегации уходят на диск (медленнее, избегает OOM).
- Причины OOM: огромный broadcast · перекошенная партиция слишком большая ·
collect()на драйвер · слишком мало партиций.
executor.memory → tune spark.executor.memoryOverhead when PySpark OOMs mysteriously.executor.memory → настраивай spark.executor.memoryOverhead, когда PySpark загадочно OOMится.11 · Spark on YARN / Kubernetes & Fault Tolerance11 · Spark на YARN / Kubernetes и отказоустойчивость
Cluster managers & modes
Менеджеры кластера и режимы
- YARN — Hadoop-native, shared queues/capacity scheduler (Front Tier world)
- Kubernetes — executors as pods, cloud-native isolation, pod templates
- Standalone / Mesos (legacy)
- client mode = driver on submit host (interactive); cluster mode = driver in cluster (prod)
- Dynamic allocation scales executors by pending tasks; needs external shuffle service (or k8s shuffle tracking)
- YARN — нативный для Hadoop, общие очереди/capacity scheduler (мир Front Tier)
- Kubernetes — executors как поды, cloud-native изоляция, шаблоны подов
- Standalone / Mesos (легаси)
- client режим = драйвер на хосте подачи (интерактивный); cluster режим = драйвер в кластере (прод)
- Dynamic allocation масштабирует executors по pending-задачам; нужен external shuffle service (или k8s shuffle tracking)
Fault tolerance
Отказоустойчивость
- Lineage (DAG): recompute lost partitions vs replicate data
- Shuffle files on disk + external shuffle service survive executor loss
- Checkpointing: truncate long lineage to reliable storage (HDFS); required for streaming state
- Task retries + speculative execution for stragglers
- Driver = SPOF → cluster mode + supervision / streaming checkpoints
- Родословная (DAG): перевычислить потерянные партиции vs реплицировать данные
- Shuffle-файлы на диске + external shuffle service переживают потерю executor
- Checkpointing: обрезать длинную родословную в надёжное хранилище (HDFS); требуется для streaming state
- Ретраи задач + speculative execution для отстающих
- Драйвер = SPOF → cluster mode + supervision / streaming checkpoints
12 · Structured Streaming Basics12 · Основы Structured Streaming
- Model: stream = unbounded table; each micro-batch appends rows; same DataFrame API, incremental execution. (Continuous mode = low-latency, experimental.)
- Triggers: default (ASAP) · fixed interval ·
availableNow/once (batch backfill) · continuous - Output modes:
append(new rows; needs watermark for aggs) ·update·complete(whole result, aggs only) - Watermark:
withWatermark("ts","10 minutes")bounds state for late data; later events dropped - State + checkpoint: stateful ops keep a state store, checkpointed → exactly-once with replayable source (Kafka offsets) + idempotent sink
- Модель: поток = неограниченная таблица; каждый микро-батч добавляет строки; тот же API DataFrame, инкрементальное выполнение. (Continuous mode = низкая задержка, экспериментально.)
- Триггеры: по умолчанию (ASAP) · фиксированный интервал ·
availableNow/once (батч-backfill) · continuous - Output modes:
append(новые строки; нужен watermark для агрегатов) ·update·complete(весь результат, только агрегаты) - Watermark:
withWatermark("ts","10 minutes")ограничивает state для поздних данных; более поздние события отбрасываются - State + checkpoint: stateful-операции держат state store, чекпоинтятся → exactly-once с перечитываемым источником (Kafka offsets) + идемпотентным приёмником
13 · Performance Tuning & Debugging Slow Jobs13 · Тюнинг производительности и отладка медленных джобов
The signature DE question: "this stage hangs at 199/200"
Фирменный DE-вопрос: «эта стадия зависла на 199/200»
- Spark UI first — stage timeline + task metrics. 199/200 = classic skew.
- Confirm skew: max vs median shuffle-read / duration (50x = skew).
- Rule out: skewed key (null/hot) · too few partitions · spill · GC · straggler node · small-files explosion.
- Fix by cause: AQE skew / salt / broadcast · filter nulls · raise partitions · more memory · speculative execution.
- Verify in the UI: compare shuffle-read distribution & stage time.
- Spark UI сначала — timeline стадии + метрики задач. 199/200 = классический перекос.
- Подтвердить перекос: макс. vs медиана shuffle-read / продолжительность (50x = перекос).
- Исключить: перекошенный ключ (null/горячий) · слишком мало партиций · spill · GC · отстающий узел · взрыв мелких файлов.
- Исправить по причине: AQE skew / salt / broadcast · фильтровать null · увеличить партиции · больше памяти · speculative execution.
- Проверить в UI: сравнить распределение shuffle-read и время стадии.
General tuning checklist
Общий чек-лист тюнинга
- Read less: column/predicate/partition pruning · columnar Parquet/ORC · compact small files
- Shuffle less: broadcast small joins · pre-aggregate · bucket on join key · avoid needless repartition
- Right-size parallelism: partitions ≈ 2–4× cores · let AQE coalesce
- Fix skew · cache only reused data
- Executor sizing: ~5 cores/executor heuristic (fat executors → GC/HDFS throughput pain) · memoryOverhead · dynamic allocation
- Avoid Python UDFs (serialization boundary) → built-in SQL fns or Arrow/pandas vectorized UDFs
- Measure before/after — never tune blind
- Читать меньше: обрезка колонок/предикатов/партиций · колоночный Parquet/ORC · компактировать мелкие файлы
- Shuffle меньше: broadcast маленьких джойнов · предагрегация · bucket по ключу джойна · избегать ненужных repartition
- Подобрать параллелизм: партиции ≈ 2–4× ядра · пусть AQE сливает
- Исправить перекос · кешировать только переиспользуемые данные
- Размер Executor: эвристика ~5 ядер/executor (жирные executors → боль GC/пропускная способность HDFS) · memoryOverhead · dynamic allocation
- Избегать Python UDFs (граница сериализации) → встроенные SQL-функции или Arrow/pandas векторизованные UDFs
- Измерять до/после — никогда не тюнить вслепую
15 · Self-Quiz — flip to reveal15 · Самопроверка — переверни, чтобы открыть
Tap a card to flip. Answer out loud first.
Нажми на карточку, чтобы перевернуть. Сначала ответь вслух.
spark.sql.autoBroadcastJoinThreshold). Disable with -1. Force with broadcast(df).10MB (spark.sql.autoBroadcastJoinThreshold). Отключить через -1. Принудительно через broadcast(df).