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

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

  1. #spark_dataframe_ops1 / 5
    Нужны сразу несколько агрегатов по группам (сумма, среднее, число). Как их посчитать за один проход?
    A)Отдельным groupBy на каждый агрегат с последующим join всех результатов по ключу группировки
    B)groupBy(ключ).agg(sum, avg, count) — несколько агрегатов считаются за один проход
    C)Собрать данные на драйвер и посчитать в цикле Python
    D)Агрегаты в Spark считаются только по одному за вызов
    показать ответ и разбор
    +B)groupBy(ключ).agg(sum, avg, count) — несколько агрегатов считаются за один проход

    // разбор: groupBy(...).agg(...) принимает сразу несколько агрегатных выражений (sum, avg, count, min/max, approx_count_distinct) — Spark посчитает их одним проходом с общим shuffle по ключу. Дробить на несколько groupBy и джойнить результаты дороже: лишние проходы и перемешивания. Для DS это стандартный способ собрать признаки по сущности (юзеру, товару) в одну строку.

  2. #spark_dataframe_ops2 / 5
    К событиям нужно приклеить профиль пользователя, но сохранить события даже без совпавшего профиля. Какой join?
    A)Inner join — он сохранит все события, а профиль просто подставит там, где он есть в справочнике
    B)Cross join по ключу
    C)left outer join событий с профилями: события сохраняются, профиль где есть, иначе null
    D)Left anti join профилей
    показать ответ и разбор
    +C)left outer join событий с профилями: события сохраняются, профиль где есть, иначе null

    // разбор: Чтобы не потерять ни одного события, левой таблицей берут события и делают left outer join с профилями по ключу пользователя: все события остаются, а поля профиля заполняются, где нашлось совпадение, иначе null. inner отбросил бы события без профиля; anti — наоборот, оставил бы только без совпадения. Выбор типа join определяется тем, чьи строки обязаны сохраниться.

  3. #spark_dataframe_ops3 / 5
    Для фичи нужно к каждому событию добавить «предыдущую покупку этого пользователя». Какой инструмент 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). Это рабочая лошадка фиче-инжиниринга по последовательностям (сессии, история покупок).

  4. #spark_dataframe_ops4 / 5
    Нужна категориальная фича по правилу «если возраст < 18 → 'junior', иначе 'adult'». Как выразить это в Spark без UDF?
    A)when(col<18,'junior').otherwise('adult') внутри withColumn — встроенное условие
    B)Обычный питоновский if прямо в withColumn
    C)Только через пользовательскую функцию: встроенных условных выражений над колонками в 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.

  5. #spark_dataframe_ops5 / 5
    Почему пользовательскую функцию (обычный Python UDF) в Spark стараются заменить встроенными функциями, когда это возможно?
    A)UDF в Spark вообще запрещены
    B)Python UDF гонит строки между JVM и Python и непрозрачен Catalyst; спасает pandas UDF
    C)UDF занимают слишком много места на диске
    D)UDF выполняются исключительно на драйвере, поэтому не задействуют исполнители и оттого медленны
    показать ответ и разбор
    +B)Python UDF гонит строки между JVM и Python и непрозрачен Catalyst; спасает pandas UDF

    // разбор: Обычный Python UDF для каждой строки сериализует данные из JVM в Python-процесс исполнителя и обратно, а для Catalyst остаётся непрозрачным — оптимизации (pushdown, перестановка) через него не проходят. Отсюда ощутимое замедление. Встроенные функции выполняются в JVM и оптимизируются. Если Python-логика необходима, берут pandas UDF (векторизованный, через Apache Arrow, обрабатывает батчи), который заметно быстрее построчного.

дальше

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

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