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

Spark: ленивость и действия

Зачем это спрашивают

Ленивость - модель исполнения Spark, и без неё не объяснить ни «почему ошибка вылетела на count», ни «почему всё пересчиталось дважды». Вопрос-детектор реального опыта.

// Мантра: трансформации строят план, действия его исполняют.

Трансформации против действий

select, filter, join, groupBy - ленивые трансформации: они только наращивают план. Исполнение запускают действия: count, collect, write, show.

Следствие для отладки: ошибка в трансформации всплывёт на первом действии - стектрейс укажет на count(), а не на кривую строку тремя выше. Отлаживаются цепочки промежуточными действиями по шагам.

// Зачем так: оптимизатор Catalyst видит план целиком и переписывает его - толкает предикаты к чтению, отсекает колонки, переупорядочивает джойны.

transformation / action
строит план / исполняет план
Catalyst
оптимизатор: переписывает весь план до исполнения

Пересчёт графа и cache

Каждое действие пересчитывает граф заново: два count() по одной цепочке - два полных прохода по данным. Цикл с действием внутри - пересчёт всей цепочки на каждом витке.

Переиспользуешь результат - cache()/persist(). Но кэш не бесплатен: ест память экзекьюторов и окупается от двух переиспользований; кэшировать всё подряд - вымывание памяти и спиллы вместо ускорения.

// unpersist() после использования - привычка, отличающая аккуратного пользователя кластера.

cache / persist
материализация для переиспользования; ≥2 чтений

explain: где шаффлы

explain() показывает физический план: где шаффлы (узлы Exchange), какой алгоритм джойна выбран, дошли ли фильтры до сканирования.

Wide-трансформации (groupBy, join, distinct, repartition) требуют шаффла и режут план на стейджи; narrow (filter, select) исполняются без движения данных.

// Читать explain - базовый навык: «фильтр по партиционной колонке не дочитался до скана» видно только там, а стоит это чтения всей таблицы.

wide / narrow
с шаффлом (границы стейджей) / без

Как отвечать: «Почему ошибка вылетела на count(), хотя кривая строка - filter выше?»

Потому что Spark ленив: filter - трансформация, она лишь добавила узел в план и ничего не исполнила. Исполнение запустил count() - первое действие в цепочке, на нём весь план и провалился, поэтому стектрейс указывает туда. Это не баг, а плата за оптимизацию: Catalyst переписывает план целиком до исполнения. Отлаживаю такие цепочки, вызывая лёгкие действия по шагам после подозрительных трансформаций. И помню обратную сторону ленивости: каждое действие пересчитывает граф заново - если дальше нужны ещё проходы, точку после дорогой части фиксирую cache()'м.

Механизм, причина дизайна и два практических приёма - полное владение моделью исполнения.

