Cost-Based Optimizer: ANALYZE TABLE и статистика колонок
Как Spark использует статистику таблиц и колонок для выбора оптимального физического плана: join strategy, join reordering, broadcast решения
Ограниченность правил без данных¶
Предыдущий урок показал, как 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 TABLE → TBLPROPERTIES → SessionCatalog.getTableMetadata().
Parquet и ORC¶
Parquet хранит собственную статистику в footer каждого файла и каждого row group:
num_rows- количество строк в row groupmin/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.