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().
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 written≈records 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