Cost-Based Optimizer: ANALYZE TABLE и статистика колонок

Как Spark использует статистику таблиц и колонок для выбора оптимального физического плана: join strategy, join reordering, broadcast решения

core internals optimization

Ограниченность правил без данных

Предыдущий урок показал, как Rule-Based Optimizer (RBO) трансформирует план запроса, не зная ни строчки о реальных данных. Predicate Pushdown проталкивает фильтр как можно ниже по дереву - но не может знать, останется после этого фильтра 10 строк или 10 миллиардов. Column Pruning убирает лишние колонки - но не знает, сколько памяти реально займёт оставшийся набор.

Проблема проявляется в момент выбора стратегии соединения:

orders  = spark.table("orders")   # 10 млрд строк
returns = spark.table("returns")  # ???

df = orders.join(returns, "order_id")

Без знания размера таблицы returns Catalyst выбирает Sort Merge Join - самый безопасный, но тяжёлый вариант с полным shuffle обеих сторон. Если returns - это 5 МБ после фильтрации, оптимальный Broadcast Hash Join сэкономил бы минуты. Но RBO не знает этого. Cost-Based Optimizer знает - если ему предоставить статистику.

Архитектура CBO в Spark

Cost-Based Optimizer - слой поверх Logical Optimization. Он получает на вход дерево логического плана, которое RBO уже обработал (фильтры проброшены вниз, лишние колонки убраны), и должен принять ключевые решения о физическом исполнении.

Что решает CBO:

Решение Без CBO С CBO и статистикой
Стратегия join Sort Merge Join (безопасно) Broadcast Hash Join (быстро, если малая сторона мала)
Порядок join Слева направо по коду Сначала наименьший промежуточный результат
Shuffle partitions spark.sql.shuffle.partitions AQE корректирует динамически
Оценка строк Фиксированный коэффициент 0.3 Вычисляется из min/max/distinct/гистограммы

Настройки:

# Включить CBO (по умолчанию false в Spark < 3.x, true в 3.x+)
spark.conf.set("spark.sql.cbo.enabled", "true")

# Включить join reordering (по умолчанию false - агрессивная оптимизация)
spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")

# Максимальное количество таблиц для перестановки
spark.conf.set("spark.sql.cbo.joinReorder.dp.threshold", "12")

ANALYZE TABLE: сбор статистики

CBO бесполезен без данных о таблицах. Источник данных - команда ANALYZE TABLE, которая сканирует таблицу и сохраняет агрегаты в Metastore.

Уровни сбора

-- Уровень 1: размер таблицы + количество строк
ANALYZE TABLE orders COMPUTE STATISTICS;

-- Уровень 2: статистика по конкретным колонкам
ANALYZE TABLE orders COMPUTE STATISTICS
    FOR COLUMNS order_id, customer_id, amount, status, created_at;

-- Уровень 3: все колонки (дорого на широких таблицах)
ANALYZE TABLE orders COMPUTE STATISTICS FOR ALL COLUMNS;

-- Партиционированные таблицы - можно по партиции
ANALYZE TABLE orders PARTITION (year=2024, month=12)
    COMPUTE STATISTICS FOR COLUMNS amount, status;

Что собирается

Где хранится статистика

Статистика хранится в Hive Metastore - в TBLPROPERTIES таблицы и в отдельных таблицах метаданных (COLUMNS_V2, TABLE_PARAMS). При создании SparkSession Catalyst читает её через SessionCatalog.

# Просмотр table-level статистики
spark.sql("DESCRIBE EXTENDED orders").show(50, truncate=False)
# Строки: Statistics   сizeInBytes=1234567890, rowCount=1000000000

# Просмотр column-level статистики
spark.sql("DESCRIBE EXTENDED orders amount").show(truncate=False)

Для Delta Lake и Iceberg статистика хранится по-другому (подробнее в разделе ниже).

Гистограммы и оценка селективности

Проблема без гистограмм

Представим колонку city в таблице customers (100 млн строк):

  • Москва: 40 млн записей (40%)
  • Санкт-Петербург: 20 млн записей (20%)
  • Остальные города: по 100–1000 записей

