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 код.
Предыдущий урок описал четыре фазы 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 → может быть упрощён или убран
Комбинация правил: как они усиливают друг друга¶
Правила работают в синергии. Один проход открывает возможности для другого:
Шаги:
- Constant Folding:
10 + 8→18. Теперь план выглядит чище - Predicate Pushdown:
Filter (age > 18)опускается нижеProjectи ближе кScan - Column Pruning: зная что нужны только
id,city(изProject) иage(дляFilter) -Scanобрезается до трёх колонок
Без FixedPoint-итерации шаг 3 мог бы не увидеть результат шагов 1 и 2.
Порядок правил имеет значение¶
Catalyst применяет правила в определённом порядке - это не случайно.
Почему Constant Folding первым: свёртка констант может открыть тавтологии для Boolean Simplification (age > 0 + 0 → age > 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: когда структуры дерева недостаточно и нужна статистика о реальных данных, чтобы выбрать лучший физический план.