Rule-Based Optimization: constant folding, predicate pushdown, column pruning

Подробный разбор RBO в Catalyst: механика tree transformations и batch-правил, constant folding, expression и boolean simplification, predicate pushdown до уровня Parquet row groups, column pruning, null propagation, dead code elimination. Когда правила не работают и как писать optimizer-friendly код.

core internals optimization

Предыдущий урок описал четыре фазы Catalyst в целом. Этот - детальный разбор фазы Logical Optimization: как именно Catalyst переписывает дерево операторов, какие правила применяет, в каком порядке и почему они работают на любом объёме данных - от ста строк до петабайта.


Что такое Rule-Based Optimization

RBO - набор детерминированных правил преобразования плана запроса. Каждое правило не знает ничего о реальных данных: сколько строк в таблице, какое распределение значений, насколько заполнен кэш. Правило смотрит только на структуру дерева операторов и если видит определённый паттерн - применяет трансформацию.

Это принципиально отличает RBO от Cost-Based Optimization (CBO), который использует статистику данных. RBO работает всегда, даже без статистики.


Архитектура: tree transformations и батчи правил

Catalyst представляет план запроса как иммутабельное дерево (TreeNode). Каждый узел - оператор (Filter, Project, Join, Aggregate) или выражение (EqualTo, Add, Cast).

Правило - это функция LogicalPlan → LogicalPlan. Оно обходит дерево сверху вниз или снизу вверх, ищет подходящий паттерн и возвращает новое дерево с применённой трансформацией. Если паттерн не найден - дерево возвращается без изменений.

Батчи и fixpoint

Правила организованы в батчи (Batches). Каждый батч содержит связанные правила и имеет параметр стратегии выполнения:

  • FixedPoint - батч применяется повторно, пока план перестаёт меняться (fixpoint). Это нужно для правил, которые открывают возможности для других правил: после pushdown одного фильтра становится виден следующий.
  • Once - батч применяется ровно один раз.

Батч Operator Optimization - главный. В нём более 50 правил, и он работает в FixedPoint-режиме: Catalyst применяет весь набор снова и снова, пока план стабилизируется. Обычно это занимает 2–5 итераций.

Как выглядит правило изнутри (Scala)

// Упрощённый пример правила Constant Folding
object ConstantFolding extends Rule[LogicalPlan] {
  def apply(plan: LogicalPlan): LogicalPlan = plan transform {
    // Паттерн: сложение двух Integer-литералов
    case Add(Literal(a: Int, IntegerType), Literal(b: Int, IntegerType)) =>
      Literal(a + b, IntegerType)
  }
}

transform - метод TreeNode, который обходит дерево bottom-up и применяет частичную функцию (pattern matching) к каждому узлу. Если узел совпал с паттерном - он заменяется результатом. Если нет - возвращается как есть.


Constant Folding: вычисление до выполнения

Идея: если выражение содержит только константы - вычислить его один раз при планировании, а не повторять для каждой строки на Executor.

Арифметические константы

# Написано:
df.filter(col("price") > 100 * 1.2 * (1 + 0.05))

# После Constant Folding в плане:
# Filter (price > 126.0)
# Вычисление 100 * 1.2 * 1.05 = 126.0 происходит один раз на Driver
# Написано:
df.withColumn("minutes_per_week", col("hours_per_day") * 24 * 60 * 7)

# После Constant Folding:
# Project (hours_per_day * 10080)
# 24 * 60 * 7 = 10080 - свёрнуто

Cast константных литералов

# Написано:
df.filter(col("event_date") > lit("2024-01-01").cast("date"))

# После Constant Folding:
# Filter (event_date > 2024-01-01)  ← Cast выполнен при планировании

Без Constant Folding cast("2024-01-01" → DateType) вызывался бы для каждой строки. С ним - один раз.

Выражения с CURRENT_DATE / CURRENT_TIMESTAMP

from pyspark.sql.functions import current_date, expr

# Написано:
df.filter(col("created_at") > current_date() - expr("INTERVAL 30 DAYS"))

# current_date() вычисляется один раз при планировании Job,
# не при обработке каждой строки

Boolean Simplification и Expression Simplification

Catalyst применяет законы алгебры логики к булевым выражениям:

Тавтологии и противоречия

# Тавтология (всегда True) → убирается
df.filter(lit(True) & (col("age") > 18))
# → df.filter(col("age") > 18)

