Зачем вообще лезть внутрь
Большинство людей, которые «знают Spark», знают его API: read, filter, groupBy, join, write. Этого хватает, пока job выполняется за приемлемое время. А потом наступает день, когда тот же самый код на тех же данных идёт час вместо пяти минут, падает с OutOfMemoryError, или — что хуже — отрабатывает успешно, но непонятно медленно, и Spark UI показывает, что 399 ядер из 400 простаивают.
В этот момент знание API заканчивается и начинается знание механики. Чтобы крутить правильные ручки, нужна корректная ментальная модель того, что физически происходит между вызовом action и записью результата на диск. Эта статья выстраивает такую модель сверху вниз: job, стадии, задачи, память executor’а и тот самый shuffle, который связывает всё воедино. Сразу предупреждаю про главный вывод, чтобы вы читали его в голове на протяжении всего текста: настоящий рычаг производительности — это количество партиций относительно объёма данных и числа ядер, а вовсе не executor.memory, который все крутят первым.
Если хочется сразу потрогать всё руками, у нас есть полный бесплатный курс по Spark с runnable-песочницей прямо в браузере — но и без него статья самодостаточна.
Шаг 1: job режется на стадии по границам shuffle
Spark ленив. Когда вы пишете df.filter(...).select(...).groupBy(...).count(), ничего не вычисляется. Строится action (count, collect, write, show). В этот момент драйвер передаёт логический план в DAGScheduler, и тот режет его на стадии.
Линия разреза проходит ровно по границам, где требуется
Narrow-преобразование — каждая выходная партиция зависит ровно от одной входной. Данные не покидают свой узел. map, filter, select, withColumn — узкие. Их можно конвейеризировать одно за другим внутри одной задачи без обмена по сети.
Wide-преобразование — выходная партиция зависит от многих (потенциально всех) входных. Чтобы собрать все записи с одинаковым ключом в одну партицию, данные надо физически перетасовать по сети. groupBy, join, distinct, repartition, оконные функции с partitionBy — широкие. Каждое такое преобразование рвёт граф и порождает новую стадию.
| Свойство | Narrow | Wide |
|---|---|---|
| Зависимость партиций | один к одному | многие ко многим |
| Обмен по сети | нет | да, через shuffle |
| Примеры | map, filter, select | groupBy, join, distinct |
| Рвёт стадию | нет | да |
| Можно конвейеризировать | да | нет, барьер |
Практический смысл: стадий в вашем job ровно столько, сколько в нём границ shuffle, плюс один. Job с тремя join и одним groupBy — это как минимум пять стадий. Каждая граница стадии — это материализация промежуточного результата на диск и передача по сети. Поэтому первый вопрос при оптимизации — не «сколько памяти дать», а «сколько здесь shuffle и можно ли их убрать или удешевить».
Шаг 2: стадия раскладывается на задачи — по одной на партицию
Стадия — это единица планирования. Единица исполнения — это
Отсюда вытекает фундаментальное ограничение, которое ломает половину «оптимизаций»: количество партиций — это потолок параллелизма. Сколько у вас партиций, столько задач может бежать одновременно, не больше. Executor с восемью ядрами может крутить восемь задач параллельно. Если у вас 400 ядер в кластере, но входные данные нарезаны на 200 партиций — бежит 200 задач, и 200 ядер просто простаивают. Вы платите за кластер вдвое больше, чем используете.
Это и есть классическая проблема «200 партиций на 400 ядер». Откуда вообще берётся число 200? Это историческое значение по умолчанию параметра spark.sql.shuffle.partitions — количество партиций на выходе любого shuffle. Если его не трогать, после первого же groupBy или join ваш датасет, каким бы большим он ни был, превращается ровно в 200 партиций. На маленьком кластере это нормально. На большом — половина ресурсов простаивает.
# Антипаттерн: 400 ядер, но shuffle нарезает на 200 партиций
spark.conf.set("spark.sql.shuffle.partitions", "200") # дефолт
# Эвристика: 2-4 задачи на ядро, чтобы выровнять
# хвосты и дать планировщику запас на retry
total_cores = 400
spark.conf.set("spark.sql.shuffle.partitions", str(total_cores * 3)) # 1200
Почему не ровно 400, а 2-4 партиции на ядро? Потому что задачи неравны по времени: данные распределены неравномерно, и одна «жирная» партиция может доделываться, когда все остальные ядра уже свободны (это
В современных версиях (Spark 3.x и выше) частично спасает
Шаг 3: задача живёт в памяти executor’а
Теперь спускаемся на уровень железа. Задача исполняется внутри
Куча executor’а (то, что вы задаёте через spark.executor.memory) делится так. Сначала откусывается reserved memory — фиксированные 300 МБ под внутренние нужды Spark. От остатка берётся доля spark.memory.fraction (по умолчанию 0.6) — это unified memory, главная арена. Оставшиеся примерно 40% — user memory: туда идут ваши структуры данных, объекты в UDF, всё, что вы создаёте руками вне управляемых Spark буферов.
Unified memory, в свою очередь, делится на два региона с подвижной границей:
| Регион | Назначение | Что будет при нехватке |
|---|---|---|
| Execution | буферы под shuffle, сортировку, join, агрегацию | spill на диск |
| Storage | кэш (cache, persist) и broadcast-переменные | вытеснение блоков из кэша |
Ключевое слово — «unified», объединённая. Граница между execution и storage не жёсткая: они занимают друг у друга память. Если кэш простаивает, execution-операция (например, сортировка для shuffle) может занять весь регион. И наоборот, если место нужно под кэш, а execution не использует свою долю, кэш расширяется. Но есть асимметрия: execution имеет приоритет. Блоки кэша можно вытеснить, чтобы освободить место под исполнение; а вот отнять память у активной execution-операции ради кэша нельзя. Поэтому df.cache() под высокой нагрузкой на shuffle ведёт себя непредсказуемо: данные могут молча вылетать из кэша, и при следующем обращении пересчитываться с нуля.
Шаг 4: spill — тихий убийца latency
Что происходит, когда задаче нужно отсортировать или агрегировать больше данных, чем влезает в её долю execution memory? Spark не падает. Он делает spill: сбрасывает часть отсортированных или агрегированных данных на локальный диск, освобождает память, обрабатывает следующую порцию, потом сливает всё обратно. Job завершится успешно — и именно поэтому spill так коварен.
Spill (Memory) и Spill (Disk). Если они не нули — у ваших задач слишком много данных на партицию. И что критично: лечится это не добавлением памяти, а увеличением числа партиций. Больше партиций — меньше данных на задачу — всё влезает в память — spill исчезает.
Это прямая иллюстрация тезиса статьи. Инженер видит медленный groupBy, видит spill, и инстинктивно поднимает executor.memory с 8 ГБ до 16 ГБ. Spill уменьшается, но не исчезает, потому что данные всё ещё нарезаны крупно. А достаточно было поднять spark.sql.shuffle.partitions — и каждая задача стала бы обрабатывать вдвое меньше, аккуратно поместившись в исходные 8 ГБ.
Если же данных на партицию реально слишком много и память кончается совсем — execution не может ни занять у storage, ни сделать spill достаточно быстро — JVM ловит OutOfMemoryError, и менеджер ресурсов (YARN или Kubernetes) убивает контейнер с диагнозом OOMKilled. Задача помечается как упавшая, и DAGScheduler перезапускает стадию. Retry — это не бесплатно: пересчитываются все задачи стадии, повторно читается shuffle-вход, и если причина (крупные партиции, перекос) не устранена, стадия будет падать снова и снова, пока job не упрётся в лимит повторов и не умрёт целиком.
Шаг 5: shuffle — самая дорогая операция
Мы несколько раз упомянули shuffle как границу стадий. Теперь посмотрим, что физически происходит внутри, потому что это самая дорогая вещь, которую делает Spark, и понимание её механики объясняет всё остальное.
Shuffle состоит из двух фаз, и между ними проходит граница стадий.
-
Map-side write. Задачи стадии-источника читают свои партиции и для каждой записи вычисляют, в какую из выходных партиций она пойдёт (обычно
hash(key) % numPartitions). Записи сортируются по этому номеру партиции и пишутся в локальные файлы на диске того узла, где бежала задача — это shuffle-файлы. Уже здесь видна цена: каждая запись минимум один раз сериализуется и ложится на диск. -
Reduce-side fetch. Задачи стадии-приёмника должны собрать свою партицию. Но данные для партиции номер 5 размазаны по shuffle-файлам всех map-задач на всех узлах. Поэтому каждая reduce-задача обращается по сети к каждому исполнителю, тянет свой кусок, складывает в памяти (или спиллит на диск, если не влезает), и только потом начинает считать.
Считаем стоимость. При M map-задачах и R reduce-задачах через сеть тянется до M умножить на R блоков. Все промежуточные данные проходят через диск дважды (запись на map-стороне, возможный spill на reduce-стороне) и через сеть один раз — это самые медленные ресурсы в системе по сравнению с памятью и CPU. Поэтому золотое правило оптимизации Spark: избегайте shuffle, а если не можете избежать — делайте его дешевле.
// Дорого: full shuffle обеих сторон join по сети
val result = bigDf.join(smallDf, "user_id")
// Дёшево: маленькую таблицу рассылаем broadcast'ом
// целиком на каждый executor — shuffle большой таблицы не нужен
import org.apache.spark.sql.functions.broadcast
val result = bigDf.join(broadcast(smallDf), "user_id")
spark.sql.autoBroadcastJoinThreshold), Spark рассылает её целиком на каждый узел, и большая таблица джойнится локально, вообще без перетасовки по сети. Один большой shuffle превращается в одну сетевую рассылку маленькой таблицы. На правильных данных это разница между десятью минутами и десятью секундами.
Собираем модель воедино
Прогоним один воображаемый job через всю цепочку, чтобы модель щёлкнула:
- Вы пишете
read→filter→join→groupBy→writeи вызываетеwrite(это action). DAGSchedulerрежет план по двум границам shuffle (joinиgroupBy) — получаем три стадии.- Каждая стадия раскладывается на задачи по числу партиций. После shuffle их ровно
spark.sql.shuffle.partitionsштук — и если это дефолтные 200 на большом кластере, половина ядер простаивает. - Задачи бегут в execution memory executor’ов. Если данных на партицию много, начинается spill на диск — job не падает, но тихо тормозит.
- На границах стадий случается shuffle: map-side write на диск, fetch по сети. Это самая дорогая фаза, и
broadcast joinлибо грамотное число партиций решают почти всё.
Когда такой job тормозит, последовательность диагностики обратная вашему инстинкту. Не лезьте сначала в executor.memory. Откройте Spark UI, найдите медленную стадию, посмотрите на распределение времени задач (перекос?), на метрики Spill (крупные партиции?) и на число задач относительно числа ядер (простой?). В девяти случаях из десяти ответ — поправить число партиций под объём данных и количество ядер, а заодно убрать лишний shuffle. Память — последнее, что стоит крутить, и почти всегда симптом, а не причина.
Куда дальше
Эта модель — фундамент. Поверх неё лежат AQE, dynamic partition pruning, способы лечения перекоса (salting, split skew join), форматы хранения, влияющие на размер партиций при чтении (
Если вы пришли из Big Data Engineering и хотите системно — посмотрите направление Data Engineering на нашей платформе: всё внутри-к-железу, с интерактивными песочницами и от руки нарисованными диаграммами.
Полный бесплатный курс с runnable-песочницей, где вы сами увидите стадии, задачи и spill в живом Spark UI: Apache Spark.