EXPLAIN FORMATTED: чтение логического и физического плана пошагово
Как читать execution plan: Catalyst pipeline, физические операторы, Exchange/Shuffle, join-стратегии, AQE и диагностика узких мест через df.explain(mode="formatted").
Почему EXPLAIN - главный инструмент Spark-инженера¶
Когда DataFrame или SQL-запрос работает медленнее ожидаемого, у разработчика есть два пути. Первый - угадывать причину и вносить изменения наугад. Второй - прочитать план выполнения и точно понять, что происходит внутри Spark.
EXPLAIN - это рентгеновский снимок вашего запроса. Он показывает, как Spark намерен выполнить вашу трансформацию: какие данные прочитает, какие операции применит, где будет перераспределять данные по сети (Shuffle), какую стратегию выберет для join.
Без понимания плана выполнения невозможно:
- Объяснить, почему один запрос занимает 10 минут, а другой - 30 секунд
- Убедиться, что Predicate Pushdown действительно сработал и Spark не читает лишние данные
- Диагностировать Data Skew - перекос нагрузки на отдельные партиции
- Понять, почему вместо ожидаемого BroadcastHashJoin Spark использует дорогой SortMergeJoin
Чтение execution plan - это навык, а не талант. Как только вы научитесь видеть ключевые операторы и понимать их стоимость, вы сможете находить проблемы за минуты, а не часы отладки.
Catalyst Optimizer: четыре стадии трансформации¶
Прежде чем читать вывод EXPLAIN, нужно понять, что именно вы читаете. Каждый запрос в Spark проходит через Catalyst Optimizer - конвейер из четырёх последовательных трансформаций:
Parsed Logical Plan - первичный разбор синтаксиса SQL или DataFrame API. Spark создаёт дерево узлов из вашего кода, но ещё не знает, существуют ли таблицы и колонки, на которые вы ссылаетесь. Имена колонок и таблиц здесь - «Unresolved» (нерезолюированные).
Analyzed Logical Plan - Spark обращается к Catalog (Hive Metastore или внутренний SparkSession catalog) и разрешает все ссылки: проверяет, существуют ли таблицы, колонки, есть ли у пользователя доступ, совпадают ли типы данных. Если на этом этапе что-то не так - получите AnalysisException.
Optimized Logical Plan - Catalyst применяет набор правил оптимизации к логическому плану. Здесь происходит магия: фильтры перемещаются ближе к источнику данных (Predicate Pushdown), ненужные колонки удаляются (Column Pruning), константные выражения вычисляются заранее (Constant Folding). Catalyst применяет правила итеративно, пока план не стабилизируется.
Physical Plan - Catalyst выбирает конкретные алгоритмы выполнения для каждого логического оператора. Логический Join превращается в BroadcastHashJoin, SortMergeJoin или ShuffleHashJoin в зависимости от размеров таблиц и конфигурации. Логический Aggregate - в HashAggregate или SortAggregate. Это именно то, что Spark будет выполнять на кластере.
WholeStageCodeGen - дополнительная фаза, где несколько физических операторов объединяются в один кусок JVM-кода. Вместо вызова метода для каждого оператора по отдельности Spark генерирует единую функцию, которая применяет все операторы за один проход по данным. Это даёт 2–5x ускорение на CPU-bound операциях.
Режимы explain: от простого к детальному¶
explain() имеет пять режимов. Выбор зависит от задачи:
df = spark.read.parquet("s3://bucket/events/") \
.filter("amount > 100") \
.groupBy("user_id") \
.agg({"amount": "sum"})
# Режим 1: только физический план (дефолт)
df.explain()
df.explain(mode="simple")
# Режим 2: все четыре плана - parsed, analyzed, optimized, physical
df.explain(mode="extended")
# Режим 3: физический план + статистика (размеры, число строк)
df.explain(mode="cost")
# Режим 4: физический план + сгенерированный JVM-код для каждого WholeStageCodeGen
df.explain(mode="codegen")
# Режим 5: структурированный вывод - отдельные секции для каждой части плана
df.explain(mode="formatted")
В SQL-запросах:
EXPLAIN SELECT user_id, SUM(amount) FROM events WHERE amount > 100 GROUP BY user_id;
EXPLAIN FORMATTED SELECT user_id, SUM(amount) FROM events WHERE amount > 100 GROUP BY user_id;
EXPLAIN EXTENDED SELECT ...;
Почему formatted лучше для диагностики? Режим simple даёт компактный вывод, но всё в одной стене текста. Режим extended показывает все планы, но тоже без структуры. Режим formatted разбивает вывод на чёткие секции: план-дерево, детали каждого оператора, схема. Это позволяет находить нужную информацию быстро.
Структура вывода EXPLAIN FORMATTED¶
Рассмотрим конкретный пример. Создадим запрос и разберём вывод:
from pyspark.sql import functions as F
events = spark.read.parquet("s3://bucket/events/")
users = spark.read.parquet("s3://bucket/users/")
result = (
events
.filter(F.col("amount") > 100)
.join(users, on="user_id", how="left")
.groupBy("user_id", "region")
.agg(F.sum("amount").alias("total_amount"))
)
result.explain(mode="formatted")
Вывод имеет три секции:
== Physical Plan ==
*(3) HashAggregate(keys=[user_id#10, region#42], functions=[sum(amount#11)])
+- Exchange hashpartitioning(user_id#10, region#42, 200), ENSURE_REQUIREMENTS, [id=#25]
+- *(2) HashAggregate(keys=[user_id#10, region#42], functions=[partial_sum(amount#11)])
+- *(2) BroadcastHashJoin [user_id#10], [user_id#55], LeftOuter, BuildRight, false
:- *(2) Filter (isnotnull(amount#11) AND (amount#11 > 100))
: +- *(2) ColumnarToRow
: +- Scan parquet s3://bucket/events/ [user_id#10,amount#11] PushedFilters: [IsNotNull(amount), ...], ReadSchema: ...
+- BroadcastExchange HashedRelationBroadcastMode(List(input[0, bigint, true])), [id=#20]
+- *(1) ColumnarToRow
+- Scan parquet s3://bucket/users/ [user_id#55,region#42] ReadSchema: ...
== Subqueries ==
(нет подзапросов)
== Details ==
(3) HashAggregate [codegen id : 3]
Input [3]: [user_id#10, region#42, amount#11]
Keys [2]: [user_id#10, region#42]
Functions [1]: [sum(amount#11)]
Aggregate Attributes [1]: [sum(amount)#88]
Results [2]: [user_id#10, region#42, sum(amount)#88 AS total_amount#90]
...
== Output Schema ==
root
|-- user_id: bigint (nullable = true)
|-- region: string (nullable = true)
|-- total_amount: double (nullable = true)
Секция Physical Plan: читаем снизу вверх¶
Физический план - это дерево, где листья (нижний уровень) - источники данных, а корень (верхний уровень) - финальный результат. Данные текут снизу вверх.
Правило чтения: всегда начинайте с нижнего уровня (Scan) и двигайтесь к верху. Это направление движения данных.
Scan parquet→ читает файлы с дискаColumnarToRow→ конвертирует из колоночного Arrow-формата в строчныйFilter→ применяет условиеamount > 100BroadcastHashJoin→ выполняет joinHashAggregate (partial)→ частичная агрегация локально на каждом исполнителеExchange→ Shuffle - перераспределение данных по сетиHashAggregate (final)→ финальная агрегация после Shuffle
Номера операторов ((1), (2), (3)) - это ID WholeStageCodeGen блоков. Операторы с одинаковым номером объединены в один генерированный JVM-метод. Смена номера означает границу между блоками кодогенерации (обычно там Exchange или InMemoryTableScan).
Символ * перед номером (*(2)) означает, что оператор участвует в WholeStageCodeGen. Оператор без * (например, BroadcastExchange) выполняется вне кодогенерации.
Секция Details¶
Здесь подробная информация о каждом операторе по ID: входные атрибуты (Input), ключи агрегации (Keys), применяемые функции (Functions), выходные атрибуты (Results). Если нужно разобраться, откуда берётся конкретная колонка или какое выражение вычисляется - смотрите сюда.
Секция Output Schema¶
Финальная схема результата с типами данных. Полезно для быстрой проверки: правильно ли Spark вывел тип, нет ли неожиданного nullable = true там, где ожидается false.
Parsed Logical Plan: первичный разбор¶
Parsed Plan - это абстрактное синтаксическое дерево (AST) вашего запроса. Spark разбирает SQL или цепочку DataFrame операций и строит дерево, где каждый узел - логическая операция.
# Посмотреть все планы: parsed, analyzed, optimized + physical
result.explain(mode="extended")
В секции == Parsed Logical Plan == вы увидите:
'Aggregate ['user_id, 'region], ['user_id, 'region, sum('amount) AS total_amount]
+- 'Join LeftOuter, ('user_id = 'user_id)
:- 'Filter ('amount > 100)
: +- 'UnresolvedRelation [events], [], false
+- 'UnresolvedRelation [users], [], false
Апостроф ' перед именами колонок и таблиц ('user_id, 'UnresolvedRelation) означает, что они ещё не разрешены - Spark не проверил, существуют ли они. Это нормально для Parsed Plan.
Когда смотреть Parsed Plan? Редко - в основном при отладке синтаксических ошибок или для понимания, как DataFrame API транслируется во внутреннее представление.
Analyzed Logical Plan: разрешение символов¶
На этапе Analyzed Plan Catalyst обращается к Catalog (SparkSession internal catalog или Hive Metastore) и разрешает все ссылки:
Aggregate [user_id#10L, region#42], [user_id#10L, region#42, sum(amount#11) AS total_amount#88]
+- Join LeftOuter, (user_id#10L = user_id#55L)
:- Filter (amount#11 > cast(100 as double))
: +- Relation [user_id#10L, amount#11, event_ts#12] parquet
+- Relation [user_id#55L, region#42, signup_date#56] parquet
Ключевые отличия от Parsed Plan:
- Нет апострофов - все символы разрешены
- Колонки получили суффиксы с ID (
user_id#10L,user_id#55L) - это уникальные идентификаторы атрибутов внутри Catalyst.#10Lозначает атрибут с ID=10, тип Long. Благодаря ID Catalyst точно знает, какаяuser_idиз какой таблицы, даже при неоднозначных именах. cast(100 as double)- Catalyst добавил неявное приведение типа, чтобы сравнениеamount (Double) > 100 (Int)было корректнымRelationс перечисленными колонками - Catalyst прочитал схему из метастора
Ошибки на этапе Analyzed: AnalysisException: Column 'xxx' does not exist, Cannot resolve column 'yyy', Reference to unresolved attribute. Если видите такую ошибку - значит Analyzed Plan не смог разрешить символ.
Optimized Logical Plan: магия Catalyst¶
Это самый интересный план. Catalyst применяет десятки правил оптимизации итеративно. Рассмотрим ключевые из них.
Predicate Pushdown¶
Фильтры перемещаются как можно ближе к источнику данных - чтобы прочитать как можно меньше данных.
Фильтр amount > 100 теперь является частью Scan events. Для Parquet-файлов это означает, что Spark передаёт условие непосредственно в Parquet-ридер, который может пропустить целые Row Groups на уровне файла - до того как данные попадут в JVM-память.
В выводе explain(mode="formatted") pushdown-фильтры видны в секции Scan:
Scan parquet s3://bucket/events/
PushedFilters: [IsNotNull(amount), GreaterThan(amount,100.0)]
ReadSchema: struct<user_id:bigint,amount:double>
PushedFilters - условия, переданные в Parquet-ридер. Они выполняются на уровне файла. ReadSchema - только те колонки, которые реально нужны запросу. Если в исходной схеме было 20 колонок, а нужны только 2 - Parquet прочитает только эти 2 (Column Pruning).
Column Pruning (Projection Pruning)¶
Catalyst убирает из плана все колонки, которые не нужны для финального результата. Для Parquet (columnar format) это означает, что ненужные колонки вообще не читаются с диска - фундаментальное преимущество колоночного хранения.
# Запрос читает только 2 колонки из 20 в файле
events.select("user_id", "amount").filter("amount > 100")
# ReadSchema будет: struct<user_id:bigint,amount:double>
# Остальные 18 колонок Parquet-ридер пропустит
Constant Folding¶
Константные выражения вычисляются один раз на стадии планирования, а не при обработке каждой строки:
# До оптимизации: вычисляется для каждой строки
events.filter(F.col("amount") > 100 * 1.2 + 20)
# После Constant Folding: 100 * 1.2 + 20 = 140.0 вычислено заранее
# В плане увидите: Filter (amount > 140.0)
Join Reordering и Filter Pushdown through Join¶
Catalyst перемещает фильтры через join-операторы, применяя их до выполнения join. Это уменьшает количество данных, участвующих в join:
# Без оптимизации: сначала join, потом фильтр
orders.join(products, "product_id").filter("price > 1000")
# После оптимизации: сначала фильтр, потом join (меньше данных в join)
# Filter может быть перемещён на сторону products: products.filter("price > 1000")
Physical Plan: выбор стратегий выполнения¶
Physical Plan - это конкретные алгоритмы. Здесь Catalyst принимает решения о том, как именно выполнять каждую логическую операцию. Разберём ключевые физические операторы.
FileScan: читаем данные с диска¶
Scan parquet s3://bucket/events/
[user_id#10L,amount#11,event_ts#12]
PushedFilters: [IsNotNull(amount), GreaterThan(amount,100.0)]
ReadSchema: struct<user_id:bigint,amount:double>
PartitionFilters: [isnotnull(event_date#13), (event_date#13 >= 2026-01-01)]
DataFilters: [isnotnull(amount#11), (amount#11 > 100.0)]
Внимательно изучайте эту секцию:
PushedFilters - условия, переданные Parquet-ридеру. Хорошо: их много, соответствуют вашим filter().
PartitionFilters - фильтрация по партиционным колонкам. Spark пропустит директории в S3/HDFS, не читая файлы. Если фильтрация по дате есть в вашем коде, но отсутствует в PartitionFilters - значит Spark читает всё, partition pruning не работает.
DataFilters - условия, применяемые к данным внутри файлов (уже прочитанных по PartitionFilters).
ReadSchema - реально читаемые колонки. Если здесь много колонок, хотя в запросе нужны единицы - возможно, Column Pruning не сработал.
Когда PushedFilters пустой? Для CSV и JSON pushdown не работает - эти форматы читаются полностью. Pushdown поддерживается для Parquet, ORC и Delta Lake. Для Iceberg - pushdown через кастомный ридер.
Filter и Project: базовые трансформации¶
*(2) Filter (isnotnull(amount#11) AND (amount#11 > 100.0))
+- *(2) ColumnarToRow
+- Scan parquet ...
Filter - применяет предикат к каждой строке. Обычно объединён с соседними операторами через WholeStageCodeGen (символ *). Дополнительный isnotnull добавлен Catalyst автоматически - он знает, что сравнение с NULL вернёт NULL, а не True.
Project - выбор и вычисление выражений для набора колонок. Появляется после withColumn, select, вычисления алиасов:
*(2) Project [user_id#10L, amount#11, (amount#11 * 1.2) AS amount_with_tax#88]
+- *(2) Filter ...
ColumnarToRow - конвертация из колоночного Arrow/Parquet формата в строчный формат Spark. Обязательная операция при чтении Parquet с включённым Arrow. После этого оператора данные становятся доступны для строчной обработки.
Exchange: главный источник дорогих операций¶
Exchange - это физический оператор Shuffle. Перераспределение данных по сети между исполнителями. Каждый Exchange означает: данные покидают память исполнителей, сериализуются, передаются по сети, десериализуются на других исполнителях. Это дорого.
Типы Exchange:
hashpartitioning(keys, N)- данные распределяются по хешу ключей. N - целевое число партиций (spark.sql.shuffle.partitions, дефолт 200). Появляется приgroupBy,joinпо ключу,windowсpartitionBy.rangepartitioning(col, N)- данные сортируются и распределяются по диапазонам. Появляется приorderByбезpartitionBy, при window функциях безpartitionBy.RoundRobinPartitioning(N)- данные распределяются равномерно без учёта ключей. Появляется приrepartition(N)без указания колонки. Не сортирует, просто перераспределяет.SinglePartition- все данные собираются в одну партицию. Появляется приcoalesce(1)или операциях, требующих полной видимости данных.
Сколько Exchange допустимо? Зависит от запроса, но как ориентир: один groupBy = один Exchange, один SortMergeJoin = два Exchange (по одному на каждую сторону). Если видите 5+ Exchange на простом запросе - это повод разобраться.
Число 200 в hashpartitioning(..., 200) - это spark.sql.shuffle.partitions. В dev/test окружении уменьшайте до 4–8, чтобы не гонять 196 пустых партиций. В production настраивайте соразмерно объёму данных: ~128–200MB на партицию.
HashAggregate и SortAggregate¶
groupBy().agg() компилируется в двухфазную агрегацию:
*(3) HashAggregate(keys=[user_id#10, region#42], functions=[sum(amount#11)])
+- Exchange hashpartitioning(user_id#10, region#42, 200)
+- *(2) HashAggregate(keys=[user_id#10, region#42], functions=[partial_sum(amount#11)])
+- ...
Фаза 1 - partial aggregation: каждый исполнитель агрегирует свою партицию локально. partial_sum означает промежуточный результат. Это уменьшает объём данных, которые нужно передать через Shuffle.
Exchange: данные с одинаковыми ключами (user_id, region) попадают на один и тот же исполнитель.
Фаза 2 - final aggregation: каждый исполнитель финализирует агрегацию по своим ключам.
HashAggregate vs SortAggregate:
HashAggregate- хранит промежуточные результаты в хеш-таблице в памяти. Быстрее, но требует памяти. Используется когда агрегатная функция поддерживает "imperative" режим (большинство стандартных функций: sum, count, avg, min, max).SortAggregate- данные сначала сортируются по ключам, потом агрегируются. Медленнее, но не требует держать всё в памяти. Используется для функций, которые не поддерживают HashAggregate (например,collect_list,collect_set, некоторые кастомные UDAF).
# HashAggregate - стандартные функции
df.groupBy("category").agg(F.sum("amount"), F.avg("price"))
# SortAggregate - collect_list требует хранения всех элементов
df.groupBy("user_id").agg(F.collect_list("event_type"))
# В плане: SortAggregate (вместо HashAggregate)
Стратегии join¶
Выбор стратегии join - одно из ключевых решений Catalyst. Неправильный выбор может привести к тому, что join занимает часы вместо минут.
BroadcastHashJoin - самый эффективный¶
*(2) BroadcastHashJoin [user_id#10], [user_id#55], LeftOuter, BuildRight, false
:- *(2) Filter ...
: +- *(2) Scan parquet events/
+- BroadcastExchange HashedRelationBroadcastMode(List(input[0, bigint, true])), [id=#20]
+- *(1) Scan parquet users/
BuildRight - правая таблица (users) хешируется и рассылается всем исполнителям. Никакого Shuffle для левой (большой) таблицы! Данные events обрабатываются на месте.
BroadcastExchange - специальный Exchange для рассылки: данные отправляются с Driver ко всем исполнителям, а не перераспределяются между исполнителями.
Как форсировать Broadcast:
from pyspark.sql import functions as F
# Если Spark не выбрал BroadcastHashJoin автоматически - явный hint
orders.join(F.broadcast(country_codes), on="country_id", how="left")
# В плане: BroadcastHashJoin ... BuildRight
# Настроить порог для авто-broadcast (дефолт 10MB)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50mb")
Когда Broadcast не срабатывает автоматически:
- Spark не смог определить размер таблицы (нет статистики)
- Размер превышает
autoBroadcastJoinThreshold - Таблица создана через стриминг
SortMergeJoin - надёжный, но дорогой¶
SortMergeJoin [order_id#10], [order_id#55], Inner
:- *(3) Sort [order_id#10 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(order_id#10, 200), ENSURE_REQUIREMENTS
: +- *(2) Filter ...
: +- *(2) Scan parquet orders/
+- *(6) Sort [order_id#55 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(order_id#55, 200), ENSURE_REQUIREMENTS
+- *(5) Scan parquet payments/
Два Exchange (по одному на каждую сторону) + два Sort. Это четыре операции с Shuffle/Sort на двух таблицах. Зато надёжен для любых размеров - работает на диске через spill, если данные не помещаются в память.
Когда видеть SortMergeJoin в плане:
- Обе таблицы большие (> порога broadcast)
- Join по нескольким ключам
- Cross join (вырождается в
CartesianProductилиBroadcastNestedLoopJoin)
ShuffleHashJoin¶
ShuffleHashJoin [key#10], [key#55], Inner, BuildRight
Промежуточный вариант: Exchange происходит, но Sort не нужен - правая сторона хешируется в памяти. Эффективен, когда одна сторона достаточно мала для хеш-таблицы, но слишком велика для Broadcast.
CartesianProduct и BroadcastNestedLoopJoin - опасные операторы. Появляются при cross join или join без условия. BroadcastNestedLoopJoin используется когда нет условия по ключу (not-equi join):
# BroadcastNestedLoopJoin: нет условия по ключу
events.join(F.broadcast(thresholds), events["amount"] > thresholds["threshold"], how="left")
# В плане: BroadcastNestedLoopJoin BuildRight, LeftOuter
# Каждая строка events проверяется против каждой строки thresholds
WholeStageCodeGen: JVM-кодогенерация¶
*(2) HashAggregate(...)
+- *(2) BroadcastHashJoin ...
:- *(2) Filter ...
: +- *(2) ColumnarToRow
Символ * и одинаковые номера (*(2)) говорят, что все эти операторы объединены в один JVM-метод. Вместо:
row → Filter.process(row) → BroadcastHashJoin.process(row) → HashAggregate.process(row)
Spark генерирует примерно:
// Сгенерированный код (упрощённо)
void processRow(Row row) {
if (row.amount > 100.0 && row.amount != null) {
// inline join lookup
Row rightRow = broadcastTable.get(row.userId);
if (rightRow != null) {
// inline partial aggregation
hashMap.merge(row.userId, row.amount, Double::sum);
}
}
}
Это устраняет накладные расходы на виртуальные вызовы методов, промежуточные объекты Row, boxing/unboxing примитивов. Для CPU-intensive операций (фильтрация, вычисления) даёт 2–5x ускорение.
Операторы, которые разрывают WholeStageCodeGen:
Exchange(очевидно - данные уходят по сети)InMemoryTableScan(чтение из кеша)BroadcastExchange- Некоторые кастомные операторы
Смена номера в *(N) всегда означает, что между блоками есть один из этих разрывов. Наличие разрыва не всегда плохо (Exchange нужен для правильности), но лишние разрывы снижают производительность.
WindowExec: оконные функции в физическом плане¶
*(4) Window [sum(amount#11) windowspecdefinition(user_id#10, date#12 ASC NULLS FIRST, ...)], [user_id#10], [date#12 ASC NULLS FIRST]
+- *(4) Sort [user_id#10 ASC NULLS FIRST, date#12 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(user_id#10, 200), ENSURE_REQUIREMENTS
+- *(3) Scan parquet events/
Оконная функция всегда требует:
- Exchange - перераспределить данные по
partitionByключу (user_id) - Sort - отсортировать данные внутри каждой партиции по
orderByколонке (date) - WindowExec - применить оконную функцию к упорядоченным данным
Это означает, что каждая оконная функция с отличным partitionBy ключом добавляет Exchange + Sort. Несколько оконных функций с одинаковым partitionBy и orderBy оптимизируются в один Exchange + Sort:
from pyspark.sql.window import Window
# Один Exchange + Sort для обоих окон (одинаковый partitionBy + orderBy)
w = Window.partitionBy("user_id").orderBy("date")
df.withColumn("cumsum", F.sum("amount").over(w)) \
.withColumn("rank", F.rank().over(w))
# В плане: один Exchange + Sort, два WindowExec
# Два Exchange + Sort (разные окна)
w1 = Window.partitionBy("user_id").orderBy("date")
w2 = Window.partitionBy("category").orderBy("date")
df.withColumn("user_cumsum", F.sum("amount").over(w1)) \
.withColumn("cat_cumsum", F.sum("amount").over(w2))
# В плане: ДВА Exchange + ДВА Sort
InMemoryTableScan: cache() в плане¶
Когда DataFrame кешируется через cache() или persist():
df_cached = events.filter("amount > 100").cache()
df_cached.count() # Action - кешируется
# Следующий запрос с df_cached:
df_cached.groupBy("user_id").count().explain(mode="formatted")
В плане вместо Scan parquet появится:
*(1) InMemoryTableScan [user_id#10, amount#11], [isnotnull(amount#11)]
+- InMemoryRelation [user_id#10, amount#11], StorageLevel(disk, memory, ...)
+- *(1) Filter (isnotnull(amount#11) AND (amount#11 > 100.0))
+- *(1) Scan parquet s3://bucket/events/
InMemoryRelation показывает исходный план, который был закеширован. InMemoryTableScan - это чтение из кеша. Обратите внимание: InMemoryTableScan разрывает WholeStageCodeGen (нет символа *).
Важно: explain() показывает план до выполнения. Если DataFrame ещё не был материализован (нет count() или другого Action), InMemoryTableScan в плане есть, но кеш пустой - Spark будет читать с диска при первом обращении.
Adaptive Query Execution (AQE): план меняется во время выполнения¶
AQE - это механизм, при котором Spark пересматривает физический план во время выполнения, используя реальную статистику о данных. Включён по умолчанию с PySpark 3.2.
Что AQE может изменить:
- Смена стратегии join: если после первой стадии оказалось, что одна из таблиц маленькая - SortMergeJoin заменяется на BroadcastHashJoin.
- Coalescing партиций: если после Shuffle большинство партиций пустые (что типично при агрессивной фильтрации) - AQE автоматически объединяет пустые партиции. Контролируется
spark.sql.adaptive.advisoryPartitionSizeInBytes. - Обработка Data Skew: если одна партиция значительно больше других - AQE разбивает её на несколько для балансировки нагрузки.
Как выглядит AdaptivePlan в EXPLAIN:
result.explain(mode="formatted")
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
=== Final Plan (after execution) ===
*(3) HashAggregate ...
+- ...BroadcastHashJoin... ← AQE переключил на Broadcast
Если isFinalPlan=false - план может ещё измениться. isFinalPlan=true - это финальный план после всех AQE-пересмотров.
Настройки AQE:
# Включить/выключить AQE (включён по умолчанию в PySpark 3.2+)
spark.conf.set("spark.sql.adaptive.enabled", "true")
# Порог для авто-broadcast в AQE (может быть больше статического threshold)
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "30mb")
# Оптимальный размер партиции после coalescing
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128mb")
# Обработка skew
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
Partition Pruning: пропуск ненужных файлов¶
Когда таблица партиционирована на диске (например, по year/month/day), Spark может полностью пропустить директории, которые не нужны запросу:
events = spark.read.parquet("s3://bucket/events/")
# Таблица партиционирована: year=2024/month=01/, year=2024/month=02/, ...
# Spark прочитает только год=2024, месяц=05
events.filter("year = 2024 AND month = 5").groupBy("user_id").count()
В плане:
Scan parquet s3://bucket/events/
PartitionFilters: [isnotnull(year#20), (year#20 = 2024), isnotnull(month#21), (month#21 = 5)]
PushedFilters: []
ReadSchema: ...
PartitionFilters - условия применены к метаданным партиций. Spark не читает файлы из других директорий. Это может ускорить запрос в 10–100 раз на большом историческом архиве.
Когда Partition Pruning не работает:
- Фильтрация по вычисляемому выражению:
filter(F.year("event_date") == 2024)- Spark не может пропушитьF.year()в партицию - Использование
repartition()без колонки - теряется информация о партиционировании - Формат не поддерживает partition discovery (некоторые кастомные форматы)
# ПЛОХО: вычисляемое выражение - нет partition pruning
events.filter(F.year("event_date") == 2024)
# ХОРОШО: прямая фильтрация по партиционной колонке
events.filter(F.col("year") == 2024)
Связь физического плана со Spark UI¶
Execution plan - статический снимок намерений Spark. Spark UI показывает реальное выполнение с метриками:
Jobs → Stages → Tasks
Связь с планом:
- Каждый
Exchangeсоздаёт границу между Stage - Количество Tasks в Stage = количество партиций на входе
- Вкладка "SQL" в Spark UI показывает интерактивный план с метриками по каждому оператору
# Spark UI URL по умолчанию
print(spark.sparkContext.uiWebUrl) # http://driver:4040
В Spark UI вкладка SQL → Description показывает физический план с числами:
- number of output rows - реальное число строк на выходе каждого оператора
- data size - размер данных после Shuffle
- spill (memory/disk) - если данные не поместились в RAM и пролились на диск
Сравнивая план из explain() с метриками из Spark UI, можно найти:
- Где оценки Catalyst сильно расходятся с реальностью (плохая статистика)
- Где происходит Spill (нужно больше памяти или меньше партиций)
- Где одна Task значительно дольше других (Data Skew)
Практика: пошаговое чтение плана¶
Пример 1: Простой select + filter¶
df = spark.read.parquet("s3://bucket/events/")
result = df.filter(F.col("country") == "RU").select("user_id", "amount")
result.explain(mode="formatted")
== Physical Plan ==
*(1) Project [user_id#10, amount#11]
+- *(1) Filter (isnotnull(country#12) AND (country#12 = RU))
+- *(1) ColumnarToRow
+- Scan parquet s3://bucket/events/
[user_id#10, amount#11, country#12]
PushedFilters: [IsNotNull(country), EqualTo(country,RU)]
ReadSchema: struct<user_id:bigint,amount:double,country:string>
Читаем снизу вверх:
Scan parquet- читает только 3 нужных колонки (Column Pruning сработал)PushedFilters: [EqualTo(country,RU)]- условие передано в Parquet-ридер (Predicate Pushdown сработал)Filter- дополнительная проверка в JVM после чтения (на случай ложных срабатываний Min/Max статистики)Project- выбирает толькоuser_idиamount- Все операторы в блоке
*(1)- объединены в WholeStageCodeGen - Ноль Exchange - отлично, никакого Shuffle
Пример 2: GroupBy + Count¶
result = df.filter(F.col("country") == "RU").groupBy("category").count()
result.explain(mode="formatted")
== Physical Plan ==
*(3) HashAggregate(keys=[category#15], functions=[count(1)])
+- Exchange hashpartitioning(category#15, 200), ENSURE_REQUIREMENTS, [id=#10]
+- *(2) HashAggregate(keys=[category#15], functions=[partial_count(1)])
+- *(2) Filter (isnotnull(country#12) AND (country#12 = RU))
+- *(2) ColumnarToRow
+- Scan parquet s3://bucket/events/
PushedFilters: [IsNotNull(country), EqualTo(country,RU)]
ReadSchema: struct<category:string,country:string>
Читаем снизу вверх:
Scanчитает толькоcategoryиcountry- Column Pruning убралuser_id,amountи др.Filterс Predicate PushdownHashAggregate (partial)- каждый исполнитель считает частичные счётчики локально- Один Exchange (hashpartitioning по
category) - данные перераспределяются по категориям HashAggregate (final)- финальный подсчёт
Число 200 в hashpartitioning - дефолт. В dev-окружении: spark.conf.set("spark.sql.shuffle.partitions", "4").
Пример 3: Join - диагностика стратегии¶
big_table = spark.read.parquet("s3://bucket/events/") # 50 GB
small_table = spark.read.parquet("s3://bucket/countries/") # 100 KB
result = big_table.join(small_table, on="country_id", how="left")
result.explain(mode="formatted")
Ожидаемый план:
*(2) BroadcastHashJoin [country_id#10], [country_id#55], LeftOuter, BuildRight, false
:- *(2) Scan parquet s3://bucket/events/
+- BroadcastExchange HashedRelationBroadcastMode(...)
+- *(1) Scan parquet s3://bucket/countries/
Нет Exchange для большой таблицы, BroadcastExchange только для маленькой. Правильно.
Проблемный план (Broadcast не сработал):
SortMergeJoin [country_id#10], [country_id#55], LeftOuter
:- *(3) Sort [country_id#10 ASC NULLS FIRST]
: +- Exchange hashpartitioning(country_id#10, 200)
: +- *(2) Scan parquet s3://bucket/events/
+- *(6) Sort [country_id#55 ASC NULLS FIRST]
+- Exchange hashpartitioning(country_id#55, 200)
+- *(5) Scan parquet s3://bucket/countries/
Два Exchange, два Sort. Spark не выбрал Broadcast - возможно, не смог определить размер countries. Решение:
from pyspark.sql import functions as F
# Явный hint
big_table.join(F.broadcast(small_table), on="country_id", how="left")
Пример 4: Window функция - признак дорогого Sort¶
from pyspark.sql.window import Window
w = Window.partitionBy("user_id").orderBy("event_ts")
result = df.withColumn("prev_amount", F.lag("amount", 1).over(w))
result.explain(mode="formatted")
*(4) Window [lag(amount#11, 1, null) windowspecdefinition(user_id#10, event_ts#12 ASC NULLS FIRST)]
+- *(4) Sort [user_id#10 ASC NULLS FIRST, event_ts#12 ASC NULLS FIRST]
+- Exchange hashpartitioning(user_id#10, 200), ENSURE_REQUIREMENTS
+- *(3) Scan parquet ...
Exchange (по user_id) + Sort (по user_id, event_ts) + WindowExec. Это минимально необходимо для корректной оконной функции. Оптимизировать здесь нечего - это правильный план.
Если видите два Exchange для разных оконных функций - проверьте, можно ли использовать одинаковый partitionBy/orderBy для всех:
# Плохо: два Exchange
w1 = Window.partitionBy("user_id").orderBy("event_ts")
w2 = Window.partitionBy("user_id").orderBy("event_ts") # тот же, но другой объект
# Хорошо: один Exchange, один Sort
w = Window.partitionBy("user_id").orderBy("event_ts")
df.withColumn("lag_amount", F.lag("amount").over(w)) \
.withColumn("lead_amount", F.lead("amount").over(w))
Диагностика узких мест¶
Избыточный Shuffle¶
Признаки: много Exchange операторов в плане, медленная стадия Shuffle в Spark UI, большой объём Shuffle Write/Read в метриках Stage.
Что смотреть в плане:
Exchange RoundRobinPartitioningпослеrepartition(N)без колонки - данные перераспределяются без смыслаExchange hashpartitioningперед join при малой таблице - Broadcast не сработал- Несколько Exchange подряд без полезных операций между ними
Решения:
# Явный broadcast вместо SortMergeJoin
big.join(F.broadcast(small), on="key")
# Уменьшить число партиций после фильтрации
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# Убрать лишний repartition()
df.repartition(10) # если не нужен - удалить
Data Skew¶
Признаки: одна Task значительно медленнее остальных (в Spark UI: Tasks по Stage с огромным разбросом времени выполнения), тип операции - Exchange hashpartitioning.
В плане Data Skew не виден - это рантайм-проблема, связанная с неравномерным распределением ключей.
Решения:
# AQE автоматически обрабатывает skew (Spark 3.2+)
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
# Ручная обработка: Salting - добавить случайный суффикс к ключу
df_salted = df.withColumn(
"salted_key",
F.concat(F.col("skewed_key"), F.lit("_"), (F.rand() * 5).cast("int"))
)
CartesianProduct (NestedLoopJoin)¶
Признаки: в плане BroadcastNestedLoopJoin или CartesianProduct, огромное число строк на выходе.
# Опасный запрос - нет условия join
df1.join(df2) # DecartProduct: |df1| × |df2| строк
# BroadcastNestedLoopJoin: non-equi условие (не =, а > или <)
events.join(F.broadcast(thresholds), events["amount"] > thresholds["min_amount"])
# Каждая строка events проверяется против каждой строки thresholds
При обнаружении CartesianProduct без явного намерения - это баг в логике запроса.
Чек-лист для быстрой проверки плана¶
После написания нового запроса или оптимизации - проверьте план по следующему списку:
Сканирование данных:
PushedFiltersсодержат ваши условия фильтрации? (Predicate Pushdown сработал?)PartitionFiltersсодержат фильтрацию по партиционным колонкам? (Partition Pruning?)ReadSchemaсодержит только нужные колонки? (Column Pruning?)
Join:
- Ожидаемая стратегия join? (BroadcastHashJoin для маленьких таблиц, SortMergeJoin для больших)
- Нет ли
CartesianProductилиBroadcastNestedLoopJoinбез умысла? - Нет ли лишнего Exchange для join, который мог бы быть Broadcast?
Shuffle:
- Сколько
Exchangeв плане? Каждый ли оправдан? - Нет ли
RoundRobinPartitioning(бессмысленныйrepartition(N))? - Значение 200 в
hashpartitioning(..., 200)- подходит ли для ваших данных?
Агрегация:
- Используется
HashAggregate(быстро) илиSortAggregate(медленно)? ДляSortAggregate- можно ли заменитьcollect_listна другую функцию?
Кеш:
- Если ожидается
InMemoryTableScan- был ли выполнен Action для материализации кеша?
AQE:
isFinalPlan=false? - план ещё может измениться- Включён ли
spark.sql.adaptive.enabled?
Антипаттерны при работе с EXPLAIN¶
Игнорирование Exchange. "Запрос работает, зачем смотреть план?" - пока данных мало, Exchange незаметен. На 1TB Exchange становится узким местом. Читайте план до запуска на production-объёмах.
Слепое доверие оптимизатору. Catalyst хорош, но не идеален. Особенно без свежей статистики (ANALYZE TABLE). Явные hints (Broadcast, Repartition) часто нужны в production.
Оптимизация по плану без метрик. explain() показывает намерения, Spark UI показывает реальность. Убедитесь, что "оптимизированный" план действительно быстрее - иногда Catalyst знает лучше.
Чтение только верхнего уровня плана. Проблема обычно внизу - в FileScan (нет Pushdown), в Exchange (лишний Shuffle). Всегда читайте до конца дерева.
Не учитывать AQE. explain() до выполнения показывает начальный план. После AQE-оптимизации план может существенно измениться. Проверяйте финальный план в Spark UI.