PySpark Performance & TuningPySpark: производительность и тюнинг

Where interviews get serious: the execution model, shuffles, join strategies, skew, partitioning, caching, UDFs, and reading a query plan. The "why is my job slow?" toolkit. Где интервью становится серьёзным: модель выполнения, shuffle, стратегии join, перекос, партиционирование, кэширование, UDF и чтение плана запроса. Набор «почему мой джоб тормозит?».
🟢 EssentialsОсновы 🔴 PerformanceПроизводительность 📘 Spark theoryТеория Spark

1. Execution model: where time goes1. Модель выполнения: куда уходит время

To reason about speed you need the runtime hierarchy. One action → a job → split into stages at shuffle boundaries → each stage is many tasks, one per partition, run by executors. The driver coordinates.

Чтобы рассуждать о скорости, нужна иерархия выполнения. Одно действиеjob → делится на stage'ы на границах shuffle → каждый stage — это много task'ов, по одному на партицию, выполняемых executor'ами. Driver координирует.

action (e.g. write) └─ JOB ├─ STAGE 0 (narrow ops) ──┐ shuffle boundary │ └─ task per partition │ └─ STAGE 1 (after shuffle) ◀┘ └─ task per partition ──▶ run on EXECUTORS (cores)
Narrow vs wide transformationsУзкие vs широкие трансформации Narrow (map, filter, withColumn): each output partition depends on one input partition — no data movement, cheap. Wide (groupBy, join, distinct, repartition, orderBy): output partitions pull from many inputs → a shuffle across the network. Wide ops are where time goes. Узкие (map, filter, withColumn): каждая выходная партиция зависит от одной входной — нет движения данных, дёшево. Широкие (groupBy, join, distinct, repartition, orderBy): выходные партиции тянут из многих входных → shuffle по сети. Широкие операции — там, где уходит время.