На чём валят

  • Ждать ошибку на строке с filter - она вылетит на первом действии.
  • Цикл с действием внутри - каждый виток пересчитывает цепочку с нуля.
  • Кэшировать всё подряд «для скорости» - вымывание памяти и спиллы.
  • Не заглянуть в explain и не заметить, что партиционный фильтр не дочитался до скана.

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

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

  1. #spark_lazy1 / 5
    В ноутбуке ты пять раз считаешь разные агрегаты над одним и тем же тяжёлым отфильтрованным DataFrame — и каждый раз ждёшь заново. Почему и что сделать?
    A)Spark кеширует всё автоматически, так что задержка — это сеть; сделать тут ничего не получится
    B)Ноутбук держит устаревшую версию DataFrame; нужно перезапускать ядро перед каждым агрегатом
    C)Каждый агрегат — новое действие; из-за ленивости Spark заново гонит цепочку от источника. Закешировать переиспользуемый DataFrame (cache/persist)
    D)Потому что несколько агрегатов не получится считать по очереди; их обязательно собрать в один вызов agg, иного корректного пути в Spark практически нет
    показать ответ и разбор
    +C)Каждый агрегат — новое действие; из-за ленивости Spark заново гонит цепочку от источника. Закешировать переиспользуемый DataFrame (cache/persist)

    // разбор: Из-за ленивости каждое действие (каждый агрегат) заново исполняет план от чтения источника — фильтр отрабатывает пять раз. cache()/persist() материализует промежуточный DataFrame, и повторные действия берут готовое.

  2. #spark_lazy2 / 5
    Почему привычка после каждого шага звать df.count() или df.show(), чтобы «убедиться, что всё ок», может резко замедлить работу на больших данных?
    A)Каждый такой вызов — действие, которое запускает полный пересчёт цепочки и чтение данных заново
    B)Потому что count() и show() портят план запроса, и после них Spark не может его оптимизировать
    C)Потому что они пишут временные файлы на диск драйвера, и он со временем переполняется
    D)Они безопасны и ничего не замедляют — это рекомендованный способ пошаговой отладки в Spark
    показать ответ и разбор
    +A)Каждый такой вызов — действие, которое запускает полный пересчёт цепочки и чтение данных заново

    // разбор: count/show — действия. При ленивой модели каждое из них заново прогоняет весь план от источника. Расставленные «для проверки» после каждого шага, они превращают один проход по данным в несколько. Отлаживай на сэмпле.

  3. #spark_lazy3 / 5
    df.collect() (или toPandas()) на многогигабайтном DataFrame роняет приложение с OutOfMemory на драйвере, хотя кластер большой. Почему размер кластера не спасает?
    A)Потому что collect() принудительно отключает всех исполнителей кластера и намеренно пересчитывает весь датафрейм в один поток на стороне драйвера, игнорируя память узлов
    B)Потому что драйвер имеет ровно 1 ГБ памяти, и это не получится поменять настройками
    C)Потому что Parquet распаковывается в сотню раз, и формат переполнит память
    D)collect()/toPandas() стягивают все партиции со всех исполнителей в память одного процесса-драйвера; кластер обрабатывает данные, но итог не влезает в одну машину
    показать ответ и разбор
    +D)collect()/toPandas() стягивают все партиции со всех исполнителей в память одного процесса-драйвера; кластер обрабатывает данные, но итог не влезает в одну машину

    // разбор: Эти действия собирают все партиции на драйвер. Кластер помогает считать распределённо, но результат материализуется в одном процессе — если он не влезает в память драйвера, будет OOM. Поэтому сначала агрегируй/сэмплируй, потом собирай.

  4. #spark_lazy4 / 5
    Ты создал DataFrame с колонкой из случайного сэмпла, потом отдельно посчитал его размер и сохранил на диск — в сохранённом оказался не тот набор строк, что показал count(). Как ленивость это объясняет?
    A)Это известный баг конкретной сборки Spark на данном кластере; после обновления рантайма до свежей версии результаты count() и записи наверняка совпадут
    B)Без кеша каждое действие заново прогоняет цепочку со случайностью — count() и write() видят разные выборки; надёжно лечит только cache()/persist()
    C)Потому что операция записи по умолчанию обращается к другому, резервному источнику данных, а не к тому, что использовал count()
    D)Потому что count() округляет число строк, а на диск пишется точное — набор строк на самом деле один и тот же
    показать ответ и разбор
    +B)Без кеша каждое действие заново прогоняет цепочку со случайностью — count() и write() видят разные выборки; надёжно лечит только cache()/persist()

    // разбор: Ленивость плюс пересчёт: DataFrame — это план, а не зафиксированные данные. Каждое действие заново исполняет план, включая недетерминированный сэмпл, поэтому count() и write() получают разные строки. Даже одинаковый seed не гарантия (маппинг строк по партициям может смениться между запусками) — надёжно только материализовать один раз: cache()/persist() либо записать на диск и читать уже оттуда.

  5. #spark_lazy5 / 5
    Почему df.show() иногда возвращается почти мгновенно, а df.count() на том же DataFrame думает минуту?
    A)Show() читает результат из оперативного кеша, а count() каждый раз идёт на диск за всеми данными
    B)Count() медленный из-за бага в оптимизаторе запросов; show() его просто обходит стороной
    C)show() обычно берёт лишь первые несколько строк (часто одну партицию), а count() обязан пройти все данные целиком
    D)Разницы быть не должно; вероятно, между двумя вызовами кто-то изменил данные под тобой
    показать ответ и разбор
    +C)show() обычно берёт лишь первые несколько строк (часто одну партицию), а count() обязан пройти все данные целиком

    // разбор: show(n) нужно вернуть лишь n строк — Spark часто материализует одну-две партиции и останавливается. count() должен пройти весь датасет, чтобы посчитать строки. Отсюда разница; для быстрой проверки предпочитай show/limit, а не count.

дальше

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

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