сеньорчикОткрыть в Telegram
← вся теориятеория к собесу · Spark и большие данные

Ядро Spark: RDD и план выполнения

Модель исполнения Spark

Чтобы чинить медленный Spark, надо понимать, как он исполняет код: ленивый DAG, границы стейджей по шаффлу, таски по партициям. Собес проверяет именно эту механику, а не знание API.

Стержень: самая дорогая операция - шаффл, и диагноз всегда в Spark UI, а не в логах драйвера.

// Формулировки: «что такое шаффл?», «чем narrow-трансформация отличается от wide?», «где смотреть, почему job медленный?».

DAG, job, stage, task

Driver строит DAG (directed acyclic graph) трансформаций, action (count, write, collect) запускает job, job делится на stages по границам шаффла, а stage это таски по числу партиций. Параллелизм равен числу партиций стейджа.

Трансформации ленивы, материализуют их только actions. Narrow-трансформации (map, filter) живут внутри стейджа, каждая партиция обрабатывается независимо; wide (groupBy, join, distinct) требуют перетасовки данных по ключу - рождают шаффл и новый стейдж.

// Отсюда правило чтения плана: границы стейджей это шаффлы. Считаешь стейджи - считаешь, сколько раз данные поедут по сети.

job / stage / task
запуск action / кусок до шаффла / работа над одной партицией
narrow / wide
без перетасовки / с шаффлом по ключу

Шаффл, память, партиции

Шаффл - запись промежуточных блоков на диск плюс перегонка по сети между экзекьюторами: самая дорогая операция, и её объём - первое, что смотрят в UI.

Память экзекьютора делится на execution (шаффлы, джойны, сортировки) и storage (кэш). Не хватает execution - данные спиллятся на диск (медленно), совсем не хватает - таск падает с OOM (out of memory).

// Число партиций и есть параллелизм: мало - ядра простаивают, много - оверхед на таски. spark.sql.shuffle.partitions по умолчанию 200 и почти всегда требует тюнинга под реальный объём - 200 на терабайте дают партиции-гиганты, на мегабайтах - 200 пустых тасков.

shuffle
перераспределение данных по ключу между экзекьюторами
spill
сброс не влезших в память данных на диск

Кэш, AQE, Spark UI

cache() и persist() окупаются, только когда датафрейм используется повторно; кэшировать одноразовый - занять память без пользы. checkpoint рубит слишком длинный lineage в итеративных алгоритмах, чтобы сбой не пересчитывал всю цепочку заново.

Catalyst переписывает план, а AQE (adaptive query execution) в рантайме по фактическим статистикам чинит число шаффл-партиций, перекошенные джойны и broadcast. В новых версиях включён по умолчанию, но не отменяет проектирования.

// Spark UI - главный инструмент диагностики: стейджи и их время, объёмы шаффлов, спиллы, перекос длительности тасков. Диагноз всегда там; тюнить конфиги наугад без UI - гадание.

AQE
adaptive query execution - коррекция плана в рантайме
lineage
граф происхождения данных для пересчёта при сбое

Как отвечать: «Что такое шаффл и почему это самое дорогое в Spark?»

Шаффл это перераспределение данных между экзекьюторами по ключу, которое случается на wide-трансформациях: groupBy, join, distinct, repartition. Дорогой он потому, что это одновременно и диск, и сеть: каждый экзекьютор пишет промежуточные блоки на локальный диск, а потом они перегоняются по сети к тем экзекьюторам, которым эти ключи достались. То есть данные материализуются и едут целиком, в отличие от narrow-операций вроде map и filter, которые остаются внутри партиции. Шаффл ещё и задаёт границу стейджа, так что по числу шаффлов я сразу вижу, сколько раз пайплайн гоняет данные туда-сюда. Поэтому первое, что я делаю при оптимизации, смотрю в Spark UI объёмы шаффлов и стараюсь их сократить: отфильтровать раньше, заменить джойн на broadcast, не тасовать лишний раз.

Почему это сильный ответ: назван корень (диск + сеть), связь с wide-операциями и границами стейджей, и переход к практике (UI, сокращение шаффлов) - механика, а не определение из документации.

На чём валят

  • collect() большого результата на драйвер - OOM драйвера.
  • Дефолтные 200 шаффл-партиций на терабайте - партиции-гиганты и спиллы; на мегабайтах - 200 пустых тасков.
  • Кэшировать одноразовый датафрейм - память занята, пользы ноль.
  • Игнорировать Spark UI и «тюнить» наугад конфигами.
  • Длинная цепочка без checkpoint в итеративном алгоритме - пересчёт lineage на каждом сбое.

