Как посмотреть план выполнения запроса (execution plan) в Apache Spark?

«Как посмотреть план выполнения запроса (execution plan) в Apache Spark?» — вопрос из категории Apache Spark, который задают на 26% собеседований Data Scientist / ML Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

В Apache Spark план выполнения запроса можно и нужно анализировать для понимания логики работы и оптимизации. Основной метод — .explain().

1. Основные режимы .explain():

  • .explain() или .explain(mode="simple"): Выводит только физический план (Physical Plan).
  • .explain(mode="extended") или .explain(True): Выводит логический (Logical Plan), оптимизированный логический план (Optimized Logical Plan) и физический план (Physical Plan).
  • .explain(mode="formatted"): Выводит план в виде отдельной таблицы с разделами, что улучшает читаемость.
  • .explain(mode="cost"): Доступен, если включена поддержка CBO (Cost-Based Optimizer). Показывает план с оценкой стоимости операций.

Пример:

val df = spark.read.parquet("data.parquet")
val resultDF = df.filter($"age" > 25).groupBy("department").agg(avg("salary"))

// Самый полезный вариант для анализа
resultDF.explain("extended")

2. Ключевые элементы плана и на что обращать внимание:

  • Scan / FileScan: Как и откуда читаются данные. Проверяйте, используются ли партиционирование и фильтрация на уровне чтения (PushedFilters).
  • Filter: Где применяется условие WHERE. Желательно, чтобы фильтрация происходила как можно раньше.
  • Project: Выбор столбцов.
  • HashAggregate или ObjectHashAggregate: Операции агрегации (GROUP BY, агрегатные функции).
  • Exchange: Ключевой момент! Это операция shuffle (перемешивание данных между узлами кластера). Shuffle — самая дорогая операция. Её наличие указывает на необходимость перераспределения данных (например, перед groupBy по ключу, отличному от партиционирования).
    • Exchange hashpartitioning(...) — shuffle для агрегаций или join.
    • Exchange singlepartitioning — сбор всех данных в один узел (опасно для больших данных!).
  • Sort: Сортировка. Может быть очень дорогой.
  • BroadcastHashJoin или SortMergeJoin: Тип соединения (join). BroadcastHashJoin (broadcast) обычно эффективнее, если одна из таблиц мала.

3. Практические шаги по оптимизации на основе плана:

  1. Ищите Exchange (shuffle): Можно ли его избежать? Например, перепартиционировав данные заранее или используя broadcast join для маленькой таблицы.
  2. Проверяйте PushedFilters в Scan: Убедитесь, что предикаты фильтрации "проталкиваются" до чтения данных (особенно для форматов Parquet/ORC).
  3. Обращайте внимание на размеры данных на каждом этапе: План может показывать примерное число строк. Резкий рост числа строк после Join может указывать на проблему (декартово произведение).
  4. Сравнивайте планы: Запустите .explain() для разных вариантов записи одного запроса (например, с перестановкой filter и join), чтобы увидеть, как меняется план.