Catalyst Optimizer: 4 фазы от AST до Physical Plan
Полный разбор Catalyst: парсинг SQL в AST, Analysis через Catalog, Rule-Based Optimization (predicate pushdown, column pruning, constant folding), Cost-Based Optimization, физическое планирование, Whole-Stage CodeGen и Tungsten. Как читать explain() и почему UDF ломает оптимизатор.
Вы пишете df.filter(col("age") > 18).groupBy("city").count() - и это работает быстрее эквивалентного кода на RDD. Почему? Потому что между вашим кодом и реальным выполнением стоит Catalyst Optimizer - компилятор запросов, который переписывает, переставляет и оптимизирует план прежде чем Spark сделает хотя бы одно обращение к данным.
Почему DataFrame быстрее RDD¶
С RDD вы диктуете как обрабатывать данные - каждый map, filter, reduceByKey выполняется именно так, как написан. Spark не вправе что-либо менять.
С DataFrame/SQL вы описываете что хотите получить. Catalyst получает это описание и сам решает, как его выполнить эффективнее. Он может переставить операции, объединить несколько шагов в один, выбросить ненужные колонки - всё то, что опытный инженер делал бы вручную.
Catalyst написан на Scala и использует функциональное программирование - деревья операторов, pattern matching и иммутабельные трансформации. Каждая фаза оптимизации - это набор правил, каждое из которых принимает дерево и возвращает новое (потенциально улучшенное) дерево.
Четыре фазы Catalyst¶
| Фаза | Вход | Выход | Что происходит |
|---|---|---|---|
| Parsing | SQL-строка / DF-операции | Unresolved Logical Plan | Синтаксический разбор в AST |
| Analysis | Unresolved Logical Plan | Resolved Logical Plan | Разрешение имён через Catalog |
| Logical Optimization | Resolved Logical Plan | Optimized Logical Plan | Применение правил оптимизации |
| Physical Planning | Optimized Logical Plan | Physical Plan(s) → лучший | Выбор алгоритмов (Join, Sort) |
| Code Generation | Physical Plan | Bytecode (RDD) | Компиляция в JVM-байткод |
Фаза 1: Parsing - от текста к дереву¶
Первая задача Catalyst - превратить SQL-строку или цепочку DataFrame-операций в Abstract Syntax Tree (AST).
AST (Abstract Syntax Tree) - дерево, где каждый узел - оператор (Filter, Aggregate, Join), а листья - колонки, константы, отношения. Имена в дереве пока неразрешены (unresolved): users - это строка, не реальная таблица; age - строка, не конкретная колонка с типом.
DataFrame API строит то же дерево - просто не через парсинг текста, а через вызовы методов. df.filter(col("age") > 18) создаёт узел Filter(age > 18, relation) без разбора текста.
Фаза 2: Analysis - разрешение имён¶
Analyzer берёт Unresolved Plan и превращает его в Resolved Plan, используя Catalog - реестр всех таблиц, вьюшек и колонок.
Что делает Analyzer:
- Resolution:
UnresolvedRelation["users"]→LogicalRelationс реальной схемой - Type coercion:
age > 18-18(Int) согласуется с типомage(Long). Еслиage- String, будет выброшена ошибка типов - Attribute binding: каждой колонке присваивается уникальный ExprId (
age#3), чтобы устранить неоднозначность при джойнах - Subquery resolution: вложенные SELECT/EXISTS разворачиваются и связываются с внешним планом
Если Analysis не удалась (таблица не существует, колонка не найдена, типы несовместимы) - вы получаете AnalysisException сразу, без запуска Job.
Фаза 3: Logical Optimization - переписывание дерева¶
Это главная интеллектуальная работа Catalyst. Optimizer применяет набор правил (Rules) к Resolved Plan. Каждое правило - функция (LogicalPlan) → LogicalPlan, которая ищет в дереве определённый паттерн и заменяет его улучшенной версией.
Правила применяются итеративно в батчах: каждый батч содержит несколько правил, которые применяются по кругу, пока план перестаёт меняться (fixpoint).
Predicate Pushdown¶
Одно из самых важных правил: фильтры опускаются вниз по дереву, как можно ближе к источнику данных.
До: джойн объединяет весь датасет users (1 млрд строк) с orders, затем применяется фильтр.
После: сначала фильтруется users (остаётся 10 млн строк из Москвы), потом джойн. Объём данных на джойне сократился в 100 раз.
Predicate Pushdown работает не только на уровне Plan - для форматов Parquet и ORC он пробрасывается ещё глубже: в Parquet Reader, который пропускает Row Groups не удовлетворяющие условию, не читая их вовсе.
Column Pruning (Projection Pruning)¶
Удаление колонок, которые не нужны для результата, как можно раньше:
Для Parquet это особенно важно - формат колончатый, и чтение двух колонок из двадцати в 10 раз быстрее чтения всей строки.
Constant Folding¶
Вычисление константных выражений на этапе оптимизации, а не при обработке каждой строки:
# Написано так (вычисляется N раз в runtime):
df.filter(col("price") > 100 * 1.2 * (1 + 0.05))
# Catalyst упрощает до (вычисляется один раз при планировании):
df.filter(col("price") > 126.0)
Другие примеры constant folding:
# Исходно
df.filter(lit(True) & (col("age") > 18)) # → df.filter(col("age") > 18)
df.filter(lit(False) | (col("age") > 18)) # → df.filter(col("age") > 18)
df.filter(col("age") > 18 & lit(True)) # → df.filter(col("age") > 18)
# CURRENT_DATE вычисляется один раз
df.filter(col("date") > current_date() - expr("INTERVAL 30 DAYS"))
Boolean Simplification¶
# NOT NOT → ничего
df.filter(~(~col("active"))) # → df.filter(col("active"))
# DeMorgan
df.filter(~(col("a") & col("b"))) # → df.filter(~col("a") | ~col("b"))
# Тавтология и противоречие
df.filter(col("age") > 0 & col("age") > 0) # → df.filter(col("age") > 0)
df.filter(col("age") > 18 & col("age") < 10) # → пустой результат (contradicts)
Другие правила оптимизации¶
| Правило | Что делает |
|---|---|
| Join Reordering | Меняет порядок джойнов чтобы маленькие таблицы джойнились первыми |
| Subquery Elimination | Преобразует коррелированные подзапросы в join |
| Limit Pushdown | Опускает LIMIT вниз чтобы быстрее остановить сканирование |
| IN → Join | WHERE id IN (SELECT id FROM ...) превращается в semi-join |
| Null Propagation | Упрощает выражения с NULL (NULL = x → NULL, убирает проверки) |
| Common Sub-Expression Elimination | Вычисляет одинаковые выражения один раз |
Фаза 4: Physical Planning - выбор алгоритмов¶
Optimized Logical Plan - это всё ещё абстрактное описание: «сджойни эти две таблицы», «сгруппируй по городу». Physical Planner выбирает конкретные алгоритмы.
Стратегии Join¶
Это самое важное решение Physical Planner:
| Стратегия Join | Когда применяется | Shuffle | Плюсы |
|---|---|---|---|
| Broadcast Hash Join | Один датасет < autoBroadcastJoinThreshold (10 MB) |
Нет | Нет shuffle, быстро |
| Shuffle Hash Join | Средние датасеты, один вписывается в память | Да (один) | Без сортировки |
| Sort Merge Join | Большие датасеты, нет памяти для хеш-таблицы | Да (оба) | Масштабируется |
| Cartesian Product | JOIN без условия (CROSS JOIN) | Да | Редко нужен |
# Явное управление стратегией
from pyspark.sql.functions import broadcast
# Принудительный Broadcast
result = large_df.join(broadcast(small_df), "user_id")
# Порог Broadcast (default 10 MB)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "50mb")
# Отключить Broadcast полностью
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
Cost-Based Optimization (CBO)¶
Rule-based оптимизация работает без данных о реальных таблицах. CBO использует статистику - размеры таблиц, количество строк, кардинальность колонок - чтобы принять более точные решения.
Без статистики Catalyst угадывает размеры и может выбрать неоптимальный план. С статистикой он точно знает, что orders содержит 100 млн строк, а order_status - 5 строк, и выбирает Broadcast для order_status.
Сбор статистики:
-- Собрать статистику для таблицы
ANALYZE TABLE users COMPUTE STATISTICS;
-- С гистограммами по колонкам (для CBO join reordering)
ANALYZE TABLE users COMPUTE STATISTICS FOR COLUMNS age, city, country;
# Включить CBO
spark.conf.set("spark.sql.cbo.enabled", "true")
spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")
# Проверить статистику
spark.sql("DESCRIBE EXTENDED users").show(50, truncate=False)
Статистика хранится в Catalog (Hive Metastore или Spark InMemoryCatalog). После ALTER TABLE или INSERT данные устаревают - нужно пересобирать статистику.
Whole-Stage Code Generation и Tungsten¶
После выбора Physical Plan Spark не интерпретирует его построчно - он компилирует в JVM-байткод.
Проблема интерпретации¶
Без Code Generation каждая запись проходит через цепочку виртуальных вызовов методов:
record → Filter.next() → call Project.next() → call Aggregate.next() → ...
Каждый next() - виртуальный вызов в JVM. JIT-компилятор не может их инлайнить, процессор не может предсказать переход. Накладные расходы на миллиард записей существенны.
Whole-Stage Code Generation¶
Spark генерирует единую Java-функцию для всего Stage:
Сгенерированный код:
- Нет виртуальных вызовов - все операции инлайнены
- JIT видит простой цикл - эффективно компилирует в машинный код
- Данные остаются в CPU-регистрах/кэше между операциями
Генерация происходит через Janino (runtime Java compiler) - строка Java-кода компилируется в класс при планировании Job.
Tungsten: управление памятью вне JVM heap¶
Параллельно с Code Generation, Tungsten управляет памятью напрямую, обходя JVM GC:
UnsafeRow - бинарное представление строки в off-heap памяти. Поле age (Int, 4 байта) лежит по фиксированному смещению от начала строки - доступ без десериализации, без создания Java-объекта, без GC.
Это объясняет, почему операции DataFrame быстрее, чем работа с обычными Java/Python объектами в UDF.
Чтение explain(): все четыре плана¶
explain() - главный инструмент диагностики плана выполнения:
df = spark.read.parquet("s3a://lake/users/") \
.filter(col("age") > 18) \
.join(spark.read.parquet("s3a://lake/orders/"), "user_id") \
.groupBy("city") \
.agg(count("*").alias("cnt"))
# Короткий вывод (только Physical Plan)
df.explain()
# Полный вывод: все 4 плана
df.explain(mode="extended")
# или
df.explain(True)
# Форматированный (рекомендуется)
df.explain(mode="formatted")
# Только cost model (с CBO)
df.explain(mode="cost")
# Кодогенерация
df.explain(mode="codegen")
Анатомия explain(mode="extended")¶
== Parsed Logical Plan ==
'Aggregate ['city], ['city, count(1) AS cnt]
+- 'Join Inner, ('user_id = 'user_id)
:- 'Filter ('age > 18)
: +- 'UnresolvedRelation [users] ← имена не разрешены, апостроф
+- 'UnresolvedRelation [orders]
== Analyzed Logical Plan ==
city: string, cnt: bigint
Aggregate [city#14], [city#14, count(1) AS cnt#22L]
+- Join Inner, (user_id#3L = user_id#31L) ← ExprId (#3, #31) разрешены
:- Filter (age#5 > 18)
: +- Relation[user_id#3L,age#5,city#14] parquet
+- Relation[user_id#31L,amount#35] parquet
== Optimized Logical Plan ==
Aggregate [city#14], [city#14, count(1) AS cnt#22L]
+- Project [city#14] ← Column Pruning: убраны user_id, age
+- Join Inner, (user_id#3L = user_id#31L)
:- Filter (isnotnull(age#5) AND (age#5 > 18)) ← isnotnull добавлен
: +- Filter (isnotnull(user_id#3L)) ← Predicate Pushdown
: +- Relation[...] parquet
+- Filter (isnotnull(user_id#31L))
+- Relation[...] parquet
== Physical Plan ==
*(3) HashAggregate(keys=[city#14], functions=[count(1)]) ← * = CodeGen
+- Exchange hashpartitioning(city#14, 200) ← Shuffle
+- *(2) HashAggregate(keys=[city#14], functions=[partial_count(1)])
+- *(2) Project [city#14]
+- *(2) SortMergeJoin [user_id#3L], [user_id#31L], Inner
:- *(1) Sort [user_id#3L ASC], false, 0
: +- Exchange hashpartitioning(user_id#3L, 200)
: +- *(1) Filter (isnotnull(age#5) AND (age#5 > 18))
: +- *(1) ColumnarToRow
: +- FileScan parquet [user_id#3L,age#5,city#14]
+- *(1) Sort [user_id#31L ASC], false, 0
+- Exchange hashpartitioning(user_id#31L, 200)
+- *(1) ColumnarToRow
+- FileScan parquet [user_id#31L] ← Column Pruning
Что читать:
*перед номером Stage → оператор участвует в Whole-Stage CodeGenExchange→ Shuffle (Stage Boundary)- Апостроф в
'UnresolvedRelation→ Unresolved (в Parsed Plan) FileScan parquet [col1, col2]→ сканируются только нужные колонки (Column Pruning)isnotnull(x)→ Catalyst добавил NULL-проверку автоматически
Ограничения Catalyst¶
UDF - чёрный ящик¶
Python UDF - главный враг оптимизатора. Catalyst не может заглянуть внутрь UDF и применить к ней правила:
from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType
# UDF - непрозрачна для Catalyst
@udf(BooleanType())
def is_adult(age):
return age > 18
# Catalyst не может опустить этот фильтр ниже join
df.filter(is_adult(col("age"))).join(orders, "user_id")
# ↑ джойн выполнится на всём датасете, затем UDF применится к результату
Вместо UDF используйте встроенные функции - они прозрачны для Catalyst:
# Catalyst видит col("age") > 18 и делает Predicate Pushdown
df.filter(col("age") > 18).join(orders, "user_id")
Pandas UDF лучше Python UDF, но тоже ограничены¶
Pandas UDF (vectorized) не делает построчных вызовов Python - данные передаются блоками через Apache Arrow. Это значительно быстрее, но Catalyst по-прежнему не может оптимизировать их семантику.
from pyspark.sql.functions import pandas_udf
import pandas as pd
@pandas_udf("boolean")
def is_adult_pd(age: pd.Series) -> pd.Series:
return age > 18
# Быстрее Python UDF, но Catalyst всё равно не видит внутрь
df.filter(is_adult_pd(col("age")))
Отсутствие статистики ломает CBO¶
Без ANALYZE TABLE Catalyst использует догадки о размерах данных. Типичная проблема - BroadcastHashJoin для больших таблиц (Catalyst думал, что таблица маленькая) или SortMergeJoin там, где можно было бы использовать BroadcastHashJoin.
# Симптом: большая таблица broadcast-ится и вызывает OOM Driver
# В explain() видно: BroadcastHashJoin для датасета без статистики
# Решение
spark.sql("ANALYZE TABLE orders COMPUTE STATISTICS")
# или отключить broadcast для подозрительных случаев
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
AQE: адаптивная оптимизация в runtime¶
Статический план - фиксирован до выполнения. AQE (Adaptive Query Execution) позволяет Catalyst корректировать план во время выполнения на основе реальных данных:
spark.conf.set("spark.sql.adaptive.enabled", "true") # default true в Spark 3.2+
# Что делает AQE:
# 1. Coalesce Partitions: объединяет мелкие shuffle-партиции в крупные
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128mb")
# 2. Skew Join: разбивает крупные skewed-партиции
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
# 3. Switch Join: переключает SMJ → BHJ если партнёр оказался маленьким
spark.conf.set("spark.sql.adaptive.localShuffleReader.enabled", "true")
Практика: оптимизация реального запроса¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, count, broadcast
spark = SparkSession.builder \
.master("local[*]") \
.appName("CatalystDemo") \
.config("spark.sql.shuffle.partitions", "4") \
.config("spark.sql.adaptive.enabled", "true") \
.getOrCreate()
spark.sparkContext.setLogLevel("WARN")
# Создаём датасеты
users = spark.createDataFrame([
(1, "Alice", "Moscow", 25),
(2, "Bob", "Berlin", 17),
(3, "Carol", "Moscow", 32),
(4, "Dave", "Paris", 15),
(5, "Eve", "Berlin", 28),
], ["user_id", "name", "city", "age"])
orders = spark.createDataFrame([
(101, 1, 1500), (102, 1, 200), (103, 3, 800),
(104, 5, 300), (105, 3, 950),
], ["order_id", "user_id", "amount"])
status = spark.createDataFrame([
("active",), ("inactive",), ("suspended",),
], ["status_name"])
# ── Плохой запрос: джойнит всё, потом фильтрует ──────────────────────
bad_query = (
users
.join(orders, "user_id") # джойн до фильтра
.filter(col("age") > 18) # фильтр после джойна
.groupBy("city")
.agg(count("*").alias("orders"))
)
print("=== ПЛОХОЙ ЗАПРОС ===")
bad_query.explain(mode="formatted")
# ── Хороший запрос: фильтр до джойна + broadcast ─────────────────────
good_query = (
users
.filter(col("age") > 18) # фильтр до джойна
.join(broadcast(orders), "user_id") # broadcast для маленькой таблицы
.groupBy("city")
.agg(count("*").alias("orders"))
)
print("=== ХОРОШИЙ ЗАПРОС ===")
good_query.explain(mode="formatted")
# Результаты одинаковы, но план разный
good_query.show()
Сравните вывод explain() обоих запросов:
- В плохом запросе:
SortMergeJoin+ фильтр после джойна - В хорошем запросе:
BroadcastHashJoin+ фильтр перед джойном (Catalyst может ещё pushdown'нуть его)
Важно: Catalyst сам применяет Predicate Pushdown - в реальности оба запроса могут дать одинаковый план после оптимизации. Проверьте explain() - не верьте только коду, смотрите на реальный план.
Best practices для Catalyst-friendly кода¶
Используйте встроенные функции вместо UDF:
# Плохо: UDF - непрозрачен для Catalyst
from pyspark.sql.functions import udf
parse_date = udf(lambda s: s[:10], StringType())
df.withColumn("date", parse_date(col("datetime")))
# Хорошо: встроенная функция - Catalyst видит и оптимизирует
from pyspark.sql.functions import substring
df.withColumn("date", substring(col("datetime"), 1, 10))
Проверяйте план перед запуском:
# Перед запуском больших Job
df.explain(mode="formatted")
# Убедитесь что видите:
# - FileScan читает только нужные колонки (Column Pruning)
# - Filter стоит ДО Join (Predicate Pushdown)
# - BroadcastHashJoin для маленьких таблиц (не SortMergeJoin)
# - * перед операторами (Whole-Stage CodeGen работает)
Собирайте статистику для CBO:
# После загрузки данных в managed таблицы
spark.sql("ANALYZE TABLE users COMPUTE STATISTICS FOR ALL COLUMNS")
Фильтруйте явно перед джойном (как страховка):
# Catalyst обычно сам делает Predicate Pushdown,
# но явный фильтр до джойна - документирует намерение
# и гарантирует поведение независимо от версии Spark
users_adult = users.filter(col("age") > 18)
result = users_adult.join(orders, "user_id")
Включите AQE для production:
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
Итог¶
Catalyst - это многофазный компилятор запросов:
- Parsing - SQL-текст или DataFrame API → AST (Unresolved Logical Plan)
- Analysis - разрешение имён через Catalog → Resolved Logical Plan
- Logical Optimization - применение правил: Predicate Pushdown, Column Pruning, Constant Folding, Join Reordering → Optimized Plan
- Physical Planning - выбор алгоритмов (BHJ vs SMJ vs SHJ), оценка стоимости через CBO → Physical Plan
- Code Generation (Tungsten) - компиляция в JVM-байткод через Whole-Stage CodeGen → выполнимые RDD
Catalyst объясняет: почему DataFrame быстрее RDD (оптимизация между слоями), почему UDF замедляет (непрозрачный чёрный ящик), почему ANALYZE TABLE важна (CBO не работает без статистики), и почему explain() - первый инструмент диагностики производительности.