Проверьте себя

Пять вопросов из банка по этой подтеме. Всего их 14, остальные разбираются в тренажёре.

  1. #spark_core1 / 5
    Чем трансформация отличается от действия (action) в Spark?
    A)Трансформации выполняются сразу же, а действия лишь строят отложенный план
    B)Трансформации ленивы — план; действие (count/collect) его запускает
    C)Между трансформацией и действием в Spark нет никакой разницы
    D)Действия выполняются на драйвере, а трансформации — на диске
    показать ответ и разбор
    +B)Трансформации ленивы — план; действие (count/collect) его запускает

    // разбор: Трансформации (map, filter, join) ленивы: они лишь достраивают план вычислений. План собирается, оптимизируется Catalyst и реально исполняется только при действии (count, collect, write). Ленивость и позволяет Spark переупорядочить и слить операции, применить pushdown и не считать лишнего.

  2. #spark_core2 / 5
    Что такое shuffle в Spark и почему это дорогая операция?
    A)Shuffle — это локальная сортировка данных внутри одной партиции без сети
    B)Shuffle просто кэширует промежуточные данные в памяти для скорости
    C)Перекладка данных между узлами по ключу; бьёт по сети и диску — дорого
    D)Shuffle — способ сжать данные перед записью в распределённое хранилище
    показать ответ и разбор
    +C)Перекладка данных между узлами по ключу; бьёт по сети и диску — дорого

    // разбор: Shuffle перераспределяет данные между узлами по ключу — это делают groupBy, join, distinct, repartition. Он пишет промежуточные данные на диск и гоняет их по сети, поэтому обычно и есть главное узкое место Spark-джобы. Минимизация числа shuffle (broadcast-join, предагрегация) — базовый приём оптимизации.

  3. #spark_core3 / 5
    Чем narrow-зависимость отличается от wide в Spark?
    A)Narrow-зависимость требует передачи данных по сети между узлами кластера перед вычислением
    B)Wide и narrow зависимости — синонимы, это одно и то же в Spark
    C)Narrow-зависимость обязательно порождает дорогой shuffle между стадиями
    D)Narrow: родитель → одна дочерняя партиция (без shuffle); wide тасует данные
    показать ответ и разбор
    +D)Narrow: родитель → одна дочерняя партиция (без shuffle); wide тасует данные

    // разбор: При narrow-зависимости (map, filter) каждая родительская партиция питает ровно одну дочернюю — данные остаются локально, без сети. Wide-зависимость (groupByKey, join) требует данных из многих партиций, то есть shuffle. По границам wide-зависимостей Spark и режет план на стадии (stages): внутри стадии — конвейер narrow-операций.

  4. #spark_core4 / 5
    Зачем нужны распределённые вычисления (кластер) для больших данных?
    A)Чтобы все данные целиком поместились в оперативную память одного мощного сервера
    B)Ради красивой визуализации результатов обработки данных
    C)Данные не влезают в одну машину — их делят между многими узлами
    D)Чтобы отказаться от использования дисков в пользу памяти
    показать ответ и разбор
    +C)Данные не влезают в одну машину — их делят между многими узлами

    // разбор: Когда данные и вычисления не помещаются в одну машину, задачу масштабируют горизонтально: разбивают данные на партиции и обрабатывают параллельно на многих узлах кластера. Это дешевле и надёжнее вертикального масштабирования (бесконечно наращивать одну машину нельзя). Spark/Hadoop — фреймворки именно для такого параллелизма.

  5. #spark_core5 / 5
    Чем driver отличается от executor в Spark?
    A)Driver управляет джобой и строит план; executor'ы считают партиции
    B)Именно executor управляет всей джобой целиком, а driver лишь хранит данные на диске
    C)Driver и executor — это просто два названия одного и того же процесса Spark
    D)Driver обрабатывает данные, а executor нужен для логирования
    показать ответ и разбор
    +A)Driver управляет джобой и строит план; executor'ы считают партиции

    // разбор: Driver — управляющий процесс: строит план из вашего кода, режет его на стадии и задачи, раздаёт executor'ам и собирает результат. Executor'ы — рабочие процессы на узлах, которые исполняют задачи над партициями данных. Поэтому сбор большого результата на driver через collect() опасен: он должен уместиться в память одной машины.

дальше

Теорию прочитали. Навык ставится повторением

В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.