df.filter((col("age") > 0) | lit(True))
# → df (фильтр убран полностью)

# Противоречие (всегда False) → пустой результат
df.filter(lit(False) & (col("age") > 18))
# → EmptyRelation (Spark даже не читает данные)

df.filter((col("age") > 18) & (col("age") < 10))
# → EmptyRelation (математически невозможно)

Двойное отрицание

df.filter(~(~col("active")))
# → df.filter(col("active"))

Законы де Моргана

df.filter(~(col("city").isNull() & col("age").isNull()))
# → df.filter(col("city").isNotNull() | col("age").isNotNull())

IF и CASE упрощение

from pyspark.sql.functions import when

# IF(true, expr_a, expr_b) → expr_a
df.withColumn("x", when(lit(True), col("a")).otherwise(col("b")))
# → df.withColumn("x", col("a"))

# IF(false, expr_a, expr_b) → expr_b
df.withColumn("x", when(lit(False), col("a")).otherwise(col("b")))
# → df.withColumn("x", col("b"))

Упрощение сравнений

# NOT (a > b) → a <= b
df.filter(~(col("age") > 18))
# → df.filter(col("age") <= 18)

# a == a → true (для not-nullable колонок)
df.filter(col("id") === col("id"))
# → df (фильтр убран, если id не nullable)

Predicate Pushdown: фильтруй как можно раньше

Самое мощное правило RBO. Принцип: если фильтр может быть применён раньше в дереве - перемести его туда, не меняя семантики запроса.

Pushdown через Project

Project не меняет число строк - фильтр безопасно перемещается через него вниз.

Pushdown через Join

Самый важный вариант. Фильтр на одной стороне джойна можно опустить до самого джойна:

Джойн обрабатывает на 920 млн строк меньше. При shuffle это означает в разы меньше данных записано на диск и передано по сети.

Когда pushdown через join запрещён

Catalyst не может слепо пушить фильтры через join - для OUTER JOIN это меняет семантику:

# LEFT JOIN: нельзя пушить фильтр на правую таблицу вниз join
# (нарушит NULL-сохраняющую семантику LEFT JOIN)
users.join(orders, "user_id", "left") \
     .filter(col("orders.amount") > 1000)
# ← Catalyst знает: этот фильтр нельзя опустить ниже LEFT JOIN

# INNER JOIN: можно пушить на обе стороны
users.join(orders, "user_id", "inner") \
     .filter(col("users.city") == "Moscow")
# ← Безопасно переносится до join

Predicate Pushdown в Parquet/ORC

Это самый глубокий уровень - фильтр пробрасывается не только в дерево Spark-плана, но и в физический Parquet Reader:

Parquet хранит min/max статистику для каждого Row Group (блок данных ~128 MB). Если значение условия не попадает в диапазон [min, max] Row Group - Parquet Reader пропускает его полностью, не читая с диска.

Для ещё более точной фильтрации Parquet поддерживает Bloom Filters - вероятностную структуру данных для проверки принадлежности значения множеству:

# Написать Parquet с Bloom Filter по колонке user_id
df.write \
  .option("parquet.bloom.filter.enabled#user_id", "true") \
  .option("parquet.bloom.filter.expected.ndv#user_id", "1000000") \
  .parquet("output/")

# При чтении с фильтром по user_id:
# Spark проверяет Bloom Filter → если значения точно нет → skip Row Group
spark.read.parquet("output/").filter(col("user_id") == 42).show()

Bloom Filter гарантирует: если значение ТОЧНО отсутствует в Row Group - filter вернёт false (без false negatives). Если присутствует - нужно читать (возможны false positives).

Проверка: что попало в pushdown

df = spark.read.parquet("s3a://lake/events/") \
    .filter(col("city") == "Moscow") \
    .filter(col("date") >= "2024-01-01")

df.explain(mode="formatted")

В выводе explain() ищем PushedFilters в секции FileScan:

FileScan parquet [...] PushedFilters: [IsNotNull(city),
  EqualTo(city,Moscow), IsNotNull(date), GreaterThanOrEqual(date,2024-01-01)],
ReadSchema: struct<city:string,date:date,amount:double>

PushedFilters - фильтры, отданные Parquet Reader. ReadSchema - только нужные колонки (Column Pruning).


Column Pruning: не читай лишнего

Идея: если в финальном результате нужны только колонки A и B - зачем читать C, D, E, ... Z?

