Выборка данных под ML
Выгрузка данных под ML - ежедневный маршрут DS в Spark, и ошибка №1 всей темы живёт здесь: toPandas() на десятках миллионов строк. Проверяют дисциплину «агрегируй там, вытягивай мало».
// Правило: фильтр и агрегация живут в Spark; наружу уезжает уже маленькое.
toPandas: посчитай, потом вызывай
toPandas() собирает всё на драйвер. Перед вызовом - прикидка размера: строки × колонки × тип; десятки миллионов строк - OOM (out of memory), классика номер один.
«Посмотреть глазами» это limit(1000).toPandas() или show(n); describe и approxQuantile дают сводки распределённо, без вытягивания.
// Обмен с pandas кратно ускоряет Arrow (spark.sql.execution.arrow.pyspark.enabled): колоночная сериализация вместо построчной.
- Arrow
- колоночный обмен Spark↔pandas - кратно быстрее
Читать меньше: предикаты и колонки
Партиционные фильтры (dt='2026-07-01') отсекают файлы до чтения это не оптимизация запроса, а нечтение лишнего. select нужных колонок - паркет колоночный, читаются только они.
select * из широкого паркета ради трёх колонок - плата чтением всех: на колоночных форматах это кратная разница.
// Запись результата - паркет с партиционированием по ключу будущего чтения; один гигантский CSV через coalesce(1) - антипаттерн.
- partition pruning
- фильтр по партиции отсекает файлы до чтения
Сэмпл для разработки
Итерации кода - на sample(fraction): быстрые циклы на доле данных, полный прогон - когда логика готова.
Если важна целостность историй (сессии, воронки) - сэмпл по юзерам, не по строкам: случайные строки рвут истории, и сессионная логика тихо ломается.
// Сэмпл по юзерам это фильтр по хешу id (например, pmod(hash(user_id), 100) < 5), стабильный между запусками.
Как отвечать: «Как безопасно вытащить датасет из Spark в pandas?»
Сначала уменьшаю на стороне Spark: партиционные фильтры при чтении, select только нужных колонок, агрегация до зерна, которое реально нужно модели. Потом прикидываю размер результата - строки на колонки на типы; если это не мегабайты-сотни мегабайт, в pandas ему рано. Включаю Arrow для быстрой сериализации и только тогда toPandas(). Посмотреть глазами - limit().toPandas(), сводки - describe/approxQuantile распределённо. Для разработки - сэмпл по юзерам, чтобы не порвать истории, полный прогон в конце.
Порядок «уменьши → посчитай → вытяни» плюс сэмпл по юзерам - дисциплина, которая и проверяется.
На чём валят
- −toPandas() на десятках миллионов строк - OOM драйвера, классика №1.
- −Читать всю таблицу и фильтровать в pandas - фильтр живёт в Spark.
- −select * из широкого паркета ради трёх колонок.
- −Сэмпл по строкам для сессионного анализа - истории юзеров порваны.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 14, остальные разбираются в тренажёре.
- Приклеиваешь к таблице на 500 млн строк маленький справочник категорий (2 тыс. строк). Как избежать дорогого перемешивания данных по всему кластеру?A)Заранее отсортировать обе таблицы по ключу перед джойном — тогда shuffle не понадобитсяB)Переписать джойн как цикл по строкам справочника — так выйдет быстрее, чем один большой joinC)Резко увеличить число партиций у большой таблицы, чтобы перемешивание распараллелилось между узлами и в итоге стало практически бесплатнымD)Сделать broadcast-join: разослать маленький справочник на всех исполнителей, тогда большую таблицу перемешивать не нужно
показать ответ и разбор
+D)Сделать broadcast-join: разослать маленький справочник на всех исполнителей, тогда большую таблицу перемешивать не нужно// разбор: Обычный join перемешивает (shuffle) обе таблицы по ключу — для 500 млн строк это дорого. Если одна таблица маленькая, её рассылают на все исполнители (broadcast), и большая джойнится локально без shuffle. Классический приём при склейке фич со справочником.
- Джойн событий с профилями завис: 199 задач готовы за секунды, одна крутится час. Что с данными и как чинить со стороны DS?A)Кластеру банально не хватает исполнителей под нагрузку; если добавить достаточно новых узлов, последняя зависшая задача ускорится ровно пропорционально их числуB)Перекос данных (skew): на один ключ приходится непропорционально много строк (часто NULL/гость по user_id); отфильтровать или «посолить» проблемный ключC)Это нормально: последняя задача идёт дольше остальных из-за завершающих операций SparkD)Повреждён Parquet-файл одной из партиций; достаточно перечитать таблицу, и задача выровняется
показать ответ и разбор
+B)Перекос данных (skew): на один ключ приходится непропорционально много строк (часто NULL/гость по user_id); отфильтровать или «посолить» проблемный ключ// разбор: Один воркер обрабатывает свою партицию по ключу джойна. Если на один ключ (например, NULL или гостевой user_id) приходится гигантская доля строк, его партиция огромна — задача по ней тянется, пока остальные простаивают. Лечится фильтрацией мусорного ключа или солью (salting).
- Хочешь быстро прикинуть распределение фичи на локальной машине, но таблица огромная. Разумный первый шаг?A)Забрать всю таблицу в pandas через toPandas() и построить гистограмму уже тамB)Отсортировать всю таблицу и взять первые 10 тыс. строк как выборку для гистограммыC)Ничего не выйдет — распределение можно смотреть на полных данных, иначе форма будет невернойD)Взять репрезентативный сэмпл в Spark (df.sample) и уже его тащить в pandas для гистограммы
показать ответ и разбор
+D)Взять репрезентативный сэмпл в Spark (df.sample) и уже его тащить в pandas для гистограммы// разбор: Для прикидки распределения полный объём не нужен — репрезентативный случайный сэмпл покажет форму. df.sample делает выборку распределённо, а в pandas уходит уже немного строк. «Первые 10 тыс.» после сортировки — смещённая выборка, не случайная.
- Из широкой таблицы на 300 колонок тебе нужны 5. Почему в колоночном формате (Parquet) важно выбрать эти 5 явно, а не читать через select *?A)Select * в Parquet вызывает ошибку — формат требует перечислять колонки вручнуюB)Разницы нет: Parquet всё равно физически читает все 300 колонок при запросеC)Колоночное хранение позволяет прочитать с диска только нужные колонки (column pruning) — меньше I/O; select * тащит все 300D)Select * почти быстрее, потому что Spark читает файл одним сплошным последовательным проходом без какой-либо дополнительной логики отбора колонок
показать ответ и разбор
+C)Колоночное хранение позволяет прочитать с диска только нужные колонки (column pruning) — меньше I/O; select * тащит все 300// разбор: В колоночном формате данные лежат по колонкам, поэтому Spark физически читает только запрошенные (column pruning) — для 5 из 300 это кратно меньше I/O. select * заставляет прочитать все колонки, даже ненужные. Всегда выбирай явный список колонок.
- Тебе для фич нужны только события по стране RU, дальше идёт несколько джойнов. Почему выгоднее поставить фильтр сразу после чтения, а не в самом конце?A)Порядок операций в Spark не влияет ровным счётом ни на что: встроенный оптимизатор в случае сам переставит все шаги в идеальную последовательность за тебяB)Фильтр в конце точнее: после всех джойнов видно, какие строки нужны моделиC)Фильтр в начале опасен: можно случайно выкинуть строки, которые понадобятся в джойнах нижеD)Ранний фильтр уменьшает число строк ещё до джойнов (и включает predicate pushdown к чтению) — через дорогой shuffle пройдёт меньше данных
показать ответ и разбор
+D)Ранний фильтр уменьшает число строк ещё до джойнов (и включает predicate pushdown к чтению) — через дорогой shuffle пройдёт меньше данных// разбор: Чем раньше отсечь лишние строки, тем меньше данных пойдёт в дорогие джойны и shuffle. Catalyst часто и сам проталкивает фильтр вниз (predicate pushdown), но писать фильтр рано — надёжно и читаемо. Фильтр по RU не выкидывает нужное: строки других стран тебе и не нужны.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.