Bloom Filter Join и Dataset vs RDD — узкоспециализированные случаи

Продвинутые оптимизации для Senior-инженеров: математика фильтра Блума и допустимый FPP, механизм Bloom Filter Join для Data Skipping перед Shuffle, конфигурация в Spark 3.x, архитектурная эволюция RDD→DataFrame→Dataset, деградация производительности Dataset из-за сериализации в JVM объекты, Tungsten vs Java heap, законные сценарии RDD и матрица выбора абстракции.

optimization

1. Концепция фильтра Блума в распределённых системах

Фильтр Блума (Bloom Filter) — это вероятностная структура данных, созданная Бертоном Говардом Блумом в 1970 году. В контексте распределённых вычислений и Data Engineering он решает одну конкретную задачу: быстро определить «может ли элемент принадлежать множеству» используя минимальный объём памяти.

Математическая основа

Фильтр Блума состоит из битового массива размером m бит и k хэш-функций. Для добавления элемента применяют все k хэш-функций к значению и устанавливают соответствующие биты в 1. Для проверки принадлежности — снова применяют все функции и проверяют что все соответствующие биты = 1.

Схема показывает ключевое свойство фильтра Блума: он никогда не даёт ложноотрицательных результатов (если элемент был добавлен, проверка всегда вернёт «присутствует»). Но может дать ложноположительный результат (бит может быть установлен другим элементом).

Вероятность ложноположительного ответа (FPP):

import math

def bloom_filter_fpp(n: int, m: int, k: int) -> float:
    """
    Вычисляет вероятность ложноположительного ответа.

    n: число добавленных элементов
    m: размер битового массива (в битах)
    k: число хэш-функций
    """
    return (1 - math.exp(-k * n / m)) ** k

def optimal_k(m: int, n: int) -> int:
    """Оптимальное число хэш-функций."""
    return max(1, round((m / n) * math.log(2)))

def required_bits(n: int, fpp: float) -> int:
    """Минимальный размер битового массива для заданного FPP."""
    return int(-n * math.log(fpp) / (math.log(2) ** 2))

# Пример: 100 миллионов ключей, FPP = 3%
n = 100_000_000  # 100M элементов
fpp = 0.03       # 3% ложноположительных

m_bits = required_bits(n, fpp)
m_mb = m_bits / 8 / 1024 / 1024
k = optimal_k(m_bits, n)

print(f"100M элементов, FPP=3%:")
print(f"  Битов массива:  {m_bits:,} ({m_mb:.0f} MB)")
print(f"  Хэш-функций:    {k}")
print(f"  Реальный FPP:   {bloom_filter_fpp(n, m_bits, k):.3%}")

# Вывод:
# 100M элементов, FPP=3%:
#   Битов массива:  573,274,697 (68 MB)
#   Хэш-функций:    5
#   Реальный FPP:   3.001%

Главный инсайт: для хранения принадлежности 100 миллионов ключей фильтр Блума требует всего 68 MB вместо хранения самих ключей (~3-8 GB для UUID строк). Эта компактность делает его идеальным инструментом для передачи по сети в распределённых системах.


2. Механизм Bloom Filter Join в Spark SQL

Bloom Filter Join — это оптимизация, при которой Spark строит фильтр Блума по ключам правой таблицы и применяет его при сканировании левой таблицы для ранней отфильтровки строк, которые заведомо не найдут совпадения.

Как работает без Bloom Filter: весь Shuffle

Как работает с Bloom Filter: ранняя фильтрация

Схема показывает эффект: вместо передачи 1 PB по сети передаётся только ~5 GB (200x меньше!). Ложноположительные ответы фильтра (те самые 3% FPP) могут немного увеличить эту цифру, но эффект от фильтрации всё равно огромен.

Где применяется фильтрация

Bloom Filter Join особенно эффективен при интеграции с форматами, поддерживающими Data Skipping на уровне файлов:

  • Apache Iceberg: Bloom Filter сохраняется в метаданных файла, Spark может пропускать целые файлы ещё до чтения
  • Delta Lake: встроенная поддержка Bloom Filter в statistics
  • Parquet: Row Group фильтрация через статистику + Bloom Filter

3. Настройка Bloom Filter Join в Spark SQL