Без гистограммы CBO знает только min='Вологда', max='Ярославль' и distinctCount=1000. На запрос WHERE city = 'Москва' Spark оценит: 100 млн / 1000 distinct = 100 000 строк - в 400 раз меньше реального!

Это приведёт к ошибочному выбору Broadcast Hash Join (думает, что таблица маленькая) там, где нужен Sort Merge Join.

Equi-Height Histograms

Spark строит гистограммы типа equi-height (равновысотные): каждая корзина содержит одинаковое количество строк, но разный диапазон значений.

С equi-height гистограммой CBO точно знает: корзина с Москвой - это 40 000 строк из 100 млн (40%). Оценка селективности становится точной.

Включение гистограмм

# Включить сбор гистограмм (по умолчанию false)
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")

# Количество корзин (по умолчанию 254)
spark.conf.set("spark.sql.statistics.histogram.numBins", "254")

# После включения - пересобрать статистику
spark.sql("ANALYZE TABLE customers COMPUTE STATISTICS FOR COLUMNS city, amount")

Формула оценки строк

Для предиката amount > 100 CBO вычисляет selectivity:

selectivity = (max - literal) / (max - min)
            = (10000 - 100) / (10000 - 0) = 0.99

estimated_rows = rowCount * selectivity
               = 1_000_000 * 0.99 = 990_000

Для city = 'Москва' с гистограммой:

selectivity = rows_in_bucket / total_rows
            = 40_000_000 / 100_000_000 = 0.40

estimated_rows = 100_000_000 * 0.40 = 40_000_000

Без гистограммы было бы 1 / distinctCount = 1/1000 = 0.001 → катастрофическая ошибка.

Как CBO меняет Physical Plan

Выбор стратегии Join

Это главная победа CBO. Catalyst поддерживает несколько стратегий join, и выбор между ними критически влияет на производительность:

Ключевая настройка:

# Порог для broadcast (10 МБ по умолчанию)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "10485760")

# Увеличить для больших таблиц, которые всё же влезают в память
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "104857600")  # 100 МБ

# Отключить broadcast (если executor OOM от слишком большого broadcast)
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

Сценарий: неправильный Join без статистики

sales   = spark.table("sales")    # 5 млрд строк, 400 ГБ
regions = spark.table("regions")  # 100 строк,  10 КБ

# Без статистики Spark не знает размер regions
result = sales.join(regions, "region_id").groupBy("region_name").sum("amount")

# Physical Plan без ANALYZE TABLE:
# *(4) HashAggregate
# +- Exchange hashpartitioning(region_name, 200)  ← shuffle!
#    +- *(3) HashAggregate
#       +- *(3) SortMergeJoin [region_id]          ← shuffle обеих сторон!
#          :- *(1) Sort [region_id]
#          :  +- Exchange hashpartitioning(region_id, 200)  ← shuffle 400 ГБ
#          +- *(2) Sort [region_id]
#             +- Exchange hashpartitioning(region_id, 200)  ← shuffle 10 КБ
ANALYZE TABLE regions COMPUTE STATISTICS;
-- Statistics: sizeInBytes=10240, rowCount=100
# Physical Plan с ANALYZE TABLE:
# *(2) HashAggregate
# +- Exchange hashpartitioning(region_name, 200)
#    +- *(1) HashAggregate
#       +- *(1) BroadcastHashJoin [region_id]      ← нет shuffle 400 ГБ!
#          :- FileScan sales
#          +- BroadcastExchange HashedRelationBroadcastMode
#             +- FileScan regions

Join Reordering

При соединении трёх и более таблиц порядок имеет критическое значение. Промежуточные результаты растут экспоненциально при неверном порядке.

spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")

# Запрос: три таблицы
result = (orders           # 1 млрд строк
    .join(customers, "customer_id")   # 100 млн строк
    .join(countries, "country_id"))   # 200 строк

В обоих случаях финальный результат одинаков, но в плохом варианте Spark дважды тащит 1 млрд строк через shuffle. С Join Reorder сначала объединяются маленькие таблицы, уменьшая промежуточный shuffle.

