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