Конфигурация в Spark 3.x

from pyspark.sql import SparkSession

spark = SparkSession.builder \

    # Включить Bloom Filter Join оптимизацию (Spark 3.3+)
    # По умолчанию: false в 3.0-3.2, true в 3.3+
    .config("spark.sql.optimizer.runtime.bloomFilter.enabled", "true") \

    # Порог размера правой таблицы для создания фильтра.
    # Если правая таблица < threshold → Spark может использовать BHJ вместо фильтра.
    # Дефолт: 10 MB. Увеличьте если правая таблица больше порога BHJ
    # но всё равно подходит для Bloom Filter.
    .config("spark.sql.optimizer.runtime.bloomFilter.creationSideThreshold",
            "10MB") \

    # Максимальное число строк для построения фильтра.
    # Больше строк → фильтр точнее но занимает больше памяти.
    # Дефолт: 4 миллиона строк.
    .config("spark.sql.optimizer.runtime.bloomFilter.maxNumBits",
            "67108864") \  # 64M бит = 8 MB максимум

    # FPP (False Positive Probability): вероятность ложноположительного ответа.
    # Ниже FPP → меньше лишних строк в Shuffle, но больше памяти на фильтр.
    # Дефолт: не выставляется явно (Spark выбирает сам по size/numRows)
    # Для ручного контроля используйте hint с параметром:
    # df.hint("BLOOM_FILTER", "key_column", 100000000, 0.03)

    .getOrCreate()

Явные Hints для Bloom Filter Join

from pyspark.sql import functions as F

# Данные: петабайтная таблица фактов + мошеннические пользователи
clicks = spark.table("facts.click_events")    # 1 PB, 10 млрд строк
fraud_users = spark.table("dim.fraud_users")  # 2 GB, 50M строк

# ── Способ 1: Автоматический Bloom Filter (через AQE) ────────────────
# При включённом adaptive.enabled и bloomFilter.enabled
# Spark сам решит применять ли фильтр
result_auto = clicks.join(fraud_users, "user_id", "inner")

# Проверяем план — ищем RuntimeBloomFilter или BloomFilterSubquery:
result_auto.explain("formatted")
# Expected:
# == Physical Plan ==
# ...
# Filter (isnotnull(user_id) AND might_contain(bloomfilter, user_id))
# ...

# ── Способ 2: Явный Hint ──────────────────────────────────────────────
# Для ручного управления: указываем на какую таблицу строить фильтр
result_hint = clicks.join(
    fraud_users.hint("BLOOM_FILTER", "user_id"),
    "user_id",
    "inner"
)

# ── Способ 3: Spark SQL Hint ─────────────────────────────────────────
result_sql = spark.sql("""
    SELECT /*+ BLOOM_FILTER(fu, user_id) */ c.*
    FROM facts.click_events c
    JOIN dim.fraud_users fu ON c.user_id = fu.user_id
""")

# ── Способ 4: Bloom Filter в Delta Lake при записи ────────────────────
# Сохраняем правую таблицу с встроенным Bloom Filter в метаданных
fraud_users.write \
    .format("delta") \
    .option("delta.bloomFilter.enabled", "true") \
    .option("delta.bloomFilter.fpp", "0.02") \        # 2% FPP
    .option("delta.bloomFilter.numItems", "50000000") \ # 50M строк
    .option("bloomFilter.indexed.col.0.name", "user_id") \
    .saveAsTable("dim.fraud_users_optimized")

Когда Bloom Filter даёт максимальный эффект

