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

Spark: от данных к модели

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

Фичи в Spark - те же утечки, что в pandas, только дороже: неделя кластерного времени на фичи, подсмотревшие таргет. Проверяют перенос дисциплины point-in-time на распределённый масштаб.

// Плюс фирменная ловушка Spark: randomSplit, который пересекается сам с собой.

Point-in-time на масштабе

Фичи считаются распределённо - groupBy-агрегаты, оконные лаги, джойны справочников, но правило то же: агрегат строго «до момента события». Оконные функции с ограничением по времени решают point-in-time.

Оконный лаг без сортировки по времени внутри партиции юзера - мусор вместо лага: порядок строк в партиции не гарантирован.

// Финальная матрица обычно компактна и уезжает в sklearn/бустинг - Spark готовит, локальный стек обучает.

point-in-time
фичи из данных до момента события - против утечки

Сплит раньше статистик

Порядок жёсткий: сначала разрез на train/test (по времени или юзерам), потом расчёт статистик-фичей только по трейну. Агрегат по всем данным до сплита - target encoding подсмотрел тест.

randomSplit в Spark недетерминирован при пересчёте графа: ленивое исполнение может дать разные сплиты в разных действиях - train и test пересекутся между запусками. Фиксация - материализация сплита записью в таблицу.

// Категории высокой кардинальности - target/count encoding агрегатами по трейн-периоду, а не one-hot на миллион колонок.

материализация сплита
сплит записан в таблицу, а не «сид в воздухе»

Фичи, доживающие до сервинга

Если фича нужна и на сервинге, у неё должен быть тот же код расчёта, а тянуть Spark в онлайн нельзя. Отсюда баланс: тяжёлые исторические агрегаты - Spark'ом в фичстор, лёгкие контекстные - кодом сервинга.

Spark ML Pipeline (Transformer/Estimator) уместен, когда и обучение распределённое: индексаторы, энкодеры и модель сериализуются одним пайплайном.

// Фича, воспроизводимая только Spark-джобой, для онлайн-модели не существует - проверяй это на этапе дизайна, а не деплоя.

Как отвечать: «Как посчитать фичи в Spark и не утечь?»

Три дисциплины. Первая - point-in-time: все агрегаты оконными функциями с ограничением «до момента события», лаги - с явной сортировкой по времени внутри партиции юзера. Вторая - порядок: сначала сплит по времени или юзерам, материализованный записью в таблицу - randomSplit при ленивом пересчёте недетерминирован и умеет пересекаться сам с собой; статистики и энкодинги считаются после, только по трейну. Третья - сервинг: фичи, нужные онлайн, обязаны иметь путь расчёта без Spark, иначе они существуют только в бэктесте.

Все три специфичные ловушки Spark-фичей закрыты, включая фирменный randomSplit - глубина, которую даёт только практика.