Catalyst анализирует весь план снизу вверх и выясняет, какие колонки реально нужны на каждом шаге. Все остальные - «обрезаются» как можно раньше.

Для таблицы с 20 колонками, где нужны только 2, это даёт:

  • В 10 раз меньше данных читается с диска (Parquet: только нужные column chunks)
  • В 10 раз меньше данных передаётся по сети при shuffle
  • Меньше памяти на Executor

Как Catalyst определяет нужные колонки

Catalyst обходит план снизу вверх, собирая множество «требуемых» атрибутов:

Aggregate [city] → требует: city
↑
Join (user_id = user_id) → требует: user_id (для join), city (из Aggregate)
↑
Filter (age > 18) → требует: age, user_id, city
↑
Scan users → читать: age, user_id, city (3 из 20)

Projection Pushdown в Parquet

Column Pruning в Parquet реализуется через projection pushdown - Parquet Reader читает только нужные column chunks, пропуская все остальные на уровне файла.

Это возможно благодаря колончатому формату Parquet: данные каждой колонки хранятся физически отдельно. Чтение двух колонок из двадцати требует чтения двух column chunks, остальные восемнадцать даже не открываются.

# В explain() видна секция ReadSchema
spark.read.parquet("users/") \
    .select("city", "user_id") \
    .filter(col("age") > 18) \
    .explain(mode="formatted")

# FileScan parquet
# ReadSchema: struct<user_id:bigint,age:int,city:string>
# ↑ только 3 колонки из всех имеющихся

Null Propagation

SQL имеет трёхзначную логику: TRUE, FALSE, NULL. Catalyst оптимизирует NULL-выражения:

# NULL в арифметике → NULL
# NULL + 1 = NULL, NULL * 0 = NULL (не 0!)
df.withColumn("x", col("a") + lit(None))
# → withColumn("x", lit(None).cast(type))  - всегда NULL

# NULL сравнение → NULL (не FALSE!)
df.filter(col("city") == lit(None))  # Всегда пустой результат
# Правильно: df.filter(col("city").isNull())

# Catalyst добавляет isNotNull при необходимости
df.filter(col("age") > 18)
# В плане: Filter (isnotnull(age) AND age > 18)
# Это оптимизация: строки с NULL age отсеиваются без вычисления > 18

Null-safe операции

# == для NULL: всегда NULL
col("a") == col("b")          # NULL если любой из них NULL

# <=> для NULL: TRUE если оба NULL, FALSE иначе (null-safe equality)
col("a").eqNullSafe(col("b")) # работает корректно с NULL

Catalyst знает об этих правилах и может упростить выражения, содержащие NULL-литералы, ещё на этапе планирования.


Dead Code Elimination

Catalyst удаляет вычисления, результат которых никогда не используется:

# Добавлена колонка, которая нигде не используется
df.withColumn("temp", col("a") * col("b") * 1000) \
  .select("id", "city")  # "temp" не выбрана

# Catalyst видит: "temp" не нужна ни в select, ни в downstream
# withColumn("temp", ...) убирается из плана полностью
# Промежуточный DataFrame создан, но его колонки не используются ниже
intermediate = df.withColumn("x", col("a") + col("b")) \
                 .withColumn("y", col("c") * col("d"))

result = intermediate.select("id", "e")
# x и y не нужны → их вычисление убирается

Пример из реальной жизни: избыточные джойны

# Joinим dimension_table, но используем только ключ джойна
result = fact.join(dim, "product_id") \
             .select("fact.event_id", "fact.amount")
# Catalyst видит: ни одна колонка из dim не используется
# SortMergeJoin → может быть упрощён или убран

Комбинация правил: как они усиливают друг друга

Правила работают в синергии. Один проход открывает возможности для другого:

Шаги:

  1. Constant Folding: 10 + 818. Теперь план выглядит чище
  2. Predicate Pushdown: Filter (age > 18) опускается ниже Project и ближе к Scan
  3. Column Pruning: зная что нужны только id, city (из Project) и age (для Filter) - Scan обрезается до трёх колонок

Без FixedPoint-итерации шаг 3 мог бы не увидеть результат шагов 1 и 2.


Порядок правил имеет значение

Catalyst применяет правила в определённом порядке - это не случайно.

Почему Constant Folding первым: свёртка констант может открыть тавтологии для Boolean Simplification (age > 0 + 0age > 0 → может упроститься дальше).

