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 координирует.
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). Он пишет промежуточные файлы на диск и шлёт их по сети — медленно, и это самая частая причина тормозов.
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 часто велик для малых данных и мал для огромных).
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
| Strategy | When Spark picks it | Cost |
|---|---|---|
| Broadcast hash join | One side is small (< autoBroadcastJoinThreshold, default 10MB) | No shuffle — small side copied to every executor. Fastest. |
| Sort-merge join | Both sides large | Shuffles & sorts both sides by key. Default for big joins. |
| Shuffle hash join | One side medium, no sort needed | Shuffles, 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")
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 час.
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) | |
|---|---|---|
| Direction | Up or down | Down only |
| Shuffle? | Yes (full) | No (merges adjacent) |
| Result balance | Even partitions | Can be uneven |
| Use for | Increase parallelism / re-key | Reduce file count before write |
repartition(n) | coalesce(n) | |
|---|---|---|
| Направление | Вверх или вниз | Только вниз |
| Shuffle? | Да (полный) | Нет (сливает соседние) |
| Баланс результата | Ровные партиции | Может быть неровным |
| Когда | Поднять параллелизм / сменить ключ | Сократить число файлов перед записью |
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 сливает их автоматически.
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
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.BroadcastHashJoinvsSortMergeJoin— 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. Считай их; каждый дорогой.BroadcastHashJoinvsSortMergeJoin— broadcast реально сработал?PushedFilters/ отсечение партиций — твой фильтр доходит до источника?- Spark UI: одна task сильно медленнее остальных → перекос. Сильный spill → нехватка памяти / нужен repartition.
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.Ответь про себя, потом нажми, чтобы перевернуть.
broadcast(small_dim) → broadcast hash join, no shuffle. The small side is copied to every executor.broadcast(small_dim) → broadcast hash join, без shuffle. Малая сторона копируется на каждый executor.F.*; if you must, use a vectorized pandas_udf (Arrow batches, 10–100× faster).Он непрозрачен для Catalyst и сериализует каждую строку JVM↔Python. Предпочитай встроенные F.*; если нужно — векторизованный pandas_udf (пачки Arrow, в 10–100× быстрее).explain() and try to reduce them.Это shuffle. Каждый Exchange — перераспределение данных по сети; считай их в explain() и старайся сократить.