Алгоритм: Spark использует динамическое программирование (DP) для перебора порядков. Для N таблиц число вариантов - N!, но DP сокращает поиск до O(2^N × N). Порог dp.threshold=12 означает максимум 12 таблиц в одном join.

ANALYZE TABLE: детальная механика

Стоимость сбора статистики

ANALYZE TABLE - это полный scan таблицы. На petabyte-таблицах это может занять дольше, чем сам запрос. Стратегии:

# Только table-level (дешево: filesystem metadata + count)
spark.sql("ANALYZE TABLE orders COMPUTE STATISTICS")

# Только нужные колонки join-ключей и фильтруемых полей (разумный баланс)
spark.sql("""
    ANALYZE TABLE orders
    COMPUTE STATISTICS FOR COLUMNS order_id, customer_id, status, created_at
""")

# Избегать FOR ALL COLUMNS на таблицах с сотнями колонок

Партиционированные таблицы

-- Статистика всей таблицы (дорого на большой таблице)
ANALYZE TABLE orders COMPUTE STATISTICS;

-- Только свежие партиции (дешевле, точнее для новых данных)
ANALYZE TABLE orders PARTITION (year=2024, month=12)
    COMPUTE STATISTICS FOR COLUMNS order_id, amount, status;

-- Просмотр статистики партиции
DESCRIBE EXTENDED orders PARTITION (year=2024, month=12);

Выборочный анализ (Spark 3.0+)

# Вместо полного scan - выборка для оценки (быстро, менее точно)
spark.conf.set("spark.sql.statistics.sampling.enabled", "true")

# Размер выборки
spark.conf.set("spark.sql.statistics.sampleRatio", "0.01")  # 1%

AQE vs CBO: разные уровни оптимизации

CBO и Adaptive Query Execution (AQE) решают смежные задачи, но на разных этапах:

Что исправляет AQE, что CBO не может:

Проблема CBO AQE
Неверный join strategy Угадывает по статистике Переключает после shuffle stage
Skew в данных Не знает о skew Обнаруживает и разбивает skewed partitions
Слишком много shuffle partitions Фиксированное число Коалесцирует малые partitions
Устаревшая статистика Принимает как правду Использует реальные размеры shuffle

Внутренний цикл AQE

AQE работает не один раз за всё выполнение задания, а итерационно - на каждом Exchange (shuffle boundary):

Шаг 1: Выполняется текущий Stage (до Exchange)
Шаг 2: Shuffle Write + AQE собирает реальную статистику
        (размер каждой партиции, количество строк)
Шаг 3: AQE повторно оптимизирует план для оставшихся (ещё не запущенных) Stages
Шаг 4: Следующий Stage запускается с обновлённым планом
        → повторить с Шага 1 для следующего Exchange

Это означает, что чем больше у Job shuffle-границ, тем больше точек, в которых AQE может исправить план. Три конкретных оптимизации, которые AQE выполняет на каждой такой точке:

Оптимизация 1: Coalescing Shuffle Partitions

# До выполнения Spark планирует N партиций (spark.sql.shuffle.partitions = 200)
# После Stage 1 AQE видит реальные размеры:
# Partition 0: 2 MB, Partition 1: 3 MB, ..., Partition 190: 1 MB
# AQE объединяет малые партиции → вместо 200 → 12 укрупнённых партиций

spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "128mb")

Оптимизация 2: Switching Join Strategies

# CBO оценил правую таблицу в 500 MB → выбрал SortMergeJoin
# После Stage 1 AQE видит: после фильтрации осталось 8 MB
# → AQE переключает на BroadcastHashJoin (без дополнительного Exchange)

spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "30mb")

Оптимизация 3: Optimizing Skew Joins

# AQE обнаруживает: partition 42 содержит 80% всех строк (skew)
# → автоматически разбивает partition 42 на несколько подзадач
# → каждая подзадача дублирует соответствующую часть правой таблицы

spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "5")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "256mb")
# Включить AQE (по умолчанию true с Spark 3.2+)
spark.conf.set("spark.sql.adaptive.enabled", "true")

# AQE: конвертировать SortMerge в Broadcast, если размер оказался мал
spark.conf.set("spark.sql.adaptive.localShuffleReader.enabled", "true")

