Spark UI: вкладка SQL — чтение Physical Plan и метрики операторов

Полный разбор вкладки SQL в Spark UI: связь кода и Physical Plan, Logical→Optimized→Physical трансформации Catalyst, чтение дерева операторов снизу вверх, FileScan и Data Skipping метрики, HashAggregate и partial pre-aggregation, Whole-Stage Code Generation, AQE в реальном времени, сравнение с .explain().

core

1. Связь кода и вкладки SQL: Catalyst как переводчик

Вкладка SQL/Dataframe в Spark UI — это уникальное окно в работу оптимизатора Catalyst. Если вкладка Jobs показывает «что» выполняется, а Stages — «как быстро», то вкладка SQL отвечает на вопрос «почему именно так» — показывает физический план, который Catalyst построил для выполнения вашего запроса.

Важно понять: когда вы пишете df.filter(...).join(...).groupBy(...), вы создаёте логическую инструкцию. Catalyst транслирует её в физический план — конкретную последовательность JVM-операций с учётом статистики таблиц, конфигурации кластера и доступных стратегий выполнения. Вкладка SQL показывает именно этот физический план вместе с реальными метриками выполнения.

Что попадает на вкладку SQL

Схема показывает полный путь от кода до вкладки SQL: любая операция над DataFrame проходит через Catalyst и материализуется в физическом плане, который вы видите в UI. Ключевой момент — в UI вы видите Physical Plan, а не Logical. Это то, что реально выполнялось, со всеми оптимизациями Catalyst.

Разница между Logical и Physical Plan

Это принципиальное различие, которое часто путают:

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.getOrCreate()

df = spark.table("bronze.events")
result = (
    df
    .filter(F.col("amount") > 100)
    .join(spark.table("dim.customers"), "customer_id")
    .groupBy("region")
    .agg(F.sum("amount").alias("total"))
)

# df.explain() показывает ВСЕ планы — начиная от самого сырого
result.explain(mode="extended")
== Parsed Logical Plan ==
'Aggregate ['region], ['region, sum('amount) AS total#42]
+- 'Join Inner, ('customer_id = 'customer_id)
   :- 'Filter ('amount > 100)
   :  +- 'UnresolvedRelation [bronze, events], []
   +- 'UnresolvedRelation [dim, customers], []
# ↑ Буквальное дерево из кода, имена не разрешены

