Skip to content
Learning Platform

Зачем вообще лезть внутрь

Большинство людей, которые «знают 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(), ничего не вычисляется. Строится DAG — граф преобразований. Реальная работа стартует только при вызове action (count, collect, write, show). В этот момент драйвер передаёт логический план в DAGScheduler, и тот режет его на стадии.

Линия разреза проходит ровно по границам, где требуется shuffle. Чтобы понять, где эти границы, надо различать два класса преобразований.

Narrow-преобразование — каждая выходная партиция зависит ровно от одной входной. Данные не покидают свой узел. map, filter, select, withColumn — узкие. Их можно конвейеризировать одно за другим внутри одной задачи без обмена по сети.

Wide-преобразование — выходная партиция зависит от многих (потенциально всех) входных. Чтобы собрать все записи с одинаковым ключом в одну партицию, данные надо физически перетасовать по сети. groupBy, join, distinct, repartition, оконные функции с partitionBy — широкие. Каждое такое преобразование рвёт граф и порождает новую стадию.

СвойствоNarrowWide
Зависимость партицийодин к одномумногие ко многим
Обмен по сетинетда, через shuffle
Примерыmap, filter, selectgroupBy, join, distinct
Рвёт стадиюнетда
Можно конвейеризироватьданет, барьер

Практический смысл: стадий в вашем job ровно столько, сколько в нём границ shuffle, плюс один. Job с тремя join и одним groupBy — это как минимум пять стадий. Каждая граница стадии — это материализация промежуточного результата на диск и передача по сети. Поэтому первый вопрос при оптимизации — не «сколько памяти дать», а «сколько здесь shuffle и можно ли их убрать или удешевить».

Шаг 2: стадия раскладывается на задачи — по одной на партицию

Стадия — это единица планирования. Единица исполнения — это task. Правило железобетонное: одна задача на одну партицию. Если на входе стадии 200 партиций, Spark создаст 200 задач. Сериализованный код стадии разошлётся по executor’ам, и каждая задача независимо прогонит свою партицию через весь конвейер narrow-преобразований этой стадии.

Отсюда вытекает фундаментальное ограничение, которое ломает половину «оптимизаций»: количество партиций — это потолок параллелизма. Сколько у вас партиций, столько задач может бежать одновременно, не больше. 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 партиции на ядро? Потому что задачи неравны по времени: данные распределены неравномерно, и одна «жирная» партиция может доделываться, когда все остальные ядра уже свободны (это data skew). Если партиций кратно больше, чем ядер, планировщик раздаёт их порциями и сглаживает хвосты. Если же партиций сильно больше, чем нужно (десятки тысяч), накладные расходы на запуск задач и метаданные начинают съедать выигрыш. Золотая середина — порядка 128-256 МБ данных на партицию.

В современных версиях (Spark 3.x и выше) частично спасает AQE — он умеет схлопывать лишние пустые партиции после shuffle по фактической статистике. Но AQE работает с тем, что вы ему дали; он не телепат и не заменяет осознанный выбор числа партиций под форму ваших данных.

Шаг 3: задача живёт в памяти executor’а

Теперь спускаемся на уровень железа. Задача исполняется внутри executor — это JVM-процесс с фиксированным куском памяти. Понимание раскладки этой памяти отделяет тех, кто «увеличивает executor.memory наугад», от тех, кто понимает, почему job всё равно тормозит.

Куча 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 не виден как ошибка. Он виден только как медленность. То, что должно было быть операцией над памятью со скоростью наносекунд, превращается в сериализацию, запись на диск, чтение обратно и десериализацию — на порядки медленнее. В Spark UI это две метрики на странице стадии: 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 состоит из двух фаз, и между ними проходит граница стадий.

  1. Map-side write. Задачи стадии-источника читают свои партиции и для каждой записи вычисляют, в какую из выходных партиций она пойдёт (обычно hash(key) % numPartitions). Записи сортируются по этому номеру партиции и пишутся в локальные файлы на диске того узла, где бежала задача — это shuffle-файлы. Уже здесь видна цена: каждая запись минимум один раз сериализуется и ложится на диск.

  2. 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")

broadcast join — каноничный пример «удешевления». Если одна из таблиц помещается в память executor’а (по умолчанию порог 10 МБ, регулируется spark.sql.autoBroadcastJoinThreshold), Spark рассылает её целиком на каждый узел, и большая таблица джойнится локально, вообще без перетасовки по сети. Один большой shuffle превращается в одну сетевую рассылку маленькой таблицы. На правильных данных это разница между десятью минутами и десятью секундами.

Собираем модель воедино

Прогоним один воображаемый job через всю цепочку, чтобы модель щёлкнула:

  1. Вы пишете readfilterjoingroupBywrite и вызываете write (это action).
  2. DAGScheduler режет план по двум границам shuffle (join и groupBy) — получаем три стадии.
  3. Каждая стадия раскладывается на задачи по числу партиций. После shuffle их ровно spark.sql.shuffle.partitions штук — и если это дефолтные 200 на большом кластере, половина ядер простаивает.
  4. Задачи бегут в execution memory executor’ов. Если данных на партицию много, начинается spill на диск — job не падает, но тихо тормозит.
  5. На границах стадий случается shuffle: map-side write на диск, fetch по сети. Это самая дорогая фаза, и broadcast join либо грамотное число партиций решают почти всё.

Когда такой job тормозит, последовательность диагностики обратная вашему инстинкту. Не лезьте сначала в executor.memory. Откройте Spark UI, найдите медленную стадию, посмотрите на распределение времени задач (перекос?), на метрики Spill (крупные партиции?) и на число задач относительно числа ядер (простой?). В девяти случаях из десяти ответ — поправить число партиций под объём данных и количество ядер, а заодно убрать лишний shuffle. Память — последнее, что стоит крутить, и почти всегда симптом, а не причина.

Куда дальше

Эта модель — фундамент. Поверх неё лежат AQE, dynamic partition pruning, способы лечения перекоса (salting, split skew join), форматы хранения, влияющие на размер партиций при чтении (Parquet и его row-group’ы), и физические планы Catalyst. Но без понимания «job → стадии → задачи → память → shuffle» все эти оптимизации — заклинания.

Если вы пришли из Big Data Engineering и хотите системно — посмотрите направление Data Engineering на нашей платформе: всё внутри-к-железу, с интерактивными песочницами и от руки нарисованными диаграммами.

Полный бесплатный курс с runnable-песочницей, где вы сами увидите стадии, задачи и spill в живом Spark UI: Apache Spark.

Ещё в направлении · Data Engineering

Все материалы направления →