На чём валят

  • Агрегат по всем данным до сплита - target encoding подсмотрел тест.
  • Оконный лаг без сортировки по времени в партиции юзера.
  • randomSplit + ленивый пересчёт - train и test пересеклись между запусками.
  • One-hot на кардинальности 10⁶ - нужен count/target encoding.
  • Фичи онлайн-модели, воспроизводимые только Spark-джобой.

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

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

  1. #spark_ml_features1 / 5
    Категориальную колонку «город» нужно подать в модель Spark ML. Какая типовая цепочка преобразований?
    A)Подать строковую колонку как есть — оценщики Spark ML сами распознают и закодируют категории
    B)Сразу OneHotEncoder по строковой колонке
    C)StringIndexer (строка→индекс), OneHotEncoder (индекс→бинарный вектор), затем VectorAssembler
    D)Захешировать город в одно число и подать его как непрерывный признак
    показать ответ и разбор
    +C)StringIndexer (строка→индекс), OneHotEncoder (индекс→бинарный вектор), затем VectorAssembler

    // разбор: Spark ML не принимает строковые категории напрямую. Стандартный конвейер: StringIndexer кодирует строки числовыми индексами по частоте; OneHotEncoder превращает индекс в разреженный бинарный вектор (чтобы модель не приняла индекс за порядковую величину); затем VectorAssembler объединяет всё в итоговый вектор. Подавать сырой индекс как непрерывный признак — ошибка: 3 не «больше» 1 по смыслу города.

  2. #spark_ml_features2 / 5
    Модель обучаешь на sklearn, но датасет 300 ГБ. Разумная стратегия для DS?
    A)Обучать sklearn прямо на всех 300 ГБ, подав их в fit на драйвере одним большим массивом в память
    B)Sklearn с данными такого размера не совместим, вариантов нет
    C)Сэмпл/агрегат в Spark под sklearn, либо распределённый MLlib на всём объёме данных
    D)Считать признаки в pandas на драйвере, стянув всё туда
    показать ответ и разбор
    +C)Сэмпл/агрегат в Spark под sklearn, либо распределённый MLlib на всём объёме данных

    // разбор: sklearn работает в памяти одного процесса, поэтому 300 ГБ ему целиком не подать. Два рабочих пути: (1) в Spark подготовить/сжать данные — отфильтровать, агрегировать признаки, взять репрезентативный сэмпл — и обучить sklearn на помещающемся куске; (2) если нужен весь объём, использовать распределённые алгоритмы Spark MLlib. Стягивать 300 ГБ на драйвер (collect/toPandas) нельзя — это OOM.

  3. #spark_ml_features3 / 5
    Обученную sklearn-модель нужно применить (inference) к миллиарду строк в Spark. Как сделать это эффективно?
    A)Собрать миллиард строк на драйвер через toPandas и применить model.predict одним общим вызовом
    B)Обучить модель заново средствами Spark на каждой партиции
    C)Broadcast модели на исполнители + pandas UDF/mapInPandas — инференс распределённо, батчами
    D)Применять модель построчно обычным Python UDF — это оптимально
    показать ответ и разбор
    +C)Broadcast модели на исполнители + pandas UDF/mapInPandas — инференс распределённо, батчами

    // разбор: Инференс распараллеливают: небольшую sklearn-модель рассылают на исполнители через broadcast, а предсказания делают векторизованно — pandas UDF или mapInPandas получают данные батчами (через Arrow) и зовут model.predict на pandas-кадре внутри исполнителя. Так миллиард строк обрабатывается распределённо, без стягивания на драйвер. Построчный обычный UDF работал бы, но медленно; toPandas всего набора — OOM.

  4. #spark_ml_features4 / 5
    Что даёт объединение шагов подготовки и модели в Spark ML Pipeline?
    A)Никакой пользы, кроме более длинного кода: те же шаги проще расписать по отдельности вручную
    B)Pipeline нужен только для визуализации шагов
    C)Pipeline ускоряет обучение, распараллеливая разные модели
    D)Pipeline связывает трансформеры и модель: те же шаги на train, valid и в проде
    показать ответ и разбор
    +D)Pipeline связывает трансформеры и модель: те же шаги на train, valid и в проде

    // разбор: Pipeline (stages: индексаторы, энкодеры, VectorAssembler, оценщик) оформляет весь препроцессинг и модель как единый объект. fit прогоняет и настраивает все стадии, transform применяет их в том же порядке. Главная выгода — воспроизводимость: одинаковые преобразования на обучении, валидации и в бою, без риска «на train нормализовали, а в проде забыли». Плюс удобно сохранять/загружать целиком.

  5. #spark_ml_features5 / 5
    При масштабировании признаков (StandardScaler) в Spark ML на что критично не наступить, чтобы не получить утечку данных (leakage)?
    A)Масштабировать нужно после обучения модели, а не до
    B)Достаточно масштабировать весь датасет один раз до разбиения на train/test — так удобнее и быстрее
    C)StandardScaler в Spark не даёт утечки по определению
    D)Масштаб (среднее/стд) учат только на train; к valid/test применяют уже обученный scaler
    показать ответ и разбор
    +D)Масштаб (среднее/стд) учат только на train; к valid/test применяют уже обученный scaler

    // разбор: Утечка возникает, когда параметры препроцессинга «видят» тестовые данные. StandardScaler запоминает среднее и стандартное отклонение на fit — их надо считать только на train, а на valid/test лишь применять transform обученного scaler. Если посчитать статистики по всему датасету до сплита, информация о распределении теста просочится в обучение, и оценка качества будет оптимистично завышена. В Pipeline это соблюдается автоматически, если fit идёт по train.

дальше

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

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