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, остальные разбираются в тренажёре.
- В ноутбуке ты пять раз считаешь разные агрегаты над одним и тем же тяжёлым отфильтрованным DataFrame — и каждый раз ждёшь заново. Почему и что сделать?A)Spark кеширует всё автоматически, так что задержка — это сеть; сделать тут ничего не получитсяB)Ноутбук держит устаревшую версию DataFrame; нужно перезапускать ядро перед каждым агрегатомC)Каждый агрегат — новое действие; из-за ленивости Spark заново гонит цепочку от источника. Закешировать переиспользуемый DataFrame (cache/persist)D)Потому что несколько агрегатов не получится считать по очереди; их обязательно собрать в один вызов agg, иного корректного пути в Spark практически нет
показать ответ и разбор
+C)Каждый агрегат — новое действие; из-за ленивости Spark заново гонит цепочку от источника. Закешировать переиспользуемый DataFrame (cache/persist)// разбор: Из-за ленивости каждое действие (каждый агрегат) заново исполняет план от чтения источника — фильтр отрабатывает пять раз. cache()/persist() материализует промежуточный DataFrame, и повторные действия берут готовое.
- Почему привычка после каждого шага звать df.count() или df.show(), чтобы «убедиться, что всё ок», может резко замедлить работу на больших данных?A)Каждый такой вызов — действие, которое запускает полный пересчёт цепочки и чтение данных зановоB)Потому что count() и show() портят план запроса, и после них Spark не может его оптимизироватьC)Потому что они пишут временные файлы на диск драйвера, и он со временем переполняетсяD)Они безопасны и ничего не замедляют — это рекомендованный способ пошаговой отладки в Spark
показать ответ и разбор
+A)Каждый такой вызов — действие, которое запускает полный пересчёт цепочки и чтение данных заново// разбор: count/show — действия. При ленивой модели каждое из них заново прогоняет весь план от источника. Расставленные «для проверки» после каждого шага, они превращают один проход по данным в несколько. Отлаживай на сэмпле.
- df.collect() (или toPandas()) на многогигабайтном DataFrame роняет приложение с OutOfMemory на драйвере, хотя кластер большой. Почему размер кластера не спасает?A)Потому что collect() принудительно отключает всех исполнителей кластера и намеренно пересчитывает весь датафрейм в один поток на стороне драйвера, игнорируя память узловB)Потому что драйвер имеет ровно 1 ГБ памяти, и это не получится поменять настройкамиC)Потому что Parquet распаковывается в сотню раз, и формат переполнит памятьD)collect()/toPandas() стягивают все партиции со всех исполнителей в память одного процесса-драйвера; кластер обрабатывает данные, но итог не влезает в одну машину
показать ответ и разбор
+D)collect()/toPandas() стягивают все партиции со всех исполнителей в память одного процесса-драйвера; кластер обрабатывает данные, но итог не влезает в одну машину// разбор: Эти действия собирают все партиции на драйвер. Кластер помогает считать распределённо, но результат материализуется в одном процессе — если он не влезает в память драйвера, будет OOM. Поэтому сначала агрегируй/сэмплируй, потом собирай.
- Ты создал 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() либо записать на диск и читать уже оттуда.
- Почему 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.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.