DataFrame API для фич
DataFrame-операции - там, где Spark тормозит у новичков: джойн без broadcast, окно без partitionBy, Python-UDF на каждый чих. Вопросы проверяют знание трёх главных тюнингов.
// Симптом-легенда темы: стейдж висит на 199/200 задач. Это не зависание это skew.
Джойны: broadcast - главный тюнинг
Большой × большой - sort-merge join с шаффлом обеих сторон: дорого, но честно. Большой × маленький - broadcast(df) малой стороны: справочник рассылается экзекьюторам, шаффл исчезает.
Джойн со справочником без broadcast - полный шаффл обеих сторон на ровном месте; это самый дешёвый крупный выигрыш в Spark-запросах.
// DataFrame API и spark.sql - один план: пиши хоть SQL'ем, оптимизатору всё равно.
- broadcast join
- малая таблица рассылается всем - шаффл исчезает
Skew: когда один ключ тащит стейдж
Перекос ключей: один гигантский ключ (null, «гость», мегаполис) собирает партицию-монстра, и последняя задача стейджа работает за всех - те самые 199/200.
Лечения: AQE (adaptive query execution) skew-join (адаптивное разбиение горячих партиций), salting - соль в ключ для размазывания, отделение горячих ключей отдельной веткой.
// Диагноз ставится по распределению ключей и Spark UI: размер задач стейджа расскажет всё.
- skew / salting
- перекос ключа / соль для размазывания горячего ключа
- AQE
- adaptive query execution: правит план по фактическим данным
Окна и UDF
Оконные функции: Window.partitionBy().orderBy() + row_number/lag - дедупликация «последняя версия записи», кумулятивы, лаги. Окно без partitionBy стягивает весь датасет в одну задачу - стейдж встаёт.
Python-UDF (user-defined function) медленные (сериализация на каждую строку) и непрозрачные для Catalyst. Порядок выбора: встроенные функции → pandas_udf → обычный UDF последним.
// NULL-семантика SQL жива и тут: сравнения с NULL дают NULL, джойн по NULL-ключам не матчится; для граничных случаев - eqNullSafe.
- pandas_udf
- векторизованный UDF через Arrow - быстрее построчного
Как отвечать: «Джойн фактов со справочником тормозит. Что сделаешь?»
Первое - broadcast справочника: малая сторона рассылается экзекьюторам, и шаффл исчезает это главный тюнинг джойнов, проверяю по explain, что план сменился на BroadcastHashJoin. Если тормозит по-прежнему и стейдж висит на последних задачах это skew: смотрю распределение ключа, null'ы и «гостевые» id собирают партиции-монстры; лечу AQE skew-join'ом или salting'ом, горячие ключи можно отделить отдельной веткой. И попутно проверяю, что нет Python-UDF в горячем месте - встроенные функции быстрее на порядок.
Два главных лечения в правильном порядке с диагностикой по explain и Spark UI - опыт, не конспект.
На чём валят
- −Python-UDF там, где есть встроенная функция, на порядок медленнее.
- −Джойн со справочником без broadcast - полный шаффл обеих сторон.
- −Окно без partitionBy - весь датасет в одной задаче.
- −dropDuplicates без orderBy-логики - «какая-то» строка вместо последней; нужен row_number.
- −Стейдж «висит на 199/200» это skew, смотри распределение ключей.
Проверьте себя
Пять вопросов из банка по этой подтеме. Всего их 14, остальные разбираются в тренажёре.
- Нужны сразу несколько агрегатов по группам (сумма, среднее, число). Как их посчитать за один проход?A)Отдельным groupBy на каждый агрегат с последующим join всех результатов по ключу группировкиB)groupBy(ключ).agg(sum, avg, count) — несколько агрегатов считаются за один проходC)Собрать данные на драйвер и посчитать в цикле PythonD)Агрегаты в Spark считаются только по одному за вызов
показать ответ и разбор
+B)groupBy(ключ).agg(sum, avg, count) — несколько агрегатов считаются за один проход// разбор: groupBy(...).agg(...) принимает сразу несколько агрегатных выражений (sum, avg, count, min/max, approx_count_distinct) — Spark посчитает их одним проходом с общим shuffle по ключу. Дробить на несколько groupBy и джойнить результаты дороже: лишние проходы и перемешивания. Для DS это стандартный способ собрать признаки по сущности (юзеру, товару) в одну строку.
- К событиям нужно приклеить профиль пользователя, но сохранить события даже без совпавшего профиля. Какой join?A)Inner join — он сохранит все события, а профиль просто подставит там, где он есть в справочникеB)Cross join по ключуC)left outer join событий с профилями: события сохраняются, профиль где есть, иначе nullD)Left anti join профилей
показать ответ и разбор
+C)left outer join событий с профилями: события сохраняются, профиль где есть, иначе null// разбор: Чтобы не потерять ни одного события, левой таблицей берут события и делают left outer join с профилями по ключу пользователя: все события остаются, а поля профиля заполняются, где нашлось совпадение, иначе null. inner отбросил бы события без профиля; anti — наоборот, оставил бы только без совпадения. Выбор типа join определяется тем, чьи строки обязаны сохраниться.
- Для фичи нужно к каждому событию добавить «предыдущую покупку этого пользователя». Какой инструмент Spark?A)Обычный groupBy — он и даёт доступ к соседним строкамB)Оконная функция Window.partitionBy(user).orderBy(time) + lag даёт предыдущую строкуC)Collect() и обход в Python-цикле по времениD)Джойном таблицы самой с собой по паре (user, time+1), чтобы подтянуть соседнюю по времени строку
показать ответ и разбор
+B)Оконная функция Window.partitionBy(user).orderBy(time) + lag даёт предыдущую строку// разбор: Оконные функции считают значения относительно упорядоченного окна, НЕ схлопывая строки (в отличие от groupBy). Определяют окно Window.partitionBy(user).orderBy(time), затем lag(col, 1) достаёт предыдущее значение, lead — следующее, rank/row_number — позицию, а rolling-суммы — по рамке (rowsBetween). Это рабочая лошадка фиче-инжиниринга по последовательностям (сессии, история покупок).
- Нужна категориальная фича по правилу «если возраст < 18 → 'junior', иначе 'adult'». Как выразить это в Spark без UDF?A)when(col<18,'junior').otherwise('adult') внутри withColumn — встроенное условиеB)Обычный питоновский if прямо в withColumnC)Только через пользовательскую функцию: встроенных условных выражений над колонками в Spark нетD)Через фильтр df.filter, он и создаст колонку
показать ответ и разбор
+A)when(col<18,'junior').otherwise('adult') внутри withColumn — встроенное условие// разбор: Условные фичи строят встроенным выражением when(...).otherwise(...) из pyspark.sql.functions — оно компилируется в план и выполняется движком на исполнителях, без выхода в Python. Питоновский if в withColumn выполнится один раз на драйвере при построении плана, а не по строкам. UDF решил бы задачу, но медленнее (сериализация Python↔JVM), поэтому для простых условий предпочитают when/otherwise.
- Почему пользовательскую функцию (обычный Python UDF) в Spark стараются заменить встроенными функциями, когда это возможно?A)UDF в Spark вообще запрещеныB)Python UDF гонит строки между JVM и Python и непрозрачен Catalyst; спасает pandas UDFC)UDF занимают слишком много места на дискеD)UDF выполняются исключительно на драйвере, поэтому не задействуют исполнители и оттого медленны
показать ответ и разбор
+B)Python UDF гонит строки между JVM и Python и непрозрачен Catalyst; спасает pandas UDF// разбор: Обычный Python UDF для каждой строки сериализует данные из JVM в Python-процесс исполнителя и обратно, а для Catalyst остаётся непрозрачным — оптимизации (pushdown, перестановка) через него не проходят. Отсюда ощутимое замедление. Встроенные функции выполняются в JVM и оптимизируются. Если Python-логика необходима, берут pandas UDF (векторизованный, через Apache Arrow, обрабатывает батчи), который заметно быстрее построчного.
дальше
Теорию прочитали. Навык ставится повторением
В Сеньорчике эта подтема идёт в ежедневных сессиях: движок возвращает её, пока ответы не станут уверенными, и ведёт прогресс отдельно по каждой подтеме. Теория внутри тоже бесплатна, лимит только на количество вопросов в день.