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

Выборка данных под 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, остальные разбираются в тренажёре.

  1. #spark_data_pull1 / 5
    Приклеиваешь к таблице на 500 млн строк маленький справочник категорий (2 тыс. строк). Как избежать дорогого перемешивания данных по всему кластеру?
    A)Заранее отсортировать обе таблицы по ключу перед джойном — тогда shuffle не понадобится
    B)Переписать джойн как цикл по строкам справочника — так выйдет быстрее, чем один большой join
    C)Резко увеличить число партиций у большой таблицы, чтобы перемешивание распараллелилось между узлами и в итоге стало практически бесплатным
    D)Сделать broadcast-join: разослать маленький справочник на всех исполнителей, тогда большую таблицу перемешивать не нужно
    показать ответ и разбор
    +D)Сделать broadcast-join: разослать маленький справочник на всех исполнителей, тогда большую таблицу перемешивать не нужно

    // разбор: Обычный join перемешивает (shuffle) обе таблицы по ключу — для 500 млн строк это дорого. Если одна таблица маленькая, её рассылают на все исполнители (broadcast), и большая джойнится локально без shuffle. Классический приём при склейке фич со справочником.

  2. #spark_data_pull2 / 5
    Джойн событий с профилями завис: 199 задач готовы за секунды, одна крутится час. Что с данными и как чинить со стороны DS?
    A)Кластеру банально не хватает исполнителей под нагрузку; если добавить достаточно новых узлов, последняя зависшая задача ускорится ровно пропорционально их числу
    B)Перекос данных (skew): на один ключ приходится непропорционально много строк (часто NULL/гость по user_id); отфильтровать или «посолить» проблемный ключ
    C)Это нормально: последняя задача идёт дольше остальных из-за завершающих операций Spark
    D)Повреждён Parquet-файл одной из партиций; достаточно перечитать таблицу, и задача выровняется
    показать ответ и разбор
    +B)Перекос данных (skew): на один ключ приходится непропорционально много строк (часто NULL/гость по user_id); отфильтровать или «посолить» проблемный ключ

    // разбор: Один воркер обрабатывает свою партицию по ключу джойна. Если на один ключ (например, NULL или гостевой user_id) приходится гигантская доля строк, его партиция огромна — задача по ней тянется, пока остальные простаивают. Лечится фильтрацией мусорного ключа или солью (salting).

  3. #spark_data_pull3 / 5
    Хочешь быстро прикинуть распределение фичи на локальной машине, но таблица огромная. Разумный первый шаг?
    A)Забрать всю таблицу в pandas через toPandas() и построить гистограмму уже там
    B)Отсортировать всю таблицу и взять первые 10 тыс. строк как выборку для гистограммы
    C)Ничего не выйдет — распределение можно смотреть на полных данных, иначе форма будет неверной
    D)Взять репрезентативный сэмпл в Spark (df.sample) и уже его тащить в pandas для гистограммы
    показать ответ и разбор
    +D)Взять репрезентативный сэмпл в Spark (df.sample) и уже его тащить в pandas для гистограммы

    // разбор: Для прикидки распределения полный объём не нужен — репрезентативный случайный сэмпл покажет форму. df.sample делает выборку распределённо, а в pandas уходит уже немного строк. «Первые 10 тыс.» после сортировки — смещённая выборка, не случайная.

  4. #spark_data_pull4 / 5
    Из широкой таблицы на 300 колонок тебе нужны 5. Почему в колоночном формате (Parquet) важно выбрать эти 5 явно, а не читать через select *?
    A)Select * в Parquet вызывает ошибку — формат требует перечислять колонки вручную
    B)Разницы нет: Parquet всё равно физически читает все 300 колонок при запросе
    C)Колоночное хранение позволяет прочитать с диска только нужные колонки (column pruning) — меньше I/O; select * тащит все 300
    D)Select * почти быстрее, потому что Spark читает файл одним сплошным последовательным проходом без какой-либо дополнительной логики отбора колонок
    показать ответ и разбор
    +C)Колоночное хранение позволяет прочитать с диска только нужные колонки (column pruning) — меньше I/O; select * тащит все 300

    // разбор: В колоночном формате данные лежат по колонкам, поэтому Spark физически читает только запрошенные (column pruning) — для 5 из 300 это кратно меньше I/O. select * заставляет прочитать все колонки, даже ненужные. Всегда выбирай явный список колонок.

  5. #spark_data_pull5 / 5
    Тебе для фич нужны только события по стране RU, дальше идёт несколько джойнов. Почему выгоднее поставить фильтр сразу после чтения, а не в самом конце?
    A)Порядок операций в Spark не влияет ровным счётом ни на что: встроенный оптимизатор в случае сам переставит все шаги в идеальную последовательность за тебя
    B)Фильтр в конце точнее: после всех джойнов видно, какие строки нужны модели
    C)Фильтр в начале опасен: можно случайно выкинуть строки, которые понадобятся в джойнах ниже
    D)Ранний фильтр уменьшает число строк ещё до джойнов (и включает predicate pushdown к чтению) — через дорогой shuffle пройдёт меньше данных
    показать ответ и разбор
    +D)Ранний фильтр уменьшает число строк ещё до джойнов (и включает predicate pushdown к чтению) — через дорогой shuffle пройдёт меньше данных

    // разбор: Чем раньше отсечь лишние строки, тем меньше данных пойдёт в дорогие джойны и shuffle. Catalyst часто и сам проталкивает фильтр вниз (predicate pushdown), но писать фильтр рано — надёжно и читаемо. Фильтр по RU не выкидывает нужное: строки других стран тебе и не нужны.

дальше

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

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