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 и матрица выбора абстракции.
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 — это хирургические инструменты для конкретных ниш, а не универсальные «ускорители».