# AQE: порог для runtime broadcast
spark.conf.set("spark.sql.adaptive.autoBroadcastJoinThreshold", "30MB")

Рекомендация: CBO и AQE дополняют друг друга. CBO принимает правильные первоначальные решения на основе статистики, AQE исправляет те решения, где статистика оказалась неверной.

Dynamic Partition Pruning (DPP)

Dynamic Partition Pruning - оптимизация звёздных схем (star schema), при которой Spark динамически инжектирует фильтр измерения в скан таблицы фактов во время выполнения - без изменения пользовательского запроса.

Как работает DPP

Классический запрос к звёздной схеме:

SELECT f.amount, d.region
FROM fact_sales f
JOIN dim_dates d ON f.date_id = d.date_id
WHERE d.year = 2024 AND d.quarter = 'Q4'

Без DPP: Spark сначала сканирует всю fact_sales (допустим, 10 ТБ), потом выполняет join с dim_dates.

С DPP:

Шаг 1: Spark сканирует dim_dates → фильтрует WHERE year=2024 AND quarter='Q4'
        → получает список date_id: {20241001, 20241002, ..., 20241231}
Шаг 2: Spark инжектирует этот список как subquery в скан fact_sales:
        WHERE date_id IN (subquery результатов из dim_dates)
Шаг 3: Скан fact_sales читает только партиции Q4 2024
        → вместо 10 ТБ → 500 ГБ (только нужный квартал)

Условия для работы DPP

DPP срабатывает автоматически, если выполнены все условия:

  • Таблица фактов партиционирована по колонке join-ключа (PARTITIONED BY (date_id))
  • Таблица измерения достаточно мала для broadcast (≤ autoBroadcastJoinThreshold) или уже broadcastится явно через broadcast()
  • В запросе есть фильтр именно по таблице измерения (не по таблице фактов)
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")  # true по умолчанию в Spark 3.x

# Проверка в explain: ищите "dynamicpruningexpression" или "PartitionFilters"
fact_sales.join(broadcast(dim_dates.filter("year = 2024 AND quarter = 'Q4'")), "date_id") \
    .explain("formatted")
# PartitionFilters: [isnotnull(date_id#0), dynamicpruningexpression(date_id#0 IN ...)]

Практический эффект

Из TPC-DS бенчмарка (1 ТБ, запрос q25):

Без DPP С DPP
Время выполнения 390 сек 20 сек
Прочитано данных 1 ТБ (вся fact) ~50 ГБ (нужные партиции)
Ускорение - 19×

DPP наиболее эффективен при высокой selectivity фильтра измерения (отбирается < 20% данных из fact).

Статистика в разных форматах

Hive Metastore (Hive таблицы)

Стандартный путь: ANALYZE TABLETBLPROPERTIESSessionCatalog.getTableMetadata().

Parquet и ORC

Parquet хранит собственную статистику в footer каждого файла и каждого row group:

  • num_rows - количество строк в row group
  • min/max - для каждой колонки в row group

CBO использует это автоматически для file-level pruning, но для join decisions нужна именно ANALYZE TABLE уровня таблицы.

# Parquet-уровень статистики виден в explain
df.filter("amount > 1000").explain("formatted")
# PushedFilters: [IsNotNull(amount), GreaterThan(amount,1000.0)]
# PartitionFilters: []
# DataFilters: [isnotnull(amount#5), (amount#5 > 1000.0)]

Delta Lake

