Что на самом деле происходит, когда вы запускаете Spark job
От ленивого DAG до памяти executor’а и shuffle — ментальная модель, которая нужна, когда job падает с OOM в 2 часа ночи. Бесплатно, без регистрации.
1. Job ленив до вызова action
Когда вы пишете read → filter → join → groupBy → write, ничего не вычисляется. Spark строит DAG. Реальная работа стартует только при action (count, collect, write). В этот момент DAGScheduler режет план на стадии — и разрез проходит ровно по каждой границе shuffle.
2. Стадия раскладывается на задачи — по одной на партицию
Стадия — единица планирования, задача — единица исполнения. Правило железобетонное: одна задача на партицию. Значит количество партиций — это потолок параллелизма. 200 партиций на 400 ядрах = 200 простаивающих ядер, вы платите вдвое. Именно это, а не executor.memory, упускают чаще всего.
1 партиция = 1 задача
3. Задачи живут в памяти executor’а — и спиллят
Каждая задача исполняется в JVM executor’а. Execution и storage делят один unified-регион и занимают друг у друга, причём execution имеет приоритет. Когда задаче нужно больше, чем влезает, Spark делает spill на диск — тихий убийца latency. Лечится это не добавлением памяти, а увеличением числа партиций, чтобы на задачу приходилось меньше данных.
4. Shuffle — самая дорогая операция
На каждой границе стадий map-задачи пишут партиционированные файлы на локальный диск, а reduce-задачи тянут свой кусок по сети с каждого executor’а. Диск дважды, сеть один раз — самые медленные ресурсы. Золотое правило: избегайте shuffle, а если нельзя — удешевляйте его (broadcast join превращает полный shuffle в одну маленькую рассылку).
Настоящий рычаг — число партиций относительно данных и ядер, а не executor.memory.
Посмотрите вживую, бесплатно
Полный курс по Spark бесплатен и крутит настоящую песочницу прямо в браузере — стадии, задачи и spill в живом Spark UI. Чтобы читать, аккаунт не нужен.
Открыть бесплатный курс по Spark → Все курсы Data Engineering