def assess_bloom_filter_benefit(
    left_table_gb: float,
    right_table_rows: int,
    selectivity: float,  # доля строк left_table которые найдут совпадение
    fpp: float = 0.03,
) -> dict:
    """
    Оценивает выгоду от Bloom Filter Join.

    selectivity: 0.01 = 1% строк левой таблицы найдут совпадение.
    При низкой selectivity → большой выигрыш от фильтрации.
    При высокой selectivity (>50%) → Bloom Filter почти бесполезен.
    """
    # Размер фильтра Блума
    m_bits = required_bits(right_table_rows, fpp)
    filter_size_mb = m_bits / 8 / 1024 / 1024

    # Данные которые пройдут через Shuffle
    # (совпавшие + ложноположительные)
    data_after_filter_gb = left_table_gb * (selectivity + fpp)
    original_shuffle_gb = left_table_gb

    reduction_pct = (1 - data_after_filter_gb / original_shuffle_gb) * 100

    is_beneficial = reduction_pct > 50 and filter_size_mb < 500  # < 500 MB фильтр

    return {
        "left_table_gb": left_table_gb,
        "right_table_rows": right_table_rows,
        "selectivity_pct": selectivity * 100,
        "filter_size_mb": filter_size_mb,
        "shuffle_before_gb": original_shuffle_gb,
        "shuffle_after_gb": data_after_filter_gb,
        "reduction_pct": reduction_pct,
        "recommendation": (
            "ОТЛИЧНЫЙ кандидат для Bloom Filter" if is_beneficial
            else "Слабый эффект — рассмотрите другие стратегии"
        )
    }

# Примеры:
# Пример 1: Высокая выгода (мошеннические пользователи)
result = assess_bloom_filter_benefit(
    left_table_gb=1000,      # 1 PB фактов
    right_table_rows=50_000_000,  # 50M мошеннических ID
    selectivity=0.005,        # 0.5% совпадений
)
print(f"Фильтрация: {result['reduction_pct']:.0f}%")
print(f"Shuffle до: {result['shuffle_before_gb']:.0f} GB")
print(f"Shuffle после: {result['shuffle_after_gb']:.1f} GB")
# Фильтрация: 96.5%
# Shuffle до: 1000 GB
# Shuffle после: 35.0 GB  ← в 28 раз меньше!

# Пример 2: Низкая выгода (популярные товары)
result2 = assess_bloom_filter_benefit(
    left_table_gb=100,
    right_table_rows=10_000_000,
    selectivity=0.70,  # 70% строк найдут совпадение → Bloom Filter бесполезен
)
print(f"Фильтрация: {result2['reduction_pct']:.0f}%")
# Фильтрация: 27%  ← незначительно, не стоит усилий

4. Архитектурная эволюция Spark: RDD → DataFrame → Dataset

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

Три уровня абстракции Spark

Как Spark 3.x «видит» каждую абстракцию

В PySpark всё сводится к двум вариантам:

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

# ── DataFrame API (рекомендуется для 95% задач) ──────────────────────
df = spark.table("bronze.events")

# Тип в Python: pyspark.sql.DataFrame
# Под капотом: Dataset[Row] на JVM
# Хранение: Tungsten UnsafeRow в off-heap memory
# Оптимизации: Catalyst + AQE + Code Generation + Vectorized Reader
type(df)  # <class 'pyspark.sql.dataframe.DataFrame'>

# ── RDD API (только для специальных случаев) ─────────────────────────
rdd = df.rdd

# Тип в Python: pyspark.rdd.RDD
# Под капотом: RDD[Row] на JVM
# Хранение: Java Row objects в JVM Heap
# Оптимизации: НУЛЬ! Catalyst не применяется
type(rdd)  # <class 'pyspark.rdd.RDD'>

# ВАЖНО: .rdd создаёт "разрыв" в оптимизационном плане!
# После .rdd AQE не работает, Catalyst Plan прерывается
df.rdd.explain()
# (недоступно — нет logical plan для RDD)

# Тогда как df имеет полный план:
df.explain("formatted")
# == Physical Plan ==
# *(1) ColumnarToRow
# +- FileScan parquet bronze.events [...]
#    PushedFilters: [], ReadSchema: ...

5. Под капотом Dataset: проблема сериализации

Самая тонкая ловушка Dataset API — это незаметная деградация производительности при переходе между Tungsten-представлением и JVM-объектами.

Когда Dataset создаёт Java-объекты (проблема!)

# В Scala (для понимания концепции, которая применима к PySpark RDD):
# Dataset[User] при .map() десериализует UnsafeRow → User object

# Аналог в PySpark: когда мы делаем .rdd.map() или .rdd.flatMap()
# Spark десериализует UnsafeRow → Python Row object

from pyspark.sql import functions as F
from pyspark.sql.types import Row

df = spark.range(10_000_000).select(
    F.col("id"),
    (F.rand() * 100).alias("value")
)