== Analyzed Logical Plan ==
region: string, total: double
Aggregate [region#12], [region#12, sum(amount#5) AS total#42]
+- Join Inner, (customer_id#3 = customer_id#20)
   :- Filter (amount#5 > cast(100 as double))
   :  +- Relation bronze.events[customer_id#3,amount#5,...]
   +- Relation dim.customers[customer_id#20,region#12,...]
# ↑ Имена разрешены через каталог, типы добавлены

== Optimized Logical Plan ==
Aggregate [region#12], [region#12, sum(amount#5) AS total#42]
+- Project [amount#5, region#12]                  ← Project Pushdown добавлен
   +- Join Inner, (customer_id#3 = customer_id#20)
      :- Filter (amount#5 > 100.0)                ← Predicate Pushdown (cast убран)
      :  +- Project [customer_id#3, amount#5]     ← Column Pruning
      :     +- Relation bronze.events[...]
      +- Project [customer_id#20, region#12]      ← Column Pruning
         +- Relation dim.customers[...]
# ↑ Catalyst применил оптимизации: убрал лишние колонки, вынес фильтр вниз

== Physical Plan ==
*(2) HashAggregate(keys=[region#12], functions=[sum(amount#5)])
+- Exchange hashpartitioning(region#12, 200), ENSURE_REQUIREMENTS
   +- *(1) HashAggregate(keys=[region#12], functions=[partial_sum(amount#5)])
      +- *(1) Project [amount#5, region#12]
         +- *(1) BroadcastHashJoin [customer_id#3], [customer_id#20], Inner
            :- *(1) Filter (amount#5 > 100.0)
            :  +- *(1) FileScan parquet bronze.events[...] PushedFilters: [IsNotNull, GreaterThan]
            +- BroadcastExchange HashedRelationBroadcastMode(...)
               +- *(1) FileScan parquet dim.customers[...]
# ↑ Конкретные алгоритмы! BroadcastHashJoin, Exchange (Shuffle), *(N) = WholeStageCodegen

Именно Physical Plan отображается на вкладке SQL в виде интерактивного графа. Всё что до него — промежуточные этапы оптимизации, видимые только через .explain(mode="extended").


2. Чтение дерева операторов: от FileScan до Action

Граф Physical Plan на вкладке SQL читается снизу вверх (или слева направо в горизонтальном представлении). Данные «течут» снизу вверх: источник данных (FileScan) внизу, финальный результат (Write/Action) вверху.

Карта основных операторов

Схема даёт полную карту операторов которые вы встретите в Physical Plan. Ключевые цвета: красный (Exchange и SortMergeJoin) = дорогие операции с Shuffle, зелёный (BroadcastHashJoin) = быстрая операция без Shuffle.

Правило «читать снизу вверх»

Запрос events.filter(amount > 100).join(customers).groupBy("region").sum("amount") выполняется так:

↑ HashAggregate[final]      ← 5. Финальная агрегация
│
↑ Exchange(region)           ← 4. Shuffle по ключу region
│
↑ HashAggregate[partial]     ← 3. Частичная агрегация ДО Shuffle
│
↑ BroadcastHashJoin          ← 2. Локальный JOIN (без Shuffle!)
│                   │
↑ Filter(amount>100)│         ↑ BroadcastExchange
│                   │             │
↑ FileScan(events)  │         ↑ FileScan(customers) ← 1. Начало: чтение файлов

Данные рождаются внизу (FileScan), проходят через Filter и BroadcastHashJoin, затем частично агрегируются, уходят на Shuffle и финально агрегируются наверху.


3. Метрики оператора FileScan: Data Skipping и Partition Pruning

FileScan — это первый и часто самый важный оператор в плане. Именно здесь решается сколько данных будет прочитано вообще — и именно здесь работают ключевые оптимизации Data Skipping.

Метрики FileScan в деталях

Нажав на узел FileScan в графе Spark UI, вы увидите список метрик. Важнейшие:

FileScan parquet bronze.events [event_id#1, customer_id#2, amount#5, dt#8]

PushedFilters: [IsNotNull(amount), GreaterThan(amount,100.0), EqualTo(dt,2024-01-15)]
PartitionFilters: [isnotnull(dt#8), (dt#8 = 2024-01-15)]
ReadSchema: struct<event_id:string,customer_id:string,amount:double>

--- Runtime Metrics ---
number of files read:           12                    ← реально прочитано файлов
metadata time:                  145ms                 ← время работы с метаданными
total size of files read:       2.3 GB               ← реальный I/O
number of output rows:          4,567,890             ← строк прошло фильтры
number of partitions read:      1 out of 365          ← Partition Pruning!
dynamic partition filters applied: true              ← AQE Runtime Filter

Три уровня фильтрации при чтении

Схема показывает три уровня фильтрации и их эффект: с 365 GB до 0.3% объёма. Именно это нужно видеть в метриках FileScan.

Диагностика проблем через FileScan метрики

def diagnose_filescan(files_read: int, files_total: int,
                       rows_output: int, rows_total: int,
                       metadata_time_ms: int) -> list[str]:
    """
    Анализирует метрики FileScan и выявляет проблемы.
    Данные берутся из раздела Runtime Metrics в Spark UI.
    """
    issues = []

    # Проверяем Partition Pruning
    partition_pruning_pct = (1 - files_read / max(files_total, 1)) * 100
    if partition_pruning_pct < 50 and files_total > 100:
        issues.append(
            f"Слабый Partition Pruning: {partition_pruning_pct:.0f}% файлов пропущено. "
            f"Убедитесь что фильтр WHERE по ключу партиционирования присутствует."
        )

    # Проверяем Row-level Filtering эффективность
    row_selectivity = rows_output / max(rows_total, 1)
    if row_selectivity > 0.8 and files_read > 100:
        issues.append(
            f"Низкая эффективность фильтрации: {row_selectivity:.0%} строк проходят фильтр. "
            f"Рассмотрите Z-Order/Liquid Clustering по фильтруемым колонкам."
        )

    # Проверяем Small Files Problem через metadata_time
    time_per_file_ms = metadata_time_ms / max(files_read, 1)
    if time_per_file_ms > 50:
        issues.append(
            f"Small Files Problem: {time_per_file_ms:.0f} мс/файл на метаданные. "
            f"Нормально < 10 мс/файл. Запустите OPTIMIZE для компакции файлов."
        )

    if not issues:
        issues.append(
            f"FileScan здоровый: {100-partition_pruning_pct:.0f}% файлов читается, "
            f"{row_selectivity:.1%} строк проходит фильтры."
        )

    return issues


# Пример: проблемный план
problems = diagnose_filescan(
    files_read=1200,       # прочитано файлов
    files_total=1200,      # всего файлов в таблице
    rows_output=50_000,    # строк прошло фильтры
    rows_total=10_000_000_000,  # строк в таблице
    metadata_time_ms=120_000    # 120 секунд на метаданные!
)
for p in problems:
    print(f"⚠️  {p}")
# ⚠️  Слабый Partition Pruning: 0% файлов пропущено.
# ⚠️  Small Files Problem: 100 мс/файл на метаданные.

4. Метрики агрегации: HashAggregate и Exchange

Операции агрегации в Spark почти всегда появляются в Physical Plan дважды: сначала как частичная (partial) агрегация до Shuffle, затем как финальная (final) после Shuffle. Понимание этой двухфазной схемы критично для диагностики.

Двухфазная агрегация и её цель

Диаграмма показывает ключевую оптимизацию: благодаря partial HashAggregate до Shuffle, по сети передаётся не 10 миллионов строк, а всего 4! В реальных данных это сокращение бывает в тысячи раз.

Когда pre-aggregation НЕ помогает

# Случай 1: Высокая кардинальность — pre-aggregation бесполезна
df.groupBy("user_id").count()  # Если 10M уникальных user_id!
# partial HashAggregate выдаст 10M строк (столько же сколько входных)
# → Shuffle получит все 10M строк, pre-aggregation = overhead без выгоды

# В Spark UI: узел HashAggregate[partial]
# number of output rows = number of input rows ← признак высокой кардинальности!
# Решение: увеличить число партиций или использовать approxCountDistinct

# Случай 2: Дорогие агрегатные функции
# Функции типа collect_set(), collect_list() нельзя частично агрегировать
# → Spark переключается на SortAgg вместо HashAgg (или выключает partial agg)

# В Spark UI: SortAggregate вместо HashAggregate
# Sort перед агрегацией → дополнительные накладные расходы

Что смотреть в метриках Exchange (Shuffle)

Exchange — узел Shuffle, самый дорогой в любом плане. Открыв его в Spark UI:

Exchange hashpartitioning(region#12, 200)

--- Runtime Metrics ---
data size:               45 MB     ← общий объём Shuffle данных
records written:         1,200     ← строк записано на Map стороне
shuffle bytes written:   45 MB
shuffle fetch time:      2.3 sec   ← время ожидания данных с других Executor

Что анализировать:

  • Если data size непропорционально велик → либо мало партиций (каждая огромная), либо нет partial aggregation
  • Если shuffle fetch time велик → медленные диски или сеть, или перегружена Map-сторона
  • Если records writtenrecords in → partial aggregation не сработала (высокая кардинальность)

5. Whole-Stage Code Generation: находим блоки WSCG

WholeStageCodegen (WSCG) — фундаментальная оптимизация движка Tungsten. Вместо того чтобы вызывать каждый оператор как отдельный метод Java (что требует дорогостоящей виртуальной диспетчеризации и передачи данных через объекты), Spark объединяет несколько операторов в одну скомпилированную функцию JIT.

Как WSCG выглядит в Physical Plan

В выводе .explain() операторы покрытые WSCG помечаются *(N):

== Physical Plan ==
*(2) HashAggregate(keys=[region#12], functions=[sum(amount#5)])   ← WSCG блок 2
+- Exchange hashpartitioning(region#12, 200)                       ← ГРАНИЦА!
   +- *(1) HashAggregate(keys=[region#12],                          ← WSCG блок 1
           functions=[partial_sum(amount#5)])
      +- *(1) Project [amount#5, region#12]                         ← часть блока 1
         +- *(1) BroadcastHashJoin [customer_id#3],                 ← часть блока 1
                  [customer_id#20], Inner
            :- *(1) Filter (amount#5 > 100.0)                       ← часть блока 1
            :  +- *(1) FileScan parquet [...]                        ← часть блока 1
            +- BroadcastExchange                                     ← НЕТ *(N) — вне WSCG
               +- *(3) FileScan parquet dim.customers [...]          ← WSCG блок 3

В этом примере три WSCG-блока: *(1) охватывает FileScan+Filter+BHJ+PartialAgg, *(2) — финальную агрегацию, *(3) — чтение dim таблицы. Exchange являются естественными границами WSCG-блоков.

Когда WSCG разрывается: Python UDF — главный виновник

from pyspark.sql import functions as F
from pyspark.sql.types import StringType

# ❌ Python UDF разрывает WSCG!
@F.udf(returnType=StringType())
def categorize(amount):
    if amount > 1000:
        return "high"
    elif amount > 100:
        return "medium"
    return "low"

df.withColumn("cat", categorize(F.col("amount"))).groupBy("cat").count()

# Physical Plan с Python UDF:
# HashAggregate[final]          ← WSCG блок X
# Exchange
# HashAggregate[partial]        ← WSCG блок Y
# ↑ ГРАНИЦА — Python UDF разрывает WSCG!
# BatchEvalPython               ← Python UDF батч исполнение (вне JVM!)
# ↑ ГРАНИЦА — Python UDF разрывает WSCG!
# FileScan

# Нет монолитного блока WSCG = каждый оператор вызывается по отдельности
# = виртуальная диспетчеризация = overhead в каждой строке!
# ✅ Заменяем Python UDF на нативные SQL функции → WSCG восстановлен

df.withColumn("cat",
    F.when(F.col("amount") > 1000, "high")
     .when(F.col("amount") > 100, "medium")
     .otherwise("low")
).groupBy("cat").count()

# Physical Plan теперь:
# *(2) HashAggregate[final]
# Exchange
# *(1) HashAggregate[partial]
#   *(1) Project (включая CASE WHEN — это НЕ Python!)
#     *(1) FileScan
# Монолитный блок *(1) охватывает всё от FileScan до PartialAgg!

Диагностика через .explain(): ищем разрывы WSCG

def find_wscg_breaks(explain_output: str) -> list[str]:
    """
    Анализирует текстовый вывод .explain() на наличие разрывов WSCG.
    Разрывы — операторы без префикса *(N).

    Returns: список строк без WSCG покрытия (потенциальные проблемы).
    """
    lines = explain_output.split('\n')
    breaks = []

    wscg_breakers = [
        "BatchEvalPython",       # Python UDF
        "ArrowEvalPython",       # Pandas UDF (через Arrow, меньше overhead)
        "ObjectHashAggregate",   # collect_set/collect_list
        "SortAggregate",         # Агрегация требующая Sort
        "Window",                # Оконные функции
        "Generate",              # explode/posexplode
    ]

    for line in lines:
        stripped = line.strip()
        # Оператор с WSCG имеет prefix *(N)
        has_wscg = stripped.startswith("*(")

        for breaker in wscg_breakers:
            if breaker in stripped and not has_wscg:
                breaks.append(f"WSCG разрыв: {stripped[:80]}")

    return breaks


# Использование:
plan = result.explain("formatted")
issues = find_wscg_breaks(str(plan))
for issue in issues:
    print(issue)

6. AQE в реальном времени: как увидеть изменения плана

Вкладка SQL — единственное место в Spark UI, где можно увидеть работу AQE в реальном времени: как план менялся по мере выполнения запроса.

Маркеры AQE на графе Physical Plan

В запущенном или завершённом запросе с AQE в Spark UI:

AdaptiveSparkPlan isFinalPlan=false   ← план ещё выполняется, может меняться

AdaptiveSparkPlan isFinalPlan=true    ← план финальный (запрос завершён)

В тексте .explain() для незавершённого AQE-запроса:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- == Current Plan ==
   SortMergeJoin [customer_id#3], [customer_id#20], Inner    ← план ПОКА ТАК
   ...

== Initial Plan ==                    ← что было запланировано изначально
...SortMergeJoin...

После завершения:

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=true
+- == Final Plan ==
   BroadcastHashJoin [customer_id#3], [customer_id#20], Inner  ← AQE ИЗМЕНИЛ SMJ → BHJ!

Три ключевых изменения которые делает AQE

spark = SparkSession.builder \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
    .config("spark.sql.adaptive.skewJoin.enabled", "true") \
    .getOrCreate()


# ─── Изменение 1: Динамическое схлопывание партиций Shuffle ─────────────
df.groupBy("category").sum("amount")
# Initial Plan: Exchange(200 партиций) [задано spark.sql.shuffle.partitions]
# Final Plan:   Exchange(8 партиций)  [AQE определил что данных мало]
# В Spark UI: Exchange узел показывает реальное число партиций

# ─── Изменение 2: Замена SMJ → BHJ ────────────────────────────────────
events.join(customers, "customer_id")
# Initial Plan: SortMergeJoin (customers кажется большой до выполнения)
# После filter на Stage 0: customers сжался с 10 GB до 50 MB
# Final Plan: BroadcastHashJoin (AQE увидел реальный размер!)

# Конфигурация порога для динамического upgrade:
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "50MB")
# Если таблица после Stage'а оказалась < 50MB → AQE переключит на BHJ

# ─── Изменение 3: Обработка перекошенных партиций (Skew Join) ──────────
events.join(users, "user_id")  # user_id=NULL = 30% данных → Skew!
# Initial Plan: SortMergeJoin
# AQE обнаруживает: Partition 42 = 50GB (остальные ~100MB)
# Final Plan: SortMergeJoin + SkewJoin optimization
# Partition 42 разбивается на N подпартиций → нет Straggler Task!

# В Spark UI → SQL → Plan: ищем "SkewJoin" или "skewed partitions" в метриках

Программная проверка что AQE применил оптимизацию

def check_aqe_optimizations(df) -> dict:
    """
    Проверяет был ли план изменён через AQE.
    Вызывайте ПОСЛЕ выполнения df.count() / .collect() / .write().
    """
    # Получаем финальный план в виде строки
    executed_plan = df._jdf.queryExecution().executedPlan().toString()
    spark_plan = df._jdf.queryExecution().sparkPlan().toString()

    result = {
        "aqe_enabled": "AdaptiveSparkPlan" in spark_plan,
        "is_final_plan": "isFinalPlan=true" in executed_plan,
        "optimizations": {
            "smj_converted_to_bhj": (
                "SortMergeJoin" in spark_plan and
                "BroadcastHashJoin" in executed_plan
            ),
            "partitions_coalesced": "CoalesceShufflePartitions" in executed_plan,
            "skew_handled": "SkewJoin" in executed_plan,
        }
    }

    return result

7. Сравнение с .explain(): текстовый план для автоматизации

Интерактивный граф удобен для ручного анализа, но в CI/CD и автоматических тестах нужен программный доступ к плану. Команда .explain() — ваш инструмент.

Четыре режима .explain()

df = spark.table("bronze.events").join(spark.table("dim.customers"), "customer_id")

# Режим 1: simple — только Physical Plan (дефолт)
df.explain()
# == Physical Plan ==
# BroadcastHashJoin [...]
# :(краткий план без деталей)

# Режим 2: extended — все 4 плана
df.explain(True)
# == Parsed Logical Plan ==
# == Analyzed Logical Plan ==
# == Optimized Logical Plan ==
# == Physical Plan ==

# Режим 3: codegen — физический план + сгенерированный Java код (Tungsten)
df.explain(mode="codegen")
# Показывает JAVA КОД, который Tungsten JIT скомпилирует!
# Полезно для понимания что именно работает быстро

# Режим 4: formatted — физический план с отступами и метриками (наиболее читаемый)
df.explain(mode="formatted")
# == Physical Plan ==
# * BroadcastHashJoin [customer_id#3], [customer_id#20], Inner, BuildRight
# :- * Filter (isnotnull(customer_id#3) AND (amount#5 > 100.0))
# :  +- * FileScan parquet bronze.events[...]
# |       PushedFilters: [IsNotNull(customer_id), IsNotNull(amount), GreaterThan(amount,100.0)]
# |       PartitionFilters: []
# +- BroadcastExchange HashedRelationBroadcastMode(List(cast(customer_id#20 as string)))
#    +- * FileScan parquet dim.customers[customer_id#20, region#12]
#          PushedFilters: [IsNotNull(customer_id)]
#          PartitionFilters: []

Как написать тест что план не деградировал

import pytest
from pyspark.sql import SparkSession

@pytest.fixture(scope="session")
def spark():
    return SparkSession.builder.master("local[2]").getOrCreate()


def assert_query_uses_broadcast_join(df, table_name: str) -> None:
    """
    Тест: проверяет что запрос использует BroadcastHashJoin
    а не более медленный SortMergeJoin.

    Добавьте в CI/CD чтобы не допускать регрессии плана.
    """
    plan = df._jdf.queryExecution().executedPlan().toString()

    assert "BroadcastHashJoin" in plan, (
        f"Ожидался BroadcastHashJoin для '{table_name}', "
        f"но план содержит SortMergeJoin. "
        f"Возможно статистика устарела. Запустите ANALYZE TABLE {table_name}."
    )


def assert_partition_pruning_applied(df, min_pruning_pct: float = 50.0) -> None:
    """
    Тест: проверяет что в плане есть PartitionFilters (Partition Pruning).
    Помогает отловить запросы читающие всю таблицу без фильтрации.
    """
    plan = df._jdf.queryExecution().executedPlan().toString()

    assert "PartitionFilters" in plan, (
        "В плане нет PartitionFilters — таблица читается полностью! "
        "Добавьте фильтр по ключу партиционирования в WHERE."
    )


def test_daily_report_plan(spark):
    """Регрессионный тест плана ежедневного отчёта."""
    events = spark.table("silver.events")
    customers = spark.table("dim.customers")

    query = (
        events
        .filter(F.col("event_date") == "2024-01-15")
        .join(customers, "customer_id")
        .groupBy("region")
        .agg(F.sum("amount"))
    )

    # Проверяем что план не деградировал
    assert_query_uses_broadcast_join(query, "dim.customers")
    assert_partition_pruning_applied(query)

    # Дополнительная проверка — нет Python UDF в плане
    plan = query._jdf.queryExecution().executedPlan().toString()
    assert "BatchEvalPython" not in plan, (
        "В плане обнаружен Python UDF (BatchEvalPython). "
        "Замените на нативные Spark SQL функции для производительности."
    )

Быстрый алгоритм анализа SQL вкладки за 5 минут

def analyze_sql_tab_quickly(df) -> dict:
    """
    Экспресс-анализ Physical Plan за 5 шагов.
    Даёт мгновенный диагноз без ручного просмотра графа.
    """
    plan = df._jdf.queryExecution().executedPlan().toString()
    optimized = df._jdf.queryExecution().optimizedPlan().toString()

    diagnostics = {}

    # ── Шаг 1: Проверяем дорогие операции ─────────────────────────────────
    diagnostics["has_sort_merge_join"] = "SortMergeJoin" in plan
    diagnostics["has_broadcast_join"]  = "BroadcastHashJoin" in plan
    diagnostics["has_exchange"]        = "Exchange" in plan
    diagnostics["exchange_count"]      = plan.count("Exchange")

    # ── Шаг 2: Проверяем применились ли оптимизации ────────────────────────
    diagnostics["partition_pruning"]   = "PartitionFilters" in plan
    diagnostics["predicate_pushdown"]  = "PushedFilters" in plan
    diagnostics["column_pruning"]      = "ReadSchema" in plan  # читаем не все колонки

    # ── Шаг 3: Проверяем WSCG покрытие ────────────────────────────────────
    wscg_count = sum(1 for line in plan.split('\n') if line.strip().startswith("*("))
    non_wscg_ops = [line for line in plan.split('\n')
                    if "EvalPython" in line or "SortAggregate" in line]
    diagnostics["wscg_blocks"]      = wscg_count
    diagnostics["wscg_breaks"]      = non_wscg_ops

    # ── Шаг 4: Проверяем AQE ────────────────────────────────────────────────
    diagnostics["aqe_active"]       = "AdaptiveSparkPlan" in plan

    # ── Шаг 5: Итоговые рекомендации ───────────────────────────────────────
    recommendations = []
    if diagnostics["has_sort_merge_join"] and not diagnostics["has_broadcast_join"]:
        recommendations.append(
            "JOIN использует SortMergeJoin. Убедитесь что статистика таблиц актуальна "
            "(ANALYZE TABLE). Если одна таблица < 200 MB — используйте F.broadcast()."
        )
    if not diagnostics["partition_pruning"]:
        recommendations.append(
            "Нет Partition Pruning. Добавьте фильтр WHERE по ключу партиционирования."
        )
    if diagnostics["wscg_breaks"]:
        recommendations.append(
            f"Найдены разрывы WSCG ({diagnostics['wscg_breaks']}). "
            "Замените Python UDF на нативные SQL функции."
        )
    if not recommendations:
        recommendations.append("План выглядит оптимально.")

    diagnostics["recommendations"] = recommendations

    return diagnostics


# Использование:
result_df = spark.table("bronze.events") \
    .filter(F.col("dt") == "2024-01-15") \
    .join(spark.table("dim.customers"), "customer_id") \
    .groupBy("region") \
    .agg(F.sum("amount"))

analysis = analyze_sql_tab_quickly(result_df)
print("=== Анализ Physical Plan ===")
for key, value in analysis.items():
    print(f"  {key}: {value}")

Итоги: методика анализа вкладки SQL за 3 минуты

Шаг 1 (30 сек): Откройте вкладку SQL → найдите самый долгий запрос → нажмите.

Шаг 2 (30 сек): Смотрите на граф Physical Plan. Ищете красные флаги:

  • SortMergeJoin с двумя Exchange → дорогой Shuffle с обеих сторон
  • BatchEvalPython → Python UDF разрывает WSCG
  • Очень длинная цепочка Exchange → много Shuffle операций

Шаг 3 (30 сек): Нажмите на узел FileScan → проверьте:

  • number of partitions read — партиционирование работает?
  • metadata time — не велика ли? (Small Files?)
  • number of output rows — сколько строк прошло фильтры?

Шаг 4 (30 сек): Найдите оператор Exchange → проверьте:

  • records written после HashAggregate[partial] — уменьшилось ли? (pre-agg эффективна?)
  • data size — соразмерно ли с объёмом исходных данных?

Шаг 5 (60 сек): Запустите .explain(mode="formatted") в коде → ищите:

  • *(N) — покрытие WSCG (чем больше операторов в одном *(N), тем лучше)
  • PushedFilters — предикатный pushdown работает?
  • PartitionFilters — партиции прунятся?
  • AdaptiveSparkPlan isFinalPlan=true — AQE завершил оптимизацию?

Итог: за 3 минуты вы можете определить:

  • Тип Join (быстрый BHJ или медленный SMJ)
  • Эффективность чтения данных (Data Skipping, Partition Pruning)
  • Наличие WSCG разрывов (Python UDF)
  • Работу AQE
  • Эффективность pre-aggregation перед Shuffle