Почему Predicate Pushdown раньше Column Pruning: после pushdown фильтра Catalyst точно знает, какие колонки нужны для этого фильтра - Column Pruning учитывает это.

Почему несколько итераций: после применения Predicate Pushdown появляется новый Filter-узел ближе к Scan - возможно, теперь можно сделать ещё один Predicate Pushdown (например, если над этим Filter теперь стоит ещё один Filter).


Когда правила не срабатывают

Python UDF - непрозрачный контейнер

from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType

@udf(BooleanType())
def is_moscow(city):
    return city == "Moscow"

# Catalyst не может заглянуть внутрь UDF
df.filter(is_moscow(col("city")))
# → Predicate Pushdown НЕ работает
# → В explain(): фильтр стоит ПОСЛЕ FileScan, не до
# → Parquet Reader читает ВСЕ строки, затем Python UDF применяется

# Правильно: встроенная функция
df.filter(col("city") == "Moscow")
# → Predicate Pushdown работает
# → Parquet Reader пропускает неподходящие Row Groups

Сложные условия через expr()

# expr() парсится заново - Catalyst может не оптимизировать
df.filter(expr("city = 'Moscow' OR (age > 18 AND status = 'active')"))

# DataFrame API: Catalyst видит структуру
df.filter(
    (col("city") == "Moscow") |
    ((col("age") > 18) & (col("status") == "active"))
)

Источники без поддержки pushdown

Не все источники данных поддерживают Predicate Pushdown:

Источник Predicate Pushdown Column Pruning
Parquet Да (min/max + Bloom Filter) Да (column chunks)
ORC Да (min/max) Да
Delta Lake Да (Data Skipping) Да
Iceberg Да (partition pruning + stats) Да
CSV Нет Нет
JSON Нет Нет
JDBC Частично (WHERE clause в SQL) Нет
Kafka Нет Нет

Для CSV и JSON FileScan читает все строки и колонки, а фильтрация и проекция происходят на стороне Spark. Это главная причина, почему CSV плох для аналитических нагрузок.

Нарушение семантики - блокировщик pushdown

# Window функции: нельзя пушить фильтр через Window
from pyspark.sql.functions import rank
from pyspark.sql.window import Window

w = Window.partitionBy("city").orderBy(col("salary").desc())

# Фильтр после Window - нельзя переставить до Window
df.withColumn("rk", rank().over(w)) \
  .filter(col("rk") <= 3)  # Top 3 по зарплате в каждом городе

# Если бы Catalyst pushdown'ул фильтр до Window,
# семантика "топ 3 из всей группы" изменилась бы

Полный разбор explain(): где искать оптимизации

spark = SparkSession.builder \
    .master("local[*]") \
    .config("spark.sql.shuffle.partitions", "4") \
    .getOrCreate()

events = spark.read.parquet("events/")
users  = spark.read.parquet("users/")

query = (
    events
    .filter(col("event_date") >= "2024-01-01")
    .filter(col("country") == "RU")
    .join(users.select("user_id", "city"), "user_id")
    .groupBy("city")
    .agg({"revenue": "sum"})
)

query.explain(mode="extended")

Что искать в каждом плане:

== Parsed Logical Plan ==
# Апостроф у имён → Unresolved: 'event_date, 'country, 'user_id
# Порядок операторов точно как в коде - никакой оптимизации

== Analyzed Logical Plan ==
# Типы разрешены: event_date#5 (date), country#8 (string)
# Никакой оптимизации, только resolution

== Optimized Logical Plan ==
# ЗДЕСЬ ищем результаты RBO:

# 1. Constant Folding: "2024-01-01" уже как DateType литерал
# 2. isnotnull() добавлен автоматически:
#    Filter (isnotnull(event_date#5) AND event_date#5 >= 2024-01-01
#            AND isnotnull(country#8) AND country#8 = RU)
# 3. Predicate Pushdown: Filter стоит НИЖЕ Join, не выше
# 4. Column Pruning: events читает только нужные колонки,
#    users.select("user_id", "city") уже задан явно

== Physical Plan ==
# Exchange → Shuffle (Stage Boundary)
# *(N) → Whole-Stage CodeGen включён
# FileScan parquet [...] PushedFilters: [IsNotNull(event_date),
#   GreaterThanOrEqual(event_date,2024-01-01), IsNotNull(country),
#   EqualTo(country,RU)]
# ReadSchema: struct<user_id:bigint,event_date:date,country:string,revenue:double>