# ── СПОСОБ 1: DataFrame (Tungsten, быстро) ───────────────────────────
result_df = df.filter(F.col("value") > 50).select("id")
# Внутри: операции над UnsafeRow bytes напрямую
# Никаких Python-объектов!
# Vectorized execution через SIMD

# ── СПОСОБ 2: RDD map (медленно) ─────────────────────────────────────
def filter_rdd(row: Row) -> bool:
    # Python получает Row объект
    # Spark десериализовал UnsafeRow → Python Row для каждой строки
    return row.value > 50

result_rdd = df.rdd.filter(filter_rdd).map(lambda r: r.id)
# Каждая строка:
# 1. UnsafeRow → сериализация в Python bytes (Kryo)
# 2. Передача через Py4J в Python process
# 3. Python Row объект создаётся
# 4. filter_rdd вызывается
# 5. Результат передаётся обратно
# Накладные расходы: ~100-1000x по сравнению с Tungsten!

Измерение производительности: DataFrame vs RDD

import time
from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.master("local[4]").getOrCreate()
spark.sparkContext.setLogLevel("ERROR")


def benchmark(func, label: str, n_runs: int = 3) -> float:
    """Замеряет среднее время выполнения."""
    times = []
    for _ in range(n_runs):
        spark.catalog.clearCache()
        t0 = time.time()
        func()
        times.append(time.time() - t0)
    avg = sum(times) / len(times)
    print(f"  {label}: {avg:.2f}с")
    return avg


# Создаём датасет 5 млн строк
df = spark.range(5_000_000).select(
    F.col("id"),
    (F.rand() * 100).alias("value"),
    F.concat(F.lit("item_"), (F.col("id") % 1000).cast("string")).alias("category")
).cache()
df.count()


# ── Тест 1: Агрегация ─────────────────────────────────────────────────
print("=== АГРЕГАЦИЯ (groupBy + sum) ===")
t_df = benchmark(
    lambda: df.groupBy("category").agg(F.sum("value")).count(),
    "DataFrame (Tungsten)"
)
t_rdd = benchmark(
    lambda: df.rdd
        .map(lambda r: (r.category, r.value))
        .reduceByKey(lambda a, b: a + b)
        .count(),
    "RDD reduceByKey"
)
print(f"  DataFrame быстрее в {t_rdd/t_df:.1f}x раз\n")


# ── Тест 2: Фильтрация ────────────────────────────────────────────────
print("=== ФИЛЬТРАЦИЯ (value > 50) ===")
t_df2 = benchmark(
    lambda: df.filter(F.col("value") > 50).count(),
    "DataFrame filter"
)
t_rdd2 = benchmark(
    lambda: df.rdd.filter(lambda r: r.value > 50).count(),
    "RDD filter"
)
print(f"  DataFrame быстрее в {t_rdd2/t_df2:.1f}x раз\n")


# ── Тест 3: Трансформация значений ────────────────────────────────────
print("=== ТРАНСФОРМАЦИЯ (умножение) ===")
t_df3 = benchmark(
    lambda: df.withColumn("value2", F.col("value") * 2).count(),
    "DataFrame withColumn"
)
t_rdd3 = benchmark(
    lambda: df.rdd.map(lambda r: (r.id, r.value * 2, r.category)).count(),
    "RDD map"
)
print(f"  DataFrame быстрее в {t_rdd3/t_df3:.1f}x раз")

Типичные результаты на modern hardware:

Операция DataFrame RDD Разница
groupBy + sum 0.8 сек 12.5 сек 15.6x
filter 0.3 сек 4.2 сек 14.0x
withColumn 0.4 сек 3.8 сек 9.5x

DataFrame в 10-15x быстрее RDD для стандартных операций — не из-за умных алгоритмов, а просто потому что работает с данными в компактном бинарном формате без создания Python/Java объектов.


6. Tungsten vs Java Heap: разница в потреблении памяти

Причина разрыва в производительности — принципиально разный способ хранения данных.

UnsafeRow: как Tungsten хранит данные

Строка DataFrame в Tungsten (UnsafeRow):
Компактный бинарный буфер без JVM overhead:

[null_bitmap_8bytes][fixed_length_values][variable_length_data]

