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 ломает оптимизатор.

core internals optimization

Вы пишете 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 = xNULL, убирает проверки)
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 CodeGen
  • Exchange → 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 - это многофазный компилятор запросов:

  1. Parsing - SQL-текст или DataFrame API → AST (Unresolved Logical Plan)
  2. Analysis - разрешение имён через Catalog → Resolved Logical Plan
  3. Logical Optimization - применение правил: Predicate Pushdown, Column Pruning, Constant Folding, Join Reordering → Optimized Plan
  4. Physical Planning - выбор алгоритмов (BHJ vs SMJ vs SHJ), оценка стоимости через CBO → Physical Plan
  5. Code Generation (Tungsten) - компиляция в JVM-байткод через Whole-Stage CodeGen → выполнимые RDD

Catalyst объясняет: почему DataFrame быстрее RDD (оптимизация между слоями), почему UDF замедляет (непрозрачный чёрный ящик), почему ANALYZE TABLE важна (CBO не работает без статистики), и почему explain() - первый инструмент диагностики производительности.