Delta хранит статистику в transaction log (_delta_log/*.json). Каждый commit содержит stats для добавленных файлов:

{
  "numRecords": 1000000,
  "minValues": {"amount": 0.01, "status": "cancelled"},
  "maxValues": {"amount": 99999.99, "status": "shipped"},
  "nullCount": {"amount": 0, "status": 1234}
}

Delta автоматически собирает эту статистику при записи (по первым 32 колонкам по умолчанию). ANALYZE TABLE на Delta обновляет Hive Metastore для совместимости.

# Настройка количества колонок для автоматической статистики в Delta
spark.conf.set("spark.databricks.delta.stats.collect.maxColumnsToCollect", "64")

# Явное обновление статистики Delta
spark.sql("ANALYZE TABLE delta.`/path/to/table` COMPUTE STATISTICS FOR COLUMNS amount, status")

Apache Iceberg

Iceberg хранит статистику в manifest files. При каждой операции write автоматически записываются lower_bound и upper_bound для всех колонок.

-- Iceberg: сбор подробной статистики
ANALYZE TABLE iceberg_catalog.db.orders COMPUTE STATISTICS FOR ALL COLUMNS;

-- Просмотр статистики файлов
SELECT file_path, record_count, file_size_in_bytes
FROM iceberg_catalog.db.orders.files
ORDER BY record_count DESC;

Особенность Iceberg: partition statistics из manifest files используются для partition pruning автоматически, без ANALYZE TABLE.

Устаревшая статистика: главная ловушка

Что происходит при statistics drift

Статистика становится неактуальной после любой операции записи в таблицу. Spark не знает об этом - он использует кэшированные значения из Metastore.

Диагностика устаревшей статистики

# Проверить, когда последний раз собиралась статистика
spark.sql("DESCRIBE EXTENDED my_table").filter("col_name = 'Statistics'").show()
# Statistics   sizeInBytes=107374182400, rowCount=1000000000

# Реальное количество строк
spark.table("my_table").count()  # 1_500_000_000 - расхождение!

# В explain() смотреть на estimatedSize и rowCount
spark.table("my_table").join(...).explain("cost")
# Statistics(sizeInBytes=100.0 GiB, rowCount=1.0E9)  ← устаревшие значения

Когда обновлять статистику

Событие Действие
Начальная загрузка таблицы ANALYZE TABLE ... FOR COLUMNS join_keys, filter_cols
Ежедневный incremental load (+5-10%) Ежедневное обновление нецелесообразно
Значительное изменение данных (+20-30%) ANALYZE TABLE на затронутых партициях
Смена профиля данных (новые города, категории) ANALYZE TABLE ... FOR COLUMNS с гистограммами
После TRUNCATE TABLE или полной перезаписи Немедленное ANALYZE TABLE
После DELETE/UPDATE/MERGE (Delta/Iceberg) ANALYZE TABLE на изменённых партициях

Стратегия обновления в production

def refresh_statistics(spark, table_name, key_columns, partition_spec=None):
    if partition_spec:
        # Только свежие партиции (быстро)
        sql = f"""
            ANALYZE TABLE {table_name}
            PARTITION ({partition_spec})
            COMPUTE STATISTICS FOR COLUMNS {', '.join(key_columns)}
        """
    else:
        # Вся таблица (медленно, но точно)
        sql = f"""
            ANALYZE TABLE {table_name}
            COMPUTE STATISTICS FOR COLUMNS {', '.join(key_columns)}
        """
    spark.sql(sql)

# После ETL-джоба обновляем только текущую партицию
refresh_statistics(
    spark,
    table_name="orders",
    key_columns=["order_id", "customer_id", "status", "amount"],
    partition_spec="year=2024, month=12"
)

Практика: диагностика и оптимизация через CBO

Шаг 1: Создание тестового сценария

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, rand, round as spark_round

spark = SparkSession.builder \
    .appName("CBO-Demo") \
    .config("spark.sql.cbo.enabled", "true") \
    .config("spark.sql.cbo.joinReorder.enabled", "true") \
    .config("spark.sql.statistics.histogram.enabled", "true") \
    .getOrCreate()

# Создаём три таблицы с разными размерами
import random
from pyspark.sql.types import *

# Большая таблица: 10 млн заказов
orders_data = [(i, random.randint(1, 100000), random.randint(1, 50),
                round(random.uniform(10, 10000), 2))
               for i in range(10_000_000)]
orders = spark.createDataFrame(orders_data, ["order_id", "customer_id", "country_id", "amount"])
orders.write.mode("overwrite").saveAsTable("orders")

# Средняя таблица: 100 000 клиентов
customers_data = [(i, f"customer_{i}", random.randint(1, 50))
                  for i in range(1, 100_001)]
customers = spark.createDataFrame(customers_data, ["customer_id", "name", "country_id"])
customers.write.mode("overwrite").saveAsTable("customers")

# Маленькая таблица: 50 стран
countries_data = [(i, f"country_{i}") for i in range(1, 51)]
countries = spark.createDataFrame(countries_data, ["country_id", "country_name"])
countries.write.mode("overwrite").saveAsTable("countries")

Шаг 2: Запрос без статистики

result = (spark.table("orders")
    .join(spark.table("customers"), "customer_id")
    .join(spark.table("countries"), "country_id")
    .groupBy("country_name")
    .sum("amount"))

# Смотрим план: CBO не знает размеров → Sort Merge Join
result.explain("cost")

Ожидаемый вывод (без статистики):

== Optimized Logical Plan ==
Aggregate [country_name], [country_name, sum(amount) AS sum(amount)]
+- Join Inner, (country_id = country_id)
   :- Statistics(sizeInBytes=8.0 EiB)   ← ∞ без статистики!
   +- Statistics(sizeInBytes=8.0 EiB)

== Physical Plan ==
SortMergeJoin [customer_id], [customer_id], Inner   ← тяжелый join

Шаг 3: Сбор статистики

ANALYZE TABLE orders COMPUTE STATISTICS FOR COLUMNS order_id, customer_id, country_id, amount;
ANALYZE TABLE customers COMPUTE STATISTICS FOR COLUMNS customer_id, country_id;
ANALYZE TABLE countries COMPUTE STATISTICS FOR ALL COLUMNS;

Шаг 4: Запрос со статистикой

result = (spark.table("orders")
    .join(spark.table("customers"), "customer_id")
    .join(spark.table("countries"), "country_id")
    .groupBy("country_name")
    .sum("amount"))

result.explain("cost")

Ожидаемый вывод (со статистикой):

== Optimized Logical Plan ==
Aggregate [country_name], [country_name, sum(amount)]
+- Join Inner, (customer_id = customer_id)
   :- Statistics(sizeInBytes=152.6 MiB, rowCount=10000000)  ← точно!
   +- Join Inner, (country_id = country_id)                  ← reorder!
      :- Statistics(sizeInBytes=2.6 MiB, rowCount=100000)
      +- Statistics(sizeInBytes=1.0 KiB, rowCount=50)        ← countries в начале

== Physical Plan ==
*(4) HashAggregate(keys=[country_name], functions=[sum(amount)])
+- *(4) BroadcastHashJoin [customer_id], [customer_id], Inner  ← broadcast!
   :- FileScan orders
   +- BroadcastExchange
      +- *(2) BroadcastHashJoin [country_id], [country_id]     ← второй broadcast!
         :- FileScan customers
         +- BroadcastExchange
            +- FileScan countries

Результат: вместо двух Sort Merge Join с полным shuffle 10 млн строк - два Broadcast Hash Join. Ускорение в 5–20x в зависимости от размера кластера.

Шаг 5: Проверка через df.explain("formatted")

result.explain("formatted")
# Ищем в выводе:
# 1. Join Strategy: BroadcastHashJoin vs SortMergeJoin
# 2. Statistics в Optimized Logical Plan: sizeInBytes, rowCount
# 3. BroadcastExchange: подтверждение что таблица будет broadcast'иться
# 4. Порядок Join в дереве плана (CBO reordering)

Ошибки CBO и их диагностика

CBO выбрал неверный Join (OOM)

ExecutorLostFailure: Exit code 137 (OOM)

Причина: CBO использовал устаревшую статистику и решил сделать Broadcast Join таблицы, которая уже выросла до 50 ГБ.

# Диагностика
spark.sql("DESCRIBE EXTENDED orders").filter(col("col_name") == "Statistics").show()
# Показывает устаревшие sizeInBytes

# Быстрое исправление без пересбора статистики
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")  # отключить broadcast

# Правильное исправление
spark.sql("ANALYZE TABLE orders COMPUTE STATISTICS")

CBO не применяет Join Reorder

# Проверить, что CBO включен
spark.conf.get("spark.sql.cbo.enabled")       # должно быть "true"
spark.conf.get("spark.sql.cbo.joinReorder.enabled")  # должно быть "true"

# Проверить, что статистика есть у всех таблиц в join
# Без статистики хотя бы одной таблицы - reorder не работает
spark.sql("DESCRIBE EXTENDED table_a").filter(col("col_name") == "Statistics").show()
spark.sql("DESCRIBE EXTENDED table_b").filter(col("col_name") == "Statistics").show()

Статистика собирается слишком долго

# Вместо полного scan - выборочная статистика
spark.conf.set("spark.sql.statistics.sampling.enabled", "true")
spark.conf.set("spark.sql.statistics.sampleRatio", "0.05")  # 5% выборка

# Или: разбить ANALYZE на партиции и запускать параллельно
from concurrent.futures import ThreadPoolExecutor

partitions = [("2024", "10"), ("2024", "11"), ("2024", "12")]

def analyze_partition(year, month):
    spark.sql(f"""
        ANALYZE TABLE orders PARTITION (year={year}, month={month})
        COMPUTE STATISTICS FOR COLUMNS order_id, customer_id, amount
    """)

with ThreadPoolExecutor(max_workers=3) as executor:
    executor.map(lambda p: analyze_partition(*p), partitions)

Best Practices

1. Собирать статистику только для нужных колонок:

-- Да: join-ключи, часто фильтруемые колонки
ANALYZE TABLE orders COMPUTE STATISTICS
    FOR COLUMNS order_id, customer_id, status, created_at;

-- Нет: все 200 колонок в wide таблице
-- ANALYZE TABLE orders COMPUTE STATISTICS FOR ALL COLUMNS;

2. Включать гистограммы для skewed колонок:

# Если в колонке неравномерное распределение (один город = 40%)
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")
spark.sql("ANALYZE TABLE customers COMPUTE STATISTICS FOR COLUMNS city")

3. Встроить обновление статистики в ETL-пайплайн:

def run_etl(spark, date):
    # 1. Загрузка данных
    load_data(spark, date)

    # 2. Обновить статистику только для текущей партиции
    year, month = date.split("-")[:2]
    spark.sql(f"""
        ANALYZE TABLE orders PARTITION (year={year}, month={month})
        COMPUTE STATISTICS FOR COLUMNS order_id, customer_id, amount, status
    """)
    # Общую статистику таблицы (sizeInBytes, rowCount) - раз в неделю/месяц

4. Комбинировать CBO и AQE:

# Оба метода вместе дают максимальный эффект
spark.conf.set("spark.sql.cbo.enabled", "true")
spark.conf.set("spark.sql.cbo.joinReorder.enabled", "true")
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")

# CBO: правильный план с самого начала на основе статистики
# AQE: исправляет ошибки CBO во время выполнения

5. Мониторинг через explain("cost"):

# Перед запуском тяжёлого job - проверить план
df.explain("cost")

# В плане искать:
# Statistics(sizeInBytes=8.0 EiB) → нет статистики, CBO не работает
# Statistics(sizeInBytes=100 MiB, rowCount=1000000) → есть статистика
# BroadcastHashJoin → CBO выбрал broadcast
# SortMergeJoin → CBO выбрал тяжёлый join (возможно, правильно)

Итого

Аспект Детали
Что решает CBO Join strategy, Join Reorder, broadcast thresholds
Источник данных ANALYZE TABLE → Hive Metastore
Table-level stat sizeInBytes, rowCount
Column-level stat min, max, distinctCount, nullCount, avgLen, гистограмма
Главная ловушка Устаревшая статистика → ошибочный OOM или лишний shuffle
Взаимодействие с AQE CBO - до выполнения, AQE - во время выполнения
Delta Lake Авто-статистика при write, ANALYZE TABLE для Metastore
Apache Iceberg Статистика в manifest files, автоматический partition pruning
Включение spark.sql.cbo.enabled=true, spark.sql.cbo.joinReorder.enabled=true

Следующий урок рассматривает Tungsten - подсистему физического исполнения Spark, которая работает под CBO и RBO: управление памятью в off-heap, бинарный формат UnsafeRow и Whole-Stage Code Generation.