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, остальные разбираются в тренажёре.
- Категориальную колонку «город» нужно подать в модель Spark ML. Какая типовая цепочка преобразований?A)Подать строковую колонку как есть — оценщики Spark ML сами распознают и закодируют категорииB)Сразу OneHotEncoder по строковой колонкеC)StringIndexer (строка→индекс), OneHotEncoder (индекс→бинарный вектор), затем VectorAssemblerD)Захешировать город в одно число и подать его как непрерывный признак
показать ответ и разбор
+C)StringIndexer (строка→индекс), OneHotEncoder (индекс→бинарный вектор), затем VectorAssembler// разбор: Spark ML не принимает строковые категории напрямую. Стандартный конвейер: StringIndexer кодирует строки числовыми индексами по частоте; OneHotEncoder превращает индекс в разреженный бинарный вектор (чтобы модель не приняла индекс за порядковую величину); затем VectorAssembler объединяет всё в итоговый вектор. Подавать сырой индекс как непрерывный признак — ошибка: 3 не «больше» 1 по смыслу города.
- Модель обучаешь на 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.
- Обученную 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.
- Что даёт объединение шагов подготовки и модели в Spark ML Pipeline?A)Никакой пользы, кроме более длинного кода: те же шаги проще расписать по отдельности вручнуюB)Pipeline нужен только для визуализации шаговC)Pipeline ускоряет обучение, распараллеливая разные моделиD)Pipeline связывает трансформеры и модель: те же шаги на train, valid и в проде
показать ответ и разбор
+D)Pipeline связывает трансформеры и модель: те же шаги на train, valid и в проде// разбор: Pipeline (stages: индексаторы, энкодеры, VectorAssembler, оценщик) оформляет весь препроцессинг и модель как единый объект. fit прогоняет и настраивает все стадии, transform применяет их в том же порядке. Главная выгода — воспроизводимость: одинаковые преобразования на обучении, валидации и в бою, без риска «на train нормализовали, а в проде забыли». Плюс удобно сохранять/загружать целиком.
- При масштабировании признаков (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.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.