Практика: до и после оптимизации

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as spark_sum

spark = SparkSession.builder \
    .master("local[*]") \
    .config("spark.sql.shuffle.partitions", "4") \
    .getOrCreate()

# Широкая таблица (имитация 20 колонок)
wide_df = spark.range(10_000).toDF("id") \
    .withColumn("city",    (col("id") % 5).cast("string")) \
    .withColumn("country", (col("id") % 10).cast("string")) \
    .withColumn("age",     (col("id") % 60).cast("int")) \
    .withColumn("salary",  col("id") * 100) \
    .withColumn("c5",  col("id")) \
    .withColumn("c6",  col("id")) \
    .withColumn("c7",  col("id")) \
    .withColumn("c8",  col("id")) \
    .withColumn("c9",  col("id")) \
    .withColumn("c10", col("id"))

# ── Наивный запрос: SELECT * → Join → Filter → GroupBy ──────────────
naive = (
    wide_df                                     # все 11 колонок
    .filter(col("age") > 10 + 8)               # константа не свёрнута вручную
    .groupBy("city")
    .agg(spark_sum("salary").alias("total"))
)

print("=== NAIVE PLAN ===")
naive.explain(mode="formatted")

# ── Что сделал Catalyst? ─────────────────────────────────────────────
# 1. Constant Folding: 10 + 8 → 18
# 2. Column Pruning:  читает только id, city, age, salary (4 из 11)
# 3. Predicate Pushdown: Filter применяется до Aggregate

naive.show(5)

Best practices: optimizer-friendly код

Встроенные функции вместо UDF:

# Плохо: UDF блокирует все оптимизации
@udf("string")
def extract_domain(email):
    return email.split("@")[-1] if email else None

df.withColumn("domain", extract_domain(col("email")))

# Хорошо: встроенная функция - прозрачна для Catalyst
from pyspark.sql.functions import split, element_at
df.withColumn("domain", element_at(split(col("email"), "@"), -1))

Не используйте SELECT * без необходимости:

# Плохо: тянет все колонки через весь pipeline
df.join(dim, "id").select("*").filter(...)

# Хорошо: укажите колонки явно
df.join(dim, "id").select("df.id", "df.amount", "dim.category")

Используйте форматы с pushdown для аналитики:

# Плохо: CSV - нет ни Predicate Pushdown, ни Column Pruning
spark.read.csv("data.csv", header=True)

# Хорошо: Parquet - оба работают
spark.read.parquet("data.parquet/")

# Ещё лучше: партиционированный Parquet - Partition Pruning
spark.read.parquet("data/") \
    .filter(col("date") == "2024-01-15")
# date=2024-01-15/ читается, остальные директории пропускаются

Проверяйте PushedFilters в explain():

df = spark.read.parquet("events/").filter(col("city") == "Moscow")
df.explain(mode="formatted")
# Если PushedFilters: [] → фильтр НЕ попал в Parquet Reader
# Если PushedFilters: [EqualTo(city,Moscow)] → всё хорошо

Итог

Rule-Based Optimization - детерминированный слой оптимизации, который работает без статистики и применяется к любому запросу:

Правило Что делает Выгода
Constant Folding Вычисляет константные выражения при планировании Нет CPU на каждую строку
Boolean Simplification Упрощает AND/OR/NOT, убирает тавтологии Меньше условий, EmptyRelation для противоречий
Predicate Pushdown Опускает Filter ближе к источнику, до join Меньше строк в join, меньше данных читается
Column Pruning Убирает ненужные колонки из Scan Меньше данных читается из Parquet
Null Propagation Оптимизирует NULL-выражения Автоматический isnotnull, корректная семантика
Dead Code Elimination Убирает неиспользуемые вычисления Нет лишних withColumn / join

Главные враги RBO: Python UDF (непрозрачный контейнер), CSV/JSON (форматы без pushdown), OUTER JOIN (ограничения семантики), Window функции (порядок операций фиксирован).

Всегда проверяйте explain(mode="formatted") - если в PushedFilters пусто или ReadSchema содержит лишние колонки - что-то мешает оптимизатору.

Следующий шаг после RBO - Cost-Based Optimization: когда структуры дерева недостаточно и нужна статистика о реальных данных, чтобы выбрать лучший физический план.