EXPLAIN FORMATTED: чтение логического и физического плана пошагово

Как читать execution plan: Catalyst pipeline, физические операторы, Exchange/Shuffle, join-стратегии, AQE и диагностика узких мест через df.explain(mode="formatted").

core

Почему 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 > 100
  • BroadcastHashJoin → выполняет join
  • HashAggregate (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/

Оконная функция всегда требует:

  1. Exchange - перераспределить данные по partitionBy ключу (user_id)
  2. Sort - отсортировать данные внутри каждой партиции по orderBy колонке (date)
  3. 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 может изменить:

  1. Смена стратегии join: если после первой стадии оказалось, что одна из таблиц маленькая - SortMergeJoin заменяется на BroadcastHashJoin.
  2. Coalescing партиций: если после Shuffle большинство партиций пустые (что типично при агрессивной фильтрации) - AQE автоматически объединяет пустые партиции. Контролируется spark.sql.adaptive.advisoryPartitionSizeInBytes.
  3. Обработка 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 показывает реальное выполнение с метриками:

JobsStagesTasks

Связь с планом:

  • Каждый 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>

Читаем снизу вверх:

  1. Scan parquet - читает только 3 нужных колонки (Column Pruning сработал)
  2. PushedFilters: [EqualTo(country,RU)] - условие передано в Parquet-ридер (Predicate Pushdown сработал)
  3. Filter - дополнительная проверка в JVM после чтения (на случай ложных срабатываний Min/Max статистики)
  4. Project - выбирает только user_id и amount
  5. Все операторы в блоке *(1) - объединены в WholeStageCodeGen
  6. Ноль 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>

Читаем снизу вверх:

  1. Scan читает только category и country - Column Pruning убрал user_id, amount и др.
  2. Filter с Predicate Pushdown
  3. HashAggregate (partial) - каждый исполнитель считает частичные счётчики локально
  4. Один Exchange (hashpartitioning по category) - данные перераспределяются по категориям
  5. 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.