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 архитектур.

optimization

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 применяется на двух уровнях:

  1. Row Group уровень (Parquet): если ни одно значение в Row Group не может пройти фильтр по min/max статистике — весь блок пропускается без чтения
  2. Строчный уровень: для строк внутри прочитанных 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-ключам

Три ключевых правила:

  1. Включите AQE (adaptive.enabled=true) — без него Runtime Filtering невозможен
  2. Собирайте статистику (ANALYZE TABLE ... COMPUTE STATISTICS) — Catalyst должен знать реальные размеры измерений
  3. Оптимизируйте физическое расположение данных в фактовой таблице — Z-Order по join-ключам позволяет Bloom Filter пропускать целые Row Groups

Чек-лист проверки: запустите df.explain("formatted") и найдите might_contain(bloomfilter, — если эта строка есть, Runtime Filtering активен и данные фактов отсеиваются до Shuffle.