2. The shuffle (the #1 cost)2. Shuffle (главная цена)

A shuffle redistributes rows across the cluster so that rows with the same key land on the same partition (needed for join/groupBy). It writes intermediate files to disk and sends them over the network — slow, and the most common reason a job drags.

Shuffle перераспределяет строки по кластеру так, чтобы строки с одинаковым ключом попали в одну партицию (нужно для join/groupBy). Он пишет промежуточные файлы на диск и шлёт их по сети — медленно, и это самая частая причина тормозов.

How to cut shuffleКак сократить shuffle Filter early (less data to move) · broadcast small joins · avoid needless distinct/orderBy · pre-aggregate before a join · reuse the same partitioning so a later join is shuffle-free · tune spark.sql.shuffle.partitions (default 200 is often too high for small data, too low for huge). Фильтруй рано (меньше данных двигать) · broadcast'и маленькие join'ы · избегай лишних distinct/orderBy · агрегируй до join · переиспользуй то же партиционирование, чтобы последующий join был без shuffle · настрой spark.sql.shuffle.partitions (дефолт 200 часто велик для малых данных и мал для огромных).
groupByKey vs reduceByKey (RDD era, still asked)groupByKey vs reduceByKey (эпоха RDD, всё ещё спрашивают) reduceByKey aggregates locally (map-side combine) before the shuffle, moving far less data than groupByKey. DataFrame groupBy().agg() already does map-side partial aggregation for you. reduceByKey агрегирует локально (map-side combine) до shuffle, перемещая намного меньше данных, чем groupByKey. DataFrame groupBy().agg() уже делает частичную агрегацию на map-стороне за тебя.

3. Join strategies3. Стратегии join

StrategyWhen Spark picks itCost
Broadcast hash joinOne side is small (< autoBroadcastJoinThreshold, default 10MB)No shuffle — small side copied to every executor. Fastest.
Sort-merge joinBoth sides largeShuffles & sorts both sides by key. Default for big joins.
Shuffle hash joinOne side medium, no sort neededShuffles, builds a hash table. Less common.
СтратегияКогда Spark её выбираетЦена
Broadcast hash joinОдна сторона маленькая (< autoBroadcastJoinThreshold, дефолт 10MB)Без shuffle — малая сторона копируется на каждый executor. Самый быстрый.
Sort-merge joinОбе стороны большиеShuffle и сортировка обеих сторон по ключу. Дефолт для больших join.
Shuffle hash joinОдна сторона средняя, сортировка не нужнаShuffle, строит хеш-таблицу. Реже.
from pyspark.sql.functions import broadcast
# force a broadcast when you KNOW one side is small
big.join(broadcast(small_dim), on="dim_id", how="left")
The go-to fix for slow joinsГлавный фикс медленных join Joining a huge fact to a small dimension and it's slow? Broadcast the dimension. It turns a shuffle-heavy sort-merge join into a shuffle-free broadcast join. Джойнишь огромный факт с маленькой размерностью и тормозит? Broadcast'ни размерность. Это превращает тяжёлый по shuffle sort-merge join в broadcast join без shuffle.

4. Data skew4. Перекос данных (skew)

Skew = one key has vastly more rows than others (e.g. null user_id, or a mega-customer). After a shuffle, one task gets a giant partition while the rest finish — that one straggler holds up the whole stage. Classic symptom: 199 tasks done in seconds, 1 task runs for an hour.

Перекос = у одного ключа сильно больше строк, чем у других (например null user_id или мега-клиент). После shuffle одной task достаётся огромная партиция, пока остальные закончили — этот «отстающий» держит весь stage. Классический симптом: 199 task за секунды, 1 task час.

FixesРешения AQE skew join (on by default in Spark 3+): spark.sql.adaptive.enabled=true auto-splits skewed partitions. · Salting: add a random suffix to the hot key on both sides to spread it across tasks, then aggregate. · Filter/handle null keys separately. · Broadcast the small side to skip the shuffle entirely. AQE skew join (вкл. по умолчанию в Spark 3+): spark.sql.adaptive.enabled=true авто-дробит перекошенные партиции. · Соление (salting): добавь случайный суффикс к «горячему» ключу с обеих сторон, чтобы разнести по task'ам, потом агрегируй. · Обрабатывай null-ключи отдельно. · Broadcast'ни малую сторону, чтобы вообще убрать shuffle.
# salting sketch: spread a hot key over N buckets
N = 16
big_salted   = big.withColumn("salt", (F.rand() * N).cast("int"))
small_salted = small.withColumn("salt", F.explode(F.array(*[F.lit(i) for i in range(N)])))
big_salted.join(small_salted, on=["key", "salt"])

5. Partitions: repartition vs coalesce5. Партиции: repartition vs coalesce

repartition(n)coalesce(n)
DirectionUp or downDown only
Shuffle?Yes (full)No (merges adjacent)
Result balanceEven partitionsCan be uneven
Use forIncrease parallelism / re-keyReduce file count before write
repartition(n)coalesce(n)
НаправлениеВверх или внизТолько вниз
Shuffle?Да (полный)Нет (сливает соседние)
Баланс результатаРовные партицииМожет быть неровным
КогдаПоднять параллелизм / сменить ключСократить число файлов перед записью
Two partition counts people confuseДва числа партиций, которые путают Input partitions come from file/block sizes. Shuffle partitions = spark.sql.shuffle.partitions (default 200) control parallelism after a wide op. Too few → giant tasks/OOM; too many → tiny tasks + scheduler overhead. AQE coalesces these automatically. Входные партиции идут от размеров файлов/блоков. Shuffle-партиции = spark.sql.shuffle.partitions (дефолт 200) задают параллелизм после широкой операции. Слишком мало → огромные task/OOM; слишком много → крошечные task + накладные расходы планировщика. AQE сливает их автоматически.
Small files problemПроблема мелких файлов Writing thousands of tiny files kills downstream read performance. Before write, coalesce (or repartition by the partition column) so each output file is a healthy ~128MB–1GB. Запись тысяч крошечных файлов убивает скорость последующего чтения. Перед write делай coalesce (или repartition по колонке партиционирования), чтобы каждый файл был здоровым ~128MB–1GB.

6. Caching & persistence6. Кэширование и persist

Because evaluation is lazy, a DataFrame is recomputed from scratch every action. If you reuse one across multiple actions, cache it.

Поскольку вычисления ленивые, DataFrame пересчитывается с нуля при каждом действии. Если переиспользуешь его в нескольких действиях — кэшируй.

clean = expensive_pipeline(df).cache()   # = persist(MEMORY_AND_DISK)
clean.count()        # materializes the cache
a = clean.filter(...).count()
b = clean.groupBy(...).agg(...)   # reuses cache, no recompute
clean.unpersist()    # free it when done
Don't over-cacheНе злоупотребляй кэшем Cache only DataFrames reused ≥2 times. Caching everything fills executor memory and triggers spills/evictions that make things slower. Cache used once = pure overhead. Кэшируй только DataFrame, переиспользуемые ≥2 раз. Кэширование всего забивает память executor'ов и вызывает spill/вытеснение, что делает медленнее. Кэш при одном использовании = чистый оверхед.

7. UDFs: avoid, or use pandas_udf7. UDF: избегай или используй pandas_udf

A plain Python UDF is a black box to Catalyst: it can't optimize through it, and every row is serialized from JVM → Python and back (slow, no pushdown). Order of preference:

Обычный Python-UDF — чёрный ящик для Catalyst: он не может оптимизировать сквозь него, и каждая строка сериализуется JVM → Python и обратно (медленно, без pushdown). Порядок предпочтения:

  • 1. Built-in F.* functions — always try first; fully optimized.
  • 2. pandas_udf (vectorized) — uses Apache Arrow to process batches of rows; 10–100× faster than a row UDF.
  • 3. Plain udf — last resort, for logic with no built-in equivalent.
  • 1. Встроенные F.* — всегда пробуй первыми; полностью оптимизированы.
  • 2. pandas_udf (векторизованный) — через Apache Arrow обрабатывает пачки строк; в 10–100× быстрее построчного UDF.
  • 3. Обычный udf — последнее средство, для логики без встроенного аналога.
from pyspark.sql.functions import pandas_udf
import pandas as pd

@pandas_udf("double")
def norm(x: pd.Series) -> pd.Series:
    return (x - x.mean()) / x.std()   # vectorized, Arrow-backed

df.withColumn("z", norm("amount"))

8. Reading the plan & "why is it slow?"8. Чтение плана и «почему тормозит?»

df.explain(True)   # parsed → analyzed → optimized → physical plan

Look for these in the physical plan and the Spark UI:

Ищи это в физическом плане и в Spark UI:

  • Exchange = a shuffle. Count them; each is expensive.
  • BroadcastHashJoin vs SortMergeJoin — did broadcast actually kick in?
  • PushedFilters / partition pruning — is your filter reaching the source?
  • Spark UI: one task far slower than the rest → skew. Heavy spill → memory pressure / repartition needed.
  • Exchange = shuffle. Считай их; каждый дорогой.
  • BroadcastHashJoin vs SortMergeJoin — broadcast реально сработал?
  • PushedFilters / отсечение партиций — твой фильтр доходит до источника?
  • Spark UI: одна task сильно медленнее остальных → перекос. Сильный spill → нехватка памяти / нужен repartition.
AQE — your free win in Spark 3+AQE — бесплатная победа в Spark 3+ Adaptive Query Execution re-optimizes the plan at runtime using real stats: coalesces shuffle partitions, switches sort-merge → broadcast when a side turns out small, and splits skewed partitions. Keep spark.sql.adaptive.enabled=true. Adaptive Query Execution переоптимизирует план во время выполнения по реальной статистике: сливает shuffle-партиции, меняет sort-merge → broadcast, когда сторона оказалась маленькой, и дробит перекошенные партиции. Держи spark.sql.adaptive.enabled=true.

9. Quick self-check9. Быстрая самопроверка

Answer in your head, then tap to flip.Ответь про себя, потом нажми, чтобы перевернуть.

Narrow vs wide transformation?Узкая vs широкая трансформация?
tapнажми
Narrow (map/filter): output partition ← one input partition, no movement. Wide (join/groupBy/distinct): pulls from many partitions → shuffle over the network. Wide = a stage boundary.Узкая (map/filter): выходная партиция ← одна входная, без движения. Широкая (join/groupBy/distinct): тянет из многих → shuffle по сети. Широкая = граница stage.
Huge fact ⋈ small dim is slow. Fix?Огромный факт ⋈ маленькая размерность тормозит. Что делать?
tapнажми
broadcast(small_dim) → broadcast hash join, no shuffle. The small side is copied to every executor.broadcast(small_dim) → broadcast hash join, без shuffle. Малая сторона копируется на каждый executor.
199 tasks finish fast, 1 runs forever. Why?199 task быстро, 1 бесконечно. Почему?
tapнажми
Skew: one key has most of the rows, so one partition is huge. Fix with AQE skew join, salting, separating null keys, or broadcast.Перекос: у одного ключа большинство строк, поэтому одна партиция огромна. Решение: AQE skew join, соление, отдельная обработка null, или broadcast.
repartition vs coalesce?repartition vs coalesce?
tapнажми
repartition = full shuffle, up or down, even partitions. coalesce = no shuffle, down only, merges existing partitions (may be uneven). Use coalesce before write to cut file count.repartition = полный shuffle, вверх/вниз, ровные партиции. coalesce = без shuffle, только вниз, сливает существующие (может быть неровно). coalesce перед write — сократить число файлов.
Why is a Python UDF slow? Alternative?Почему Python-UDF медленный? Альтернатива?
tapнажми
It's opaque to Catalyst and serializes each row JVM↔Python. Prefer built-in F.*; if you must, use a vectorized pandas_udf (Arrow batches, 10–100× faster).Он непрозрачен для Catalyst и сериализует каждую строку JVM↔Python. Предпочитай встроенные F.*; если нужно — векторизованный pandas_udf (пачки Arrow, в 10–100× быстрее).
When should you cache?Когда кэшировать?
tapнажми
When the same DataFrame feeds ≥2 actions (it's recomputed each action otherwise). Don't cache single-use data — it just wastes executor memory.Когда один DataFrame питает ≥2 действия (иначе пересчитывается каждый раз). Не кэшируй одноразовые данные — это лишь тратит память executor'ов.
What does AQE do?Что делает AQE?
tapнажми
Adaptive Query Execution re-optimizes at runtime with real stats: coalesces shuffle partitions, flips to broadcast join, and splits skewed partitions. On by default in Spark 3+.Adaptive Query Execution переоптимизирует на лету по реальной статистике: сливает shuffle-партиции, переключается на broadcast join, дробит перекошенные партиции. Вкл. по умолчанию в Spark 3+.
What does Exchange mean in a plan?Что значит Exchange в плане?
tapнажми
It's a shuffle. Each Exchange means data is redistributed across the network — count them in explain() and try to reduce them.Это shuffle. Каждый Exchange — перераспределение данных по сети; считай их в explain() и старайся сократить.