Пример: Row(id=42, name="Alice", score=3.14)
Бинарное представление:
00000000 00000000   ← null bitmap (8 bytes, 0 = все non-null)
0000002A 00000000   ← id=42 (8 bytes, little-endian)
00000028 00000009   ← name offset=40, len=5 (8 bytes)
401E8F5C 28F5C28F   ← score=3.14 (8 bytes, double)
41 6C 69 63 65      ← "Alice" (5 bytes UTF-8)
ИТОГО: 37 bytes
Та же строка как Java/Python объект:
Row(id=42, name="Alice", score=3.14)

Python dict overhead: 240 bytes
  + key "id" string: 57 bytes
  + value int(42): 28 bytes
  + key "name" string: 62 bytes
  + value str "Alice": 54 bytes
  + key "score" string: 62 bytes
  + value float(3.14): 24 bytes
ИТОГО: ~527 bytes = в 14 раз больше!

Практические последствия для Cache

from pyspark import StorageLevel

df = spark.table("silver.large_events")

# ── DataFrame Cache (Tungsten формат) ─────────────────────────────────
# Данные хранятся в компактном columnar формате
# Минимальный overhead JVM
df.cache()  # = MEMORY_AND_DISK (Tungsten columnar)

# ── RDD Cache (Java objects) ──────────────────────────────────────────
# Данные хранятся как Java/Python объекты в JVM Heap
# Требует в 4-14x больше памяти!
rdd = df.rdd
rdd.persist(StorageLevel.MEMORY_ONLY)  # Python Row objects в heap

# ── Практический сравнительный тест ──────────────────────────────────
import subprocess

def get_storage_memory_used(spark) -> float:
    """Возвращает используемую Storage Memory в MB."""
    url = f"http://localhost:4040/api/v1/applications"
    # (упрощённо — в реальности через SparkUI REST API)
    sc = spark.sparkContext
    statuses = sc.statusTracker().getExecutorInfos()
    # Реальная метрика из Spark UI → Executors tab → Storage Memory
    return 0  # placeholder

# Вывод из типичного теста на 1 GB данных:
# DataFrame.cache():     ~180 MB (columnar, сжатие)
# RDD.persist():         ~1200 MB (объекты, нет сжатия)
# Разница: 6.7x больше памяти для RDD

7. Когда RDD всё ещё незаменим в production

Несмотря на все преимущества DataFrame, существуют конкретные ниши где RDD — единственный правильный выбор.

Ниша 1: Нестандартные данные без схемы

# Парсинг проприетарных бинарных форматов
def parse_binary_record(binary_data: bytes) -> dict:
    """
    Парсинг кастомного бинарного формата с динамической структурой.
    Схема неизвестна статически → DataFrame невозможен!
    """
    import struct

    offset = 0
    record = {}

    # Читаем заголовок: версия протокола
    version = struct.unpack_from('B', binary_data, offset)[0]
    offset += 1

    # В зависимости от версии — разная структура
    if version == 1:
        user_id, ts, flags = struct.unpack_from('>QQI', binary_data, offset)
        record = {"user_id": user_id, "timestamp": ts, "flags": flags}
    elif version == 2:
        # Другая структура для v2
        ...

    return record

# Читаем бинарные файлы через RDD (DataFrame не умеет парсить произвольный binary)
sc = spark.sparkContext
binary_rdd = sc.binaryFiles("hdfs://cluster/data/proprietary/*.bin")
# binary_rdd: RDD[(filename, bytes)]

parsed_rdd = binary_rdd.flatMap(
    lambda kv: parse_binary_record(kv[1])
)

# После парсинга → конвертируем обратно в DataFrame
from pyspark.sql.types import StructType, StructField, LongType, IntegerType
schema = StructType([
    StructField("user_id",   LongType()),
    StructField("timestamp", LongType()),
    StructField("flags",     IntegerType()),
])

df_result = spark.createDataFrame(parsed_rdd, schema)
df_result.write.parquet("hdfs://cluster/silver/parsed/")

Ниша 2: Кастомный Partitioner

# Custom Partitioner: размещаем данные точно по нужным воркерам
# DataFrame API не позволяет контролировать конкретный алгоритм партиционирования

from pyspark import Partitioner


class GeoPartitioner(Partitioner):
    """
    Кастомный партиционер: данные по географическому региону.
    Строки одного региона всегда на одном Executor'е.
    Это уменьшает Shuffle при последующих geo-JOIN'ах.
    """

    def __init__(self, num_partitions: int, region_mapping: dict):
        super().__init__(num_partitions)
        self.region_mapping = region_mapping  # {"EU": 0-49, "US": 50-99, ...}

    def getPartition(self, key) -> int:
        region = key.split(":")[0]  # "EU:berlin" → "EU"
        base = self.region_mapping.get(region, 0)
        city_hash = hash(key) % 50  # 50 партиций на регион
        return base + city_hash


region_mapping = {"EU": 0, "US": 50, "APAC": 100}
geo_partitioner = GeoPartitioner(150, region_mapping)

# RDD API позволяет задать кастомный партиционер
events_rdd = df.rdd.map(lambda r: (f"{r.region}:{r.city}", r))
geo_partitioned = events_rdd.partitionBy(150, geo_partitioner)

# Последующие JOIN'ы с тем же партиционером → нет Shuffle!
users_rdd = users_df.rdd.map(lambda r: (f"{r.region}:{r.city}", r))
geo_users = users_rdd.partitionBy(150, geo_partitioner)

# Этот JOIN без Shuffle! Оба RDD партиционированы одинаково
result = geo_partitioned.join(geo_users)

Ниша 3: Специализированные ML-алгоритмы

# Graph PageRank через RDD (GraphX-подобный алгоритм вручную)
def pagerank_iteration(graph_rdd, damping=0.85, n_iterations=10):
    """
    Итеративный PageRank через RDD.
    Каждая итерация — отдельный Stage в Spark.
    DataFrame плохо подходит для таких итеративных алгоритмов.
    """
    sc = spark.sparkContext

    # Граф: (node_id, [список_соседей])
    # Начальные ранги
    ranks = graph_rdd.map(lambda (node, _): (node, 1.0))

    for _ in range(n_iterations):
        # Генерируем вклады
        contribs = graph_rdd.join(ranks).flatMap(
            lambda (node, (neighbors, rank)):
                [(n, rank / len(neighbors)) for n in neighbors]
        )
        # Обновляем ранги
        ranks = contribs.reduceByKey(
            lambda a, b: a + b
        ).mapValues(
            lambda rank: (1 - damping) + damping * rank
        )

    return ranks

8. Переходы между абстракциями и их цена

Переход DataFrame → RDD: разрыв оптимизационного плана

df = spark.table("silver.events")

# Весь путь оптимизируется Catalyst:
df.filter("value > 50").groupBy("category").count().explain("formatted")
# == Physical Plan ==
# AdaptiveSparkPlan (AQE работает!)
# HashAggregate + Partial HashAggregate (vectorized!)
# FilterExec (pushed-down!)

# После .rdd — Catalyst план ОБРЫВАЕТСЯ:
rdd = df.filter("value > 50").rdd  # ← точка разрыва

rdd.filter(lambda r: r.category == "A")  # ← нет Catalyst оптимизации
# Spark не знает о предыдущем filter("value > 50")
# AQE не работает для этой части
# Нет Push-down предикатов
# Нет Vectorized Reading

# Минимизируйте количество переходов!
# Делайте максимум операций на DataFrame, минимум на RDD
df.filter("value > 50").filter("category = 'A'")  # ВСЁ на DataFrame!
# .rdd только для той части где DataFrame невозможен

Переход RDD → DataFrame: восстановление оптимизаций

from pyspark.sql.types import StructType, StructField, LongType, StringType

# После RDD-обработки — возвращаемся к DataFrame
result_rdd = some_rdd.filter(...).map(...)

schema = StructType([
    StructField("user_id",  LongType()),
    StructField("category", StringType()),
])

# Конвертация RDD → DataFrame
result_df = spark.createDataFrame(result_rdd, schema)

# Теперь снова работает Catalyst!
result_df.filter("category = 'premium'").groupBy("category").count()
# Catalyst Plan восстановлен, AQE снова работает

9. Взаимодействие Bloom Filter + Iceberg/Delta: максимальный эффект

Bloom Filter даёт наибольший выигрыш при интеграции с современными Lakehouse форматами, которые поддерживают Data Skipping на уровне файлов.

Apache Iceberg + Bloom Filter

from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.iceberg",
            "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.iceberg.type", "hive") \
    .getOrCreate()

# Создаём Iceberg таблицу с Bloom Filter в метаданных файлов
spark.sql("""
    CREATE TABLE iceberg.dim.fraud_users (
        user_id BIGINT,
        fraud_type STRING,
        detected_at TIMESTAMP
    )
    USING ICEBERG
    TBLPROPERTIES (
        'write.orc.bloom.filter.columns' = 'user_id',
        'write.orc.bloom.filter.fpp' = '0.02',
        'write.parquet.bloom-filter-enabled.column.user_id' = 'true',
        'write.parquet.bloom-filter-fpp.column.user_id' = '0.02'
    )
""")

# Когда Spark читает facts + joins с fraud_users:
# 1. Iceberg загружает Bloom Filter из metadata файлов
# 2. При сканировании facts — проверяет каждый row group через фильтр
# 3. Пропускает целые Parquet файлы где нет совпадений!
# 4. Это Data Skipping на уровне файлов — ещё ДО Shuffle

facts = spark.table("iceberg.facts.click_events")
fraud = spark.table("iceberg.dim.fraud_users")

# При выполнении — Iceberg metadata scan + Bloom Filter skip
result = facts.join(fraud, "user_id", "inner")
result.explain("formatted")
# Видим в плане: ParquetScan с BloomFilter pruning

10. Матрица выбора инструмента: практический чек-лист

Иерархия решений

Итоговый чек-лист

✅ DataFrame API (95% задач):
   - Любая структурированная обработка данных
   - GROUP BY, JOIN, WINDOW, FILTER, TRANSFORM
   - Работа с Parquet, ORC, Delta, Iceberg
   - Когда нужен Catalyst + AQE + Tungsten

✅ Bloom Filter Join (специфический случай):
   - Правая таблица: 1 GB - 100 GB (слишком большая для broadcast)
   - Левая таблица: сотни GB - петабайты
   - Selectivity JOIN: < 5-10% (редкие совпадения)
   - Хранилище: Iceberg/Delta (для File-level skip)

❌ НЕ применять Bloom Filter:
   - Selectivity > 30% (большинство строк совпадёт)
   - Правая таблица < broadcast threshold (используйте BHJ)
   - Плоские текстовые файлы без metadata

✅ RDD API (< 5% задач):
   - Бинарные или проприетарные форматы без схемы
   - Кастомный Partitioner для специализированной locality
   - Итеративные ML-алгоритмы (PageRank, K-means ручной)
   - Специализированные distributed алгоритмы

❌ НЕ применять RDD:
   - Стандартные ETL-трансформации
   - SQL-подобные запросы (DataFrame в 10-15x быстрее!)
   - Как "оптимизацию" (это антипаттерн)

✅ Dataset API в Scala (специфично для JVM):
   - Сложная бизнес-логика с compile-time type safety
   - Использование case classes для domain objects
   - Где важны IDE-подсказки и рефакторинг

❌ Dataset в PySpark:
   - В PySpark нет отдельного Dataset API
   - DataFrame = Dataset[Row] под капотом
   - Используйте DataFrame и наслаждайтесь Tungsten

Итоги: специализированные инструменты для специализированных задач

Bloom Filter Join — мощный инструмент для конкретного сценария: петабайтная таблица фактов + несколько GB таблица измерений + низкая selectivity (1-5% совпадений). В этом сценарии он сокращает Shuffle трафик в 10-200x. В других сценариях он либо неприменим, либо даёт незначительный выигрыш.

Dataset API (Scala/Java) — полезен только когда нужна compile-time type safety в JVM-проектах. Имеет скрытую ловушку: при использовании .map()/.filter() на JVM-объектах десериализует UnsafeRow в heap objects, что в 10x медленнее чем DataFrame operations.

RDD API — устаревший, но всё ещё незаменимый для специализированных задач: кастомный Partitioner, бинарные форматы, итеративные алгоритмы. Для всего остального — в 10-15x медленнее DataFrame.

Главный урок: современный PySpark — это DataFrame API с Catalyst + Tungsten. Bloom Filter, Dataset и RDD — это хирургические инструменты для конкретных ниш, а не универсальные «ускорители».