Runtime Filtering: AQE Bloom Filter push к сканированию и dynamic filtering в Star-Schema joins
Полный разбор Runtime Filtering в Spark 3.x: проблема полного сканирования фактовой таблицы, механизм Dynamic Filtering через AQE, Bloom Filter как компактный фильтр ключей, pushdown к File Scan, Data Skipping в Parquet/Delta/Iceberg, конфигурационные параметры, чтение плана выполнения и Best Practices для Lakehouse Star-Schema архитектур.
1. Проблема классического Join в Star-Schema¶
Архитектура «Звезда» (Star Schema) — стандарт Data Warehouse и современных Lakehouse. В центре находится гигантская таблица фактов (Fact Table), а вокруг неё — небольшие таблицы измерений (Dimension Tables). Запросы к такой схеме выглядят просто, но скрывают фундаментальную проблему производительности.
Почему стандартный подход читает лишние данные¶
Рассмотрим типичный аналитический запрос: «Найти все продажи европейских магазинов за последний квартал категории Электроника».
SELECT
s.sale_date,
s.amount,
st.store_name,
p.product_name
FROM fact_sales s
JOIN dim_stores st ON s.store_id = st.id
JOIN dim_products p ON s.product_id = p.id
WHERE
st.region = 'Europe'
AND p.category = 'Electronics'
AND s.sale_date >= '2024-10-01'
Без Runtime Filtering классический Spark выполнит это так:
Схема показывает катастрофу: кластер читает 1 PB данных из fact_sales, гоняет их по сети в Shuffle, а затем на этапе JOIN отбрасывает 95% прочитанного. Вся эта работа была напрасной — мы заранее знали что нужны только европейские магазины и Electronics.
Масштаб проблемы¶
Для понимания реального эффекта:
fact_sales: 1 PB, 10 миллиардов строкdim_storesпосле фильтра: 500 строк (магазины Европы)dim_productsпосле фильтра: 1,200 строк (Electronics)- Строки фактов, участвующих в JOIN: ~50,000 (0.0005% от 10B!)
- 99.9995% прочитанных данных — мусор
Если бы мы знали заранее какие store_id нужны, можно было бы применить предикат к fact_sales ДО чтения с диска. Именно это делает Runtime Filtering.
2. Концепция Dynamic Filtering: фильтр из одной таблицы в другую¶
Runtime Filtering — оптимизация, при которой Spark динамически генерирует фильтрующие условия на основе результатов выполненных Stage'ов и применяет их к ещё не запущенным Stage'ам.
Физика процесса¶
Ключевые изменения по сравнению с классическим подходом:
- Измерения обрабатываются до сканирования фактов
- Bloom Filter строится из ключей измерений (несколько килобайт)
- Факты сканируются с применением фильтра — читается 5 GB вместо 1 PB
- Shuffle уменьшается с 1 PB до 5 GB
3. Механизм AQE: связующее звено для динамической фильтрации¶
Adaptive Query Execution — это именно тот механизм, который делает Runtime Filtering возможным. Без AQE план строится статически до выполнения и не может «на лету» внедрять фильтры.
Как AQE инициирует Runtime Filtering¶
Диаграмма показывает ключевой момент: AQE получает реальную статистику из Stage 1 (500 строк вместо неизвестного) и модифицирует план Stage 2, внедряя Bloom Filter прямо в операцию Scan.
Без AQE: почему статический оптимизатор не справляется¶
# Без AQE: план строится ДО выполнения
# Catalyst не знает сколько строк останется после фильтрации dim_stores
# Поэтому:
# 1. dim_stores после WHERE region='Europe': Catalyst оценивает ~5000 строк
# (реально: 500 строк — ошибка в 10x!)
# 2. Catalyst решает: 5000 строк > broadcast threshold → SortMergeJoin
# 3. Нет Bloom Filter → нет Runtime Filtering → полный Scan фактов
# С AQE: план может быть изменён после Stage 1
# AQE видит реальные 500 строк → строит Bloom Filter → применяет к Scan Stage 2
4. Bloom Filter: компактный фильтр ключей¶
Когда фильтрующих ключей больше нескольких тысяч, передавать их как обычный список WHERE store_id IN (101, 205, 308, ...) становится неэффективно. Именно здесь Bloom Filter показывает свои преимущества.
Почему Bloom Filter идеально подходит для Runtime Filtering¶
import math
def calculate_bloom_filter_for_runtime_filtering(
n_keys: int, # число уникальных ключей из таблицы измерений
fpp: float = 0.03, # acceptable false positive rate
) -> dict:
"""
Показывает почему Bloom Filter лучше явного списка IN для Runtime Filtering.
"""
# Размер явного IN-листа (хранить как int64)
explicit_list_bytes = n_keys * 8 # 8 байт на int64
explicit_list_mb = explicit_list_bytes / 1024 / 1024
# Размер Bloom Filter для того же FPP
bits_needed = int(-n_keys * math.log(fpp) / (math.log(2) ** 2))
bloom_bytes = bits_needed // 8
bloom_mb = bloom_bytes / 1024 / 1024
compression_ratio = explicit_list_bytes / bloom_bytes
return {
"n_keys": n_keys,
"explicit_list_mb": f"{explicit_list_mb:.1f} MB",
"bloom_filter_mb": f"{bloom_mb:.1f} MB",
"compression_ratio": f"{compression_ratio:.1f}x компактнее",
"fpp": f"{fpp*100:.1f}%",
"verdict": (
"BF выгоден" if bloom_mb < explicit_list_mb * 0.3
else "Разница небольшая"
)
}
# Примеры из реального Star-Schema:
cases = [
(500, "500 европейских магазинов"),
(50_000, "50K активных клиентов"),
(1_000_000,"1M ID транзакций за месяц"),
]
print("Сравнение IN-список vs Bloom Filter:")
for n, desc in cases:
r = calculate_bloom_filter_for_runtime_filtering(n)
print(f"\n{desc}:")
print(f" IN-list: {r['explicit_list_mb']}")
print(f" Bloom Filter: {r['bloom_filter_mb']} ({r['fpp']} FPP)")
print(f" Компактнее: {r['compression_ratio']}")
# Вывод:
# 500 европейских магазинов:
# IN-list: 0.0 MB
# Bloom Filter: 0.0 MB ← оба маленькие, разница незначительна
#
# 50K активных клиентов:
# IN-list: 0.4 MB
# Bloom Filter: 0.1 MB ← в 4x компактнее
#
# 1M ID транзакций:
# IN-list: 7.6 MB
# Bloom Filter: 1.9 MB ← в 4x компактнее
Bloom Filter и Data Skipping в Parquet¶
Синергия Bloom Filter + Parquet Row Groups даёт дополнительный эффект:
Важно: Bloom Filter применяется на двух уровнях:
- Row Group уровень (Parquet): если ни одно значение в Row Group не может пройти фильтр по min/max статистике — весь блок пропускается без чтения
- Строчный уровень: для строк внутри прочитанных Row Groups каждая строка проверяется через Bloom Filter
5. Конфигурация: включение и настройка Runtime Filtering¶
Основные параметры¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
# ── Главные переключатели ────────────────────────────────────────
# AQE должен быть включён — без него Runtime Filtering не работает!
.config("spark.sql.adaptive.enabled", "true") \
# Runtime Bloom Filter Join (Spark 3.3+, по умолчанию true в 3.3+)
.config("spark.sql.optimizer.runtime.bloomFilter.enabled", "true") \
# ── Параметры Bloom Filter ────────────────────────────────────────
# Порог: строить BF только если build-сторона > этого значения строк.
# Если build-сторона маленькая → автоматически BroadcastHashJoin (лучше)
# Дефолт: 10,000,000 строк. Уменьшите если хотите BF для меньших таблиц.
.config("spark.sql.optimizer.runtime.bloomFilter.creationSideThreshold",
"10000000") \
# Порог размера для создания BF (в байтах):
# Если build-сторона меньше этого → BroadcastHashJoin без BF
# Дефолт: 10 MB
.config("spark.sql.optimizer.runtime.bloomFilter.creationSideThreshold",
"10485760") \
# Максимальный размер Bloom Filter (бит):
# Дефолт: 67108864 (= 64 Mbit = 8 MB максимум на один BF)
.config("spark.sql.optimizer.runtime.bloomFilter.maxNumBits",
"67108864") \
# Ожидаемый FPP (False Positive Probability):
# Дефолт: 0.03 (3%). Уменьшите для более точной фильтрации (больший BF)
# Увеличьте для экономии памяти (менее точная фильтрация)
# ── Dynamic Partition Pruning (DPP) ───────────────────────────────
# Дополнительная оптимизация: прунинг партиций на основе JOIN результата
# Применяется когда fact_table партиционирована по join_key
.config("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") \
.config("spark.sql.optimizer.dynamicPartitionPruning.useStats", "true") \
.getOrCreate()
Полная конфигурация для Star Schema Production¶
spark = SparkSession.builder \
.appName("star-schema-optimized") \
# ── AQE: основа всей динамической оптимизации ──────────────────
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.adaptive.coalescePartitions.enabled", "true") \
.config("spark.sql.adaptive.skewJoin.enabled", "true") \
.config("spark.sql.adaptive.autoBroadcastJoinThreshold", "50MB") \
# ── Runtime Bloom Filter ────────────────────────────────────────
.config("spark.sql.optimizer.runtime.bloomFilter.enabled", "true") \
# ── Dynamic Partition Pruning ───────────────────────────────────
.config("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") \
# ── CBO для лучшей оценки размеров таблиц ──────────────────────
.config("spark.sql.cbo.enabled", "true") \
.config("spark.sql.statistics.histogram.enabled", "true") \
# ── Shuffle оптимизация ─────────────────────────────────────────
# Много начальных партиций + AQE схлопнёт пустые
.config("spark.sql.shuffle.partitions", "400") \
.getOrCreate()
6. Синергия с форматами хранения¶
Runtime Filtering даёт максимальный эффект при интеграции с современными форматами хранения, которые хранят статистику на уровне файловых блоков.
Delta Lake: Z-Ordering + Runtime Filtering¶
# Создаём оптимизированную Delta Lake таблицу фактов
spark.sql("""
CREATE TABLE fact_sales (
sale_id BIGINT,
store_id INT,
product_id INT,
customer_id BIGINT,
sale_date DATE,
amount DECIMAL(18,2)
)
USING DELTA
PARTITIONED BY (sale_date) -- Партиционирование по дате
""")
# Z-ORDER по часто используемым JOIN-ключам
# ПОСЛЕ партиции по дате, данные внутри партиции сортируются
# так чтобы похожие store_id и product_id были рядом
spark.sql("""
OPTIMIZE fact_sales
ZORDER BY (store_id, product_id)
""")
# Эффект для Runtime Filtering:
# 1. BF probe находит релевантные store_ids
# 2. Z-Order гарантирует что строки с store_id=101..500
# сконцентрированы в ограниченном числе Row Groups
# 3. Вместо читать 95% файла → читаем только 5%!
Apache Iceberg: встроенные Bloom Filters в metadata¶
spark = SparkSession.builder \
.config("spark.sql.extensions",
"org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
.config("spark.sql.catalog.iceberg",
"org.apache.iceberg.spark.SparkCatalog") \
.getOrCreate()
# Iceberg поддерживает Bloom Filters в метаданных файлов
spark.sql("""
CREATE TABLE iceberg.warehouse.fact_sales (
sale_id BIGINT,
store_id INT,
product_id INT,
sale_date DATE,
amount DOUBLE
)
USING ICEBERG
PARTITIONED BY (days(sale_date))
TBLPROPERTIES (
'write.parquet.bloom-filter-enabled.column.store_id' = 'true',
'write.parquet.bloom-filter-enabled.column.product_id' = 'true',
'write.parquet.bloom-filter-fpp.column.store_id' = '0.02',
'write.parquet.bloom-filter-fpp.column.product_id' = '0.02'
)
""")
Parquet: Row Group Statistics и BF¶
# Для максимального эффекта Data Skipping:
# 1. Записывайте данные отсортированными по join-ключу
# 2. Это обеспечивает что каждый Row Group содержит компактный диапазон store_id
from pyspark.sql import functions as F
fact_df = spark.table("bronze.fact_sales")
# Сортировка при записи = компактные Row Groups
fact_df \
.repartition(200, F.col("store_id")) \ # хэш-партиционирование по store_id
.sortWithinPartitions("store_id", "product_id") \ # сортировка внутри партиции
.write \
.format("parquet") \
.partitionBy("sale_date") \ # файловое партиционирование по дате
.option("parquet.block.size", "134217728") \ # 128 MB Row Group
.option("parquet.enable.dictionary", "true") \ # словарная компрессия для store_id
.saveAsTable("silver.fact_sales_optimized")
7. Пошаговый разбор: Star Schema Join с Runtime Filtering¶
Проследим полную последовательность выполнения оптимизированного запроса.
Демонстрация с реальным запросом¶
from pyspark.sql import SparkSession, functions as F
import time
spark = SparkSession.builder \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.optimizer.runtime.bloomFilter.enabled", "true") \
.config("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") \
.getOrCreate()
# Создаём тестовые данные, имитирующие реальную Star Schema
def create_star_schema_data(spark, fact_rows=10_000_000):
"""Создаёт Star Schema с реалистичными пропорциями."""
from pyspark.sql import functions as F
# dim_stores: 10,000 магазинов, 15% в Европе
stores = spark.range(10_000).select(
F.col("id").alias("store_id"),
F.when(F.col("id") < 1500, F.lit("Europe"))
.when(F.col("id") < 4000, F.lit("Americas"))
.when(F.col("id") < 7000, F.lit("Asia"))
.otherwise(F.lit("Other")).alias("region"),
F.concat(F.lit("Store_"), F.col("id")).alias("store_name"),
)
# dim_products: 50,000 SKU, 3% Electronics
products = spark.range(50_000).select(
F.col("id").alias("product_id"),
F.when(F.col("id") < 1500, F.lit("Electronics"))
.when(F.col("id") < 20000, F.lit("Clothing"))
.otherwise(F.lit("Other")).alias("category"),
F.concat(F.lit("Product_"), F.col("id")).alias("product_name"),
)
# fact_sales: 10M строк, равномерно распределённые
facts = spark.range(fact_rows).select(
F.col("id").alias("sale_id"),
(F.rand() * 10_000).cast("int").alias("store_id"),
(F.rand() * 50_000).cast("int").alias("product_id"),
(F.rand() * 1000).alias("amount"),
F.date_add(F.lit("2024-01-01"), (F.rand() * 365).cast("int")).alias("sale_date"),
)
return facts, stores, products
facts, stores, products = create_star_schema_data(spark)
facts.cache()
stores.cache()
products.cache()
facts.count()
stores.count()
products.count()
def run_star_schema_query(facts, stores, products,
label: str, with_aqe: bool) -> float:
"""Запускает Star Schema JOIN и измеряет время."""
spark.conf.set("spark.sql.adaptive.enabled",
"true" if with_aqe else "false")
spark.conf.set("spark.sql.optimizer.runtime.bloomFilter.enabled",
"true" if with_aqe else "false")
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled",
"true" if with_aqe else "false")
# Star Schema запрос
query = (
facts
.join(stores.filter(F.col("region") == "Europe"), "store_id")
.join(products.filter(F.col("category") == "Electronics"), "product_id")
.groupBy("store_name", "product_name")
.agg(F.sum("amount").alias("total_revenue"))
)
t0 = time.time()
result = query.count()
elapsed = time.time() - t0
print(f"{label}: {elapsed:.1f}с, результат: {result:,} строк")
return elapsed
print("Сравнение с/без Runtime Filtering:")
t_without = run_star_schema_query(facts, stores, products,
"Без AQE/Runtime Filtering", with_aqe=False)
t_with = run_star_schema_query(facts, stores, products,
"С AQE + Runtime Filtering", with_aqe=True)
print(f"\nУскорение: {t_without/t_with:.1f}x")
8. Анализ планов выполнения: читаем Runtime Filtering в EXPLAIN¶
Как найти Runtime Filter в Physical Plan¶
facts, stores, products = create_star_schema_data(spark, fact_rows=1_000_000)
# Строим запрос
query = (
facts
.join(stores.filter(F.col("region") == "Europe"), "store_id")
.join(products.filter(F.col("category") == "Electronics"), "product_id")
.groupBy("store_name")
.agg(F.sum("amount"))
)
# Смотрим полный план
spark.conf.set("spark.sql.adaptive.enabled", "true")
query.explain("formatted")
Что искать в выводе EXPLAIN:
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
*HashAggregate(keys=[store_name#12], functions=[sum(amount#5)])
+- Exchange hashpartitioning(store_name#12, 200)
+- *HashAggregate(keys=[store_name#12], functions=[partial_sum(amount#5)])
+- *Project [...]
+- SortMergeJoin [store_id#2], [store_id#15], Inner
:- *Filter might_contain(bloomfilter, store_id#2) ← ВОТ ОН!
: +- *FileScan parquet fact_sales[...]
+- BroadcastExchange ...
+- *Filter (region#18 = Europe)
+- *FileScan parquet dim_stores[...]
Ключевые узлы для поиска:
def check_runtime_filtering_active(df) -> dict:
"""
Проверяет активирован ли Runtime Filtering в плане запроса.
Ищет следующие маркеры:
1. 'might_contain' или 'bloomfilter' → Bloom Filter применён к Scan
2. 'dynamicpruningexpression' → Dynamic Partition Pruning
3. 'RuntimeFilter' → явный узел RuntimeFilter
4. 'AdaptiveSparkPlan' с 'isFinalPlan=true' → AQE перестроил план
"""
# Получаем текстовое представление плана
plan_str = df._jdf.queryExecution().executedPlan().toString()
markers = {
"bloom_filter_active": "might_contain" in plan_str or "bloomfilter" in plan_str.lower(),
"dynamic_pruning": "dynamicpruningexpression" in plan_str.lower(),
"runtime_filter_node": "RuntimeFilter" in plan_str,
"aqe_active": "AdaptiveSparkPlan" in plan_str,
"is_final_plan": "isFinalPlan=true" in plan_str,
}
total_optimizations = sum(1 for v in markers.values() if v)
markers["optimization_count"] = total_optimizations
markers["well_optimized"] = total_optimizations >= 2
return markers
# До выполнения:
print("=== ПЛАН ДО ВЫПОЛНЕНИЯ ===")
pre_execution = check_runtime_filtering_active(query)
print(f"AQE активен: {pre_execution['aqe_active']}")
print(f"Bloom Filter: {pre_execution['bloom_filter_active']}")
# После выполнения (финальный план):
query.collect()
print("\n=== ФИНАЛЬНЫЙ ПЛАН (после AQE) ===")
post_execution = check_runtime_filtering_active(query)
print(f"Bloom Filter применён: {post_execution['bloom_filter_active']}")
print(f"Dynamic Pruning: {post_execution['dynamic_pruning']}")
print(f"Оптимизаций найдено: {post_execution['optimization_count']}")
Dynamic Partition Pruning: отдельный механизм¶
Runtime Filtering включает не только Bloom Filter, но и Dynamic Partition Pruning (DPP) — когда данные партиционированы по ключу JOIN:
# Если fact_sales партиционирован по store_id:
spark.sql("""
CREATE TABLE fact_sales_partitioned
USING PARQUET
PARTITIONED BY (store_id)
AS SELECT * FROM fact_sales
""")
# Dynamic Partition Pruning:
# Spark сначала получает список store_id из dim_stores WHERE region='Europe'
# Затем читает ТОЛЬКО партиции с этими store_id
# Остальные партиции полностью пропускаются без чтения!
# В плане:
# FileScan parquet [...]
# PartitionFilters: [isnotnull(store_id#2),
# dynamicpruningexpression(store_id#2 IN dynamicpruning#42)]
# ← DPP фильтр внедрён в PartitionFilters!
9. Ограничения: когда Runtime Filtering не работает¶
Сценарии неэффективности¶
# ── Сценарий 1: Build-сторона слишком большая ─────────────────────────
# Нет смысла строить Bloom Filter для 1B ключей
# BF размер: 1B × 8 bytes / compression ≈ 1+ GB → не broadcast!
large_dim = spark.range(1_000_000_000) # гигантское "измерение"
# AQE не будет строить BF — нет выгоды
result = facts.join(large_dim, "store_id")
# ── Сценарий 2: Низкая селективность ─────────────────────────────────
# Если 90% факт-строк совпадают с измерением — BF отсеет только 10%
# Overhead на BF > выигрыш от фильтрации
popular_stores = stores.filter(F.col("region") != "Antarctica")
# 99% магазинов остаётся → BF фильтрует только 1% → почти бесполезно
# ── Сценарий 3: Несовместимые типы данных JOIN-ключей ─────────────────
# Если store_id в facts = INT, а в dim_stores = STRING
# Spark не может применить BF без конвертации
# Неявная конвертация разрушает pushdown
facts_str = facts.withColumn("store_id", F.col("store_id").cast("string"))
stores_int = stores.withColumn("store_id", F.col("store_id").cast("int"))
# Тип конфликт → нет BF pushdown!
# РЕШЕНИЕ: явное приведение типов перед JOIN
# ── Сценарий 4: AQE отключён ─────────────────────────────────────────
spark.conf.set("spark.sql.adaptive.enabled", "false")
# Без AQE нет Runtime Filtering
# Spark строит план статически, не зная размеры после фильтрации
# ── Сценарий 5: Сложные UDF разрушают pushdown ───────────────────────
@F.udf(returnType="int")
def complex_transform(store_id):
return store_id * 2 + 1 # произвольная трансформация
facts_udf = facts.withColumn("join_key", complex_transform(F.col("store_id")))
# Catalyst не может применить BF к произвольному UDF
# РЕШЕНИЕ: использовать нативные Spark SQL функции вместо UDF
Диагностика: почему BF не применился¶
def diagnose_missing_runtime_filter(facts_df, dim_df, join_key: str) -> None:
"""Диагностирует почему Runtime Filter может не применяться."""
print("=== Диагностика Runtime Filtering ===\n")
# 1. Проверяем AQE
aqe_enabled = spark.conf.get("spark.sql.adaptive.enabled") == "true"
bf_enabled = spark.conf.get(
"spark.sql.optimizer.runtime.bloomFilter.enabled", "false"
) == "true"
print(f"1. AQE включён: {'✅' if aqe_enabled else '❌'}")
print(f" BF включён: {'✅' if bf_enabled else '❌'}")
# 2. Проверяем размер join-стороны
try:
dim_size_bytes = dim_df._jdf.queryExecution().optimizedPlan().stats().sizeInBytes()
threshold = int(spark.conf.get(
"spark.sql.optimizer.runtime.bloomFilter.creationSideThreshold",
"10485760"
))
dim_size_mb = dim_size_bytes / 1024 / 1024
print(f"\n2. Размер build-side: {dim_size_mb:.1f} MB")
print(f" BF threshold: {threshold/1024/1024:.0f} MB")
if dim_size_bytes > threshold:
print(" ⚠️ Build-side БОЛЬШЕ threshold → BF может не создаться")
else:
print(" ✅ Build-side в пределах threshold")
except Exception:
print("\n2. Не удалось получить статистику размера")
# 3. Проверяем типы ключей
facts_type = dict(facts_df.dtypes).get(join_key)
dim_type = dict(dim_df.dtypes).get(join_key)
print(f"\n3. Тип ключа в facts: {facts_type}")
print(f" Тип ключа в dim: {dim_type}")
if facts_type != dim_type:
print(f" ⚠️ НЕСОВПАДЕНИЕ ТИПОВ! Может блокировать BF pushdown")
print(f" Решение: явное CAST перед JOIN")
else:
print(f" ✅ Типы совпадают")
10. Best Practices для Lakehouse Star-Schema архитектур¶
Проектирование для максимальной эффективности Runtime Filtering¶
from pyspark.sql import SparkSession, functions as F
def setup_optimized_star_schema(spark: SparkSession) -> None:
"""
Настройка оптимизированной Star Schema для максимального
эффекта от Runtime Filtering.
"""
# ── 1. Правильные типы данных для JOIN-ключей ───────────────────
# Используйте INT или BIGINT для join-ключей, не STRING!
# STRING → медленнее хэш, больше BF, хуже сжатие
# ── 2. Партиционирование fact-таблицы ────────────────────────────
# По временному измерению (sale_date) — самые частые фильтры
# Dynamic Partition Pruning пропустит ненужные партиции
# ── 3. Z-Order/Liquid Clustering для join-ключей ─────────────────
spark.sql("""
OPTIMIZE silver.fact_sales
ZORDER BY (store_id, product_id)
-- После этого строки с одинаковым store_id находятся в
-- одних и тех же Row Groups → BF + Row Group skip максимально эффективны
""")
# ── 4. ANALYZE TABLE для точной статистики ───────────────────────
# CBO использует статистику для выбора правильной стратегии JOIN
# AQE использует статистику для порогов BF creation
for table in ["dim_stores", "dim_products", "dim_customers", "dim_dates"]:
print(f"Collecting statistics for {table}...")
spark.sql(f"ANALYZE TABLE silver.{table} COMPUTE STATISTICS")
spark.sql(f"""
ANALYZE TABLE silver.{table}
COMPUTE STATISTICS FOR ALL COLUMNS
""")
# ── 5. Включаем все нужные оптимизации ───────────────────────────
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.optimizer.runtime.bloomFilter.enabled", "true")
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")
spark.conf.set("spark.sql.cbo.enabled", "true")
spark.conf.set("spark.sql.statistics.histogram.enabled", "true")
print("Star Schema оптимизирована для Runtime Filtering ✅")
def star_schema_query_template(
spark: SparkSession,
fact_table: str,
dimension_filters: dict[str, str], # {dim_table: filter_condition}
join_keys: dict[str, str], # {dim_table: join_key}
select_cols: list[str],
agg_expr: str,
) -> "DataFrame":
"""
Шаблон запроса к Star Schema с максимальной поддержкой Runtime Filtering.
Порядок JOIN важен: сначала самые селективные измерения!
AQE и Runtime Filtering работают лучше когда маленькие
и высокоселективные измерения обрабатываются первыми.
"""
# Начинаем с фактовой таблицы
result = spark.table(fact_table)
# JOIN с измерениями (сортируем по ожидаемой селективности)
for dim_table, filter_cond in dimension_filters.items():
join_key = join_keys[dim_table]
dim_df = spark.table(dim_table)
if filter_cond:
dim_df = dim_df.filter(filter_cond)
result = result.join(dim_df, join_key, "inner")
return result
# Использование:
result = star_schema_query_template(
spark,
fact_table="silver.fact_sales",
dimension_filters={
"silver.dim_stores": "region = 'Europe'",
"silver.dim_products": "category = 'Electronics'",
"silver.dim_dates": "year = 2024 AND quarter = 4",
},
join_keys={
"silver.dim_stores": "store_id",
"silver.dim_products": "product_id",
"silver.dim_dates": "sale_date",
},
select_cols=["store_name", "product_name", "amount"],
agg_expr="SUM(amount)",
)
result.groupBy("store_name", "product_name").sum("amount").show(20)
Мониторинг эффективности Runtime Filtering¶
def measure_runtime_filtering_impact(
spark: SparkSession,
query_df,
label: str = "",
) -> dict:
"""
Измеряет реальный эффект Runtime Filtering через History Server API.
"""
import time
import requests
app_id = spark.sparkContext.applicationId
history_url = f"http://localhost:18080/api/v1/applications/{app_id}"
t0 = time.time()
count = query_df.count()
elapsed = time.time() - t0
print(f"\n{'='*50}")
print(f"Запрос: {label}")
print(f"Результат: {count:,} строк за {elapsed:.1f}с")
print(f"\nДля детальной диагностики:")
print(f" Spark UI: http://localhost:4040/SQL/")
print(f" Смотрите: Scan rows read vs Input rows")
print(f" Идеально: Input rows << Scan rows (= Runtime Filter работает)")
return {
"label": label,
"elapsed_seconds": elapsed,
"result_rows": count,
}
# Сравнение с и без Runtime Filtering:
results = []
# Без оптимизаций
spark.conf.set("spark.sql.adaptive.enabled", "false")
results.append(measure_runtime_filtering_impact(spark, query, "БЕЗ AQE/RF"))
# С оптимизациями
spark.conf.set("spark.sql.adaptive.enabled", "true")
results.append(measure_runtime_filtering_impact(spark, query, "С AQE + Runtime Filtering"))
print(f"\nИтоговое ускорение: {results[0]['elapsed_seconds']/results[1]['elapsed_seconds']:.1f}x")
Итоги: когда Runtime Filtering даёт максимальный эффект¶
Runtime Filtering максимально эффективен когда:
- Star Schema или Snowflake Schema с явной асимметрией размеров (fact >> dim)
- Измерения высокоселективны (фильтры убирают > 70% ключей)
- Таблица фактов большая (гигабайты — петабайты)
- Данные хранятся в Parquet/ORC/Delta/Iceberg с Row Group статистикой
- Таблица фактов Z-Ordered или сортирована по join-ключам
Три ключевых правила:
- Включите AQE (
adaptive.enabled=true) — без него Runtime Filtering невозможен - Собирайте статистику (
ANALYZE TABLE ... COMPUTE STATISTICS) — Catalyst должен знать реальные размеры измерений - Оптимизируйте физическое расположение данных в фактовой таблице — Z-Order по join-ключам позволяет Bloom Filter пропускать целые Row Groups
Чек-лист проверки: запустите df.explain("formatted") и найдите might_contain(bloomfilter, — если эта строка есть, Runtime Filtering активен и данные фактов отсеиваются до Shuffle.