Bloom Filter в Spark SQL: probabilistic dedup и ускорение join на больших таблицах
Вероятностные структуры данных в Spark: механика Bloom Filter, Runtime Filtering в Catalyst, stat.bloomFilter() API, Delta Lake/Iceberg индексы, pre-filter deduplication больших ETL-пайплайнов
Введение: почему JOIN и дедупликация становятся bottleneck¶
В мире распределённых вычислений существует закон, который нельзя обойти: чтобы найти пересечение двух наборов данных, их нужно сравнить. Для маленьких датасетов это тривиально. Но когда один датасет содержит 500 млн строк исторических транзакций, а другой - 50 млн записей суточного инкремента, сравнение превращается в операцию, которая может занять часы и положить кластер.
Рассмотрим реальный production-сценарий: у вас есть Silver-таблица кликстрима с историей за 2 года (2 млрд строк), и каждый час приходит новый батч - 10 млн событий. Ваша задача - дедупликация: выбрать из батча только те события, которых ещё нет в Silver-таблице. Классический подход - LEFT ANTI JOIN. Посмотрим, что при этом происходит:
- Spark читает весь инкремент (10 млн строк) -
Scan - Spark читает всю историю или её значительную часть (2 млрд строк) -
Scan - Выполняется Shuffle: обе стороны перераспределяются по ключу - сотни гигабайт сетевого трафика
- На каждом экзекьюторе строится HashTable из исторических данных
- Строки инкремента проверяются по этой HashTable
Самый дорогой шаг - Shuffle и построение HashTable. Именно здесь Bloom Filter меняет игру: он позволяет до Shuffle на каждом экзекьюторе быстро отсечь большинство строк инкремента, которых гарантированно нет в истории. Оставшиеся строки - потенциально новые - идут дальше в точную проверку.
Часть 1. Теория: механика Bloom Filter¶
1.1. Что такое Bloom Filter¶
Bloom Filter - вероятностная структура данных, предложенная Бёртоном Блумом в 1970 году. Она предназначена для быстрой проверки принадлежности элемента множеству и работает на основе битового массива фиксированного размера и набора хэш-функций.
Ключевые свойства:
- Компактность: для хранения фильтра на 1 млрд элементов с FPR=1% нужно всего ~1.2 ГБ против ~8 ГБ для HashSet
. - Скорость: проверка принадлежности - O(k), где k - количество хэш-функций (обычно 7–10). Никакого обхода дерева или цепочек коллизий.
- Гарантированное отсутствие False Negatives: если фильтр говорит «элемента НЕТ» - это истина. Строка, которую фильтр отклонил, точно отсутствует в исходном множестве.
- Возможны False Positives: если фильтр говорит «элемент ЕСТЬ» - он может ошибаться с вероятностью FPR (False Positive Rate).
Именно отсутствие False Negatives делает Bloom Filter безопасным для использования в качестве pre-filter: мы никогда не выбросим строку, которую должны были оставить. Максимум - оставим лишнюю строку (false positive), которую отсечёт точная проверка следующего шага.
1.2. Механика: как работает битовый массив¶
Bloom Filter - это битовый массив размером m бит, изначально заполненный нулями, и k независимых хэш-функций, каждая из которых отображает элемент в индекс от 0 до m-1.
Добавление элемента (при построении фильтра):
- Вычисляем
kхэшей от элемента:h1(x), h2(x), ..., hk(x) - Устанавливаем биты в позициях
h1(x) mod m,h2(x) mod m, ...,hk(x) mod mравными1
Проверка принадлежности:
- Вычисляем те же
kхэшей от проверяемого элемента - Проверяем биты в тех же позициях
- Если все биты равны
1→ «Возможно есть» (may be present) - Если хотя бы один бит равен
0→ «Точно нет» (definitely not present)
1.3. Математика: False Positive Rate¶
Вероятность ложноположительного срабатывания зависит от трёх параметров:
m- размер битового массива (бит)n- количество добавленных элементовk- количество хэш-функций
Формула FPR:
FPR ≈ (1 - e^(-kn/m))^k
Оптимальное k (при фиксированных m и n):
k_opt = (m/n) * ln(2) ≈ 0.693 * (m/n)
Размер массива для заданного n и FPR:
m = -n * ln(FPR) / (ln(2))^2 ≈ -1.44 * n * log2(FPR)
На практике это означает:
| FPR | Бит на элемент | Байт на 1 млн элементов |
|---|---|---|
| 10% (0.10) | 4.8 бит | ~0.6 МБ |
| 1% (0.01) | 9.6 бит | ~1.2 МБ |
| 0.1% (0.001) | 14.4 бит | ~1.8 МБ |
| 0.01% (0.0001) | 19.2 бит | ~2.4 МБ |
Вывод: для 1 млрд элементов с FPR=1% нам нужно около 1.2 ГБ памяти. Это ничто по сравнению с HashSet, который хранил бы сами ключи (~8 ГБ для BIGINT).
1.4. False Positives vs False Negatives¶
Понимание разницы критично для правильного применения:
| Тип ошибки | Что означает | Следствие |
|---|---|---|
| False Positive | Фильтр говорит «есть», но элемента нет | Лишняя строка пройдёт в точный JOIN |
| False Negative | Фильтр говорит «нет», но элемент есть | Невозможно для Bloom Filter |
Именно гарантия отсутствия False Negatives делает Bloom Filter безопасным для ETL-дедупликации: мы никогда не потеряем строку, которую должны были обработать. Максимальный ущерб от False Positive - небольшой overhead в последующем точном JOIN.
Часть 2. Под капотом Spark: как Bloom Filter встроен в движок¶
2.1. Два сценария применения¶
В Spark существуют два разных способа применения Bloom Filter:
Сценарий A: Runtime Bloom Filter Join Optimization (Spark 3.3+)
Автоматически создаётся оптимизатором Catalyst как часть Dynamic Filtering. Spark строит фильтр по меньшей стороне JOIN, транслирует его на экзекьюторы и применяет к большей стороне до Shuffle. Это называется Runtime Filtering или Dynamic Partition Pruning with Bloom Filter.
Сценарий B: Явное создание через stat.bloomFilter() API
Вы сами создаёте фильтр из DataFrame, сериализуете, рассылаете и применяете в коде. Полный контроль над параметрами.
Сценарий C: Bloom Filter Index в Parquet/Iceberg/Delta
Bloom Filter хранится как часть метаданных Parquet-файлов. При чтении Spark проверяет фильтр каждого файла и пропускает те, в которых точно нет нужных значений. Это file-level skipping.
2.2. Анатомия Runtime Bloom Filter Join¶
Когда Catalyst видит JOIN между большой таблицей (Build side ~10 ГБ) и маленькой (Probe side ~200 МБ), он может автоматически:
- Строит Bloom Filter по ключевой колонке меньшей стороны (Probe side)
- Транслирует фильтр всем экзекьюторам через широковещательную рассылку (Broadcast) - объект весит килобайты или мегабайты
- Применяет фильтр к большой стороне (Build side) на стадии сканирования - до Shuffle
- Строки, отклонённые фильтром, никогда не попадают в сеть
- Оставшиеся строки идут в финальный JOIN
Экономия: если 90% строк большой таблицы гарантированно не совпадают по ключу с маленькой, то Shuffle-трафик сокращается в 10 раз. Это прямое ускорение всего JOIN.
2.3. Конфигурация Runtime Bloom Filter¶
# Включение Runtime Bloom Filter в Spark (по умолчанию включён с Spark 3.3)
spark.conf.set("spark.sql.runtime.filter.enabled", "true")
# Максимальный размер фильтра для Broadcast (по умолчанию 4 МБ)
spark.conf.set("spark.sql.runtime.filter.bloomFilter.maxNumBits", "33554432") # 4 МБ в битах
# Максимальное число строк для сбора ключей (для построения фильтра)
spark.conf.set("spark.sql.runtime.filter.bloomFilter.maxNumItems", "4000000")
# Порог для Broadcast (если probe side > порога, BF не строится)
spark.conf.set("spark.sql.runtime.filter.semiJoinReductionThreshold", "0.5")
# Adaptive Query Execution должен быть включён для Runtime Filtering
spark.conf.set("spark.sql.adaptive.enabled", "true")
Посмотреть план с Runtime Filter можно так:
df_large = spark.table("silver.clickstream") # 2 млрд строк
df_small = spark.table("silver.valid_sessions") # 10 млн строк
result = df_large.join(df_small, on="session_id", how="inner")
result.explain(mode="formatted")
В плане вы увидите BloomFilterAggregate и BloomFilterMightContain:
== Physical Plan ==
BroadcastHashJoin [session_id], [session_id], Inner, BuildRight, false
:- Filter isnotnull(session_id)
: +- BloomFilterMightContain(bloomFilter=<BF>, value=session_id) ← применение
: +- FileScan parquet silver.clickstream [session_id, ...]
+- BroadcastExchange HashedRelationBroadcastMode(...)
+- BloomFilterAggregate(estimatedNumItems=10000000, numBits=94906265) ← построение
+- FileScan parquet silver.valid_sessions [session_id, ...]
Часть 3. Bloom Filter vs другие методы оптимизации JOIN¶
Прежде чем погружаться в код, важно понять, где Bloom Filter занимает своё место среди других инструментов оптимизации Spark.
3.1. Сравнительная матрица¶
| Метод | Когда применять | Ограничения |
|---|---|---|
| Broadcast Hash Join | Probe side < 10–100 МБ | Маленькая таблица должна помещаться в RAM экзекьютора |
| Partition Pruning | Данные партиционированы по ключу JOIN | Требует физического партиционирования таблицы |
| Bucketing | Регулярные JOIN по одному ключу | Требует предварительного bucketing обеих таблиц |
| Z-Ordering | Диапазонные запросы по нескольким колонкам | Только file-level skipping, не Shuffle reduction |
| Bloom Filter | Probe side 100 МБ - 10 ГБ, high cardinality | False positives, memory для фильтра |
| Sort-Merge Join | Обе стороны большие, нет других вариантов | Самый медленный, но работает всегда |
3.2. Когда Bloom Filter выигрывает¶
Bloom Filter особенно эффективен в комбинации «большая таблица + средняя таблица с высокой кардинальностью ключей»:
Высокая кардинальность ключей - критичный фактор. Если в большой таблице есть 2 млрд строк по 100 млн уникальных session_id, и в маленькой - 10 млн session_id, то только 10 из 100 = 10% строк большой таблицы совпадут. Bloom Filter отсеет 90% до Shuffle - это огромный выигрыш.
Если же кардинальность низкая (например, в big join нет department_id - 50 значений), большинство строк и так совпадёт, и BF почти ничего не отсечёт.
3.3. Bloom Filter как pre-filter перед Anti-Join¶
Классический паттерн для инкрементальной дедупликации:
Без BF: Инкремент (50 млн) Anti-Join History (2 млрд) = дорого
С BF: Инкремент (50 млн) → BF pre-filter → 5 млн → Anti-Join History (2 млрд)
При этом History тоже меньше нужно читать (за счёт file skipping)
Bloom Filter здесь работает как «первый барьер» - быстрый, вероятностный. За ним стоит «второй барьер» - медленный, точный (Anti-Join). Двухуровневая фильтрация кратно эффективнее одноуровневой точной.
Часть 4. Практика PySpark: явное управление Bloom Filter через API¶
4.1. stat.bloomFilter() - создание фильтра из DataFrame¶
PySpark предоставляет встроенный метод DataFrame.stat.bloomFilter() для создания фильтра из значений указанной колонки:
from pyspark.sql import SparkSession, functions as F
spark = SparkSession.builder.appName("bloom_filter_demo").getOrCreate()
# Читаем историческую таблицу (данные за 90 дней)
df_history = spark.table("silver.transactions").filter(
F.col("transaction_date") >= F.date_sub(F.current_date(), 90)
)
# Создаём Bloom Filter по ключевой колонке
# Параметры:
# colName - колонка, по которой строится фильтр
# expectedNumItems - ожидаемое число уникальных значений в колонке
# fpp - допустимый False Positive Rate (0.0 до 1.0)
bloom_filter = df_history.stat.bloomFilter(
colName="transaction_id",
expectedNumItems=500_000_000, # 500M уникальных transaction_id
fpp=0.01 # FPR = 1%
)
print(f"Bloom Filter object: {bloom_filter}")
print(f"Expected bits: ~{500_000_000 * 9.6 / 8 / 1024 / 1024:.0f} МБ")
Важные параметры:
expectedNumItems: переоцените, а не недооцените. Если реальных элементов окажется больше - FPR вырастет. Лучше взять с запасом 20–30%.fpp: значение 0.01 (1%) - хороший баланс. 0.1 - фильтр пропустит каждую 10-ю лишнюю строку. 0.001 - нужно в 1.5 раза больше памяти.
4.2. Применение фильтра через UDF¶
После создания фильтра его нужно применить к DataFrame инкремента. Самый надёжный способ - через UDF с Broadcast:
from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType
# Делаем Broadcast: Spark рассылает объект фильтра на все экзекьюторы
# Это критично - без broadcast каждый экзекьютор создаст свою копию (дорого)
bloom_filter_broadcast = spark.sparkContext.broadcast(bloom_filter)
@udf(returnType=BooleanType())
def bloom_filter_udf(value):
"""
UDF: проверяет, может ли значение находиться в историческом наборе.
Возвращает:
True - значение ВОЗМОЖНО есть в истории (нужна точная проверка)
False - значение ТОЧНО отсутствует в истории (можно вставлять)
ВАЖНО: False Negative невозможен - мы никогда не потеряем новую строку.
"""
if value is None:
# NULL-ключи нельзя хэшировать корректно - пропускаем их дальше
return True
return bloom_filter_broadcast.value.mightContain(value)
# Читаем инкрементальный батч
df_increment = spark.table("bronze.transactions_incoming")
# ШАГ 1: Bloom Filter pre-filter - быстро, вероятностно
# Оставляем только строки, которые МОГУТ быть новыми (+ FPR мусор)
df_maybe_new = df_increment.filter(~bloom_filter_udf(F.col("transaction_id")))
# Знак ~ (NOT): нас интересуют строки, которых НЕТ в истории
# bloom_filter_udf возвращает True = "есть в истории" → мы их отбрасываем
# bloom_filter_udf возвращает False = "точно новые" → мы их оставляем
print(f"Инкремент до BF: {df_increment.count():,}")
print(f"После BF pre-filter: {df_maybe_new.count():,}")
4.3. Двухуровневая фильтрация: BF + точный Anti-Join¶
После Bloom Filter pre-filter применяем точный Anti-Join для устранения False Positives:
# ШАГ 2: Точный Anti-Join для устранения False Positives
# На этом этапе df_maybe_new уже в 10–20 раз меньше исходного инкремента
# → Shuffle в Anti-Join в 10–20 раз дешевле
# Берём только ключи из истории для Anti-Join (минимизируем читаемые данные)
history_keys = (
spark.table("silver.transactions")
.filter(F.col("transaction_date") >= F.date_sub(F.current_date(), 90))
.select("transaction_id")
.distinct()
)
# Финальный Anti-Join: гарантированно новые строки
df_truly_new = df_maybe_new.join(
history_keys,
on="transaction_id",
how="left_anti"
)
print(f"Гарантированно новые строки: {df_truly_new.count():,}")
4.4. Измерение False Positive Rate¶
Чтобы понять, насколько эффективен ваш фильтр, можно измерить реальный FPR:
# Шаг 1: Сколько строк прошло BF pre-filter?
count_after_bf = df_maybe_new.count()
# Шаг 2: Сколько из них - настоящие новые (после точного Anti-Join)?
count_truly_new = df_truly_new.count()
# Шаг 3: Разница - это False Positives (строки уже есть в истории, но BF пропустил)
false_positives = count_after_bf - count_truly_new
fpr_actual = false_positives / count_after_bf if count_after_bf > 0 else 0
print(f"Строк после BF: {count_after_bf:,}")
print(f"Истинно новых строк: {count_truly_new:,}")
print(f"False positives: {false_positives:,}")
print(f"Фактический FPR: {fpr_actual:.3%}")
# Строки инкремента, отсечённые BF (True Negatives - гарантированно верно)
count_original = df_increment.count()
true_negatives_by_bf = count_original - count_after_bf
print(f"Отсечено фильтром (TN): {true_negatives_by_bf:,}")
print(f"Эффективность pre-filter: {true_negatives_by_bf / count_original:.1%}")
4.5. Сериализация фильтра для переиспользования¶
В production-пайплайнах фильтр часто нужно сохранить между запусками (например, строить раз в час, применять каждую минуту):
import io
from pyspark.ml.feature import BucketedRandomProjectionLSH # не нужен, просто пример импорта
# Сериализация фильтра в байты
def serialize_bloom_filter(bf) -> bytes:
"""Сохраняет Bloom Filter в байтовый массив для хранения в HDFS/S3."""
buffer = io.BytesIO()
# Spark BloomFilter имеет метод writeTo(OutputStream)
# Используем Java-обёртку через py4j
bf._jvm_bf.writeTo(buffer) # псевдокод - в реальности нужен Java IO bridge
return buffer.getvalue()
def deserialize_bloom_filter(spark: SparkSession, data: bytes):
"""Загружает Bloom Filter из байтового массива."""
buffer = io.BytesIO(data)
return spark._jvm.org.apache.spark.util.sketch.BloomFilter.readFrom(buffer)
# Практичный паттерн: кэшируем фильтр в памяти между запусками
import pickle
# ВАРИАНТ A: Pickle через pyspark util (работает с Python Bloom Filter)
# ВАРИАНТ B: Хранить в Delta/Iceberg-таблице как BINARY колонку
# ВАРИАНТ C: Пересоздавать фильтр каждый запуск (чаще всего - самый простой)
# Пересоздание фильтра - наиболее надёжный подход в production
def build_bloom_filter(spark, table: str, key_col: str, days: int = 90, fpp: float = 0.01):
"""
Строит свежий Bloom Filter из ключей таблицы за последние N дней.
Рекомендуется запускать при старте пайплайна или раз в час
для поддержания актуальности фильтра.
"""
df = (
spark.table(table)
.filter(F.col("transaction_date") >= F.date_sub(F.current_date(), days))
.select(key_col)
.dropna()
.distinct()
)
# Считаем количество уникальных ключей для правильного sizing
n_unique = df.count()
expected_items = int(n_unique * 1.3) # Запас 30% для новых данных
bf = df.stat.bloomFilter(key_col, expectedNumItems=expected_items, fpp=fpp)
print(f"BF built: {n_unique:,} unique keys, expected={expected_items:,}, fpp={fpp}")
return bf
Часть 5. Probabilistic deduplication для incremental ETL¶
5.1. Паттерн: Multi-layer deduplication¶
В реальных ETL-пайплайнах дедупликация имеет несколько уровней. Bloom Filter занимает первый, самый быстрый уровень:
5.2. Полная реализация паттерна¶
from pyspark.sql import SparkSession, functions as F, DataFrame
from pyspark.sql.functions import udf
from pyspark.sql.types import BooleanType
from typing import Optional
import logging
logger = logging.getLogger(__name__)
class BloomFilterDeduplicator:
"""
Трёхуровневый дедупликатор для инкрементальных ETL-пайплайнов.
Уровни:
1. Intra-batch dedup (dropDuplicates)
2. Bloom Filter pre-filter (вероятностный)
3. Точный Anti-Join (детерминированный)
Гарантии:
- False Negative невозможен: ни одна реально новая строка не будет потеряна
- False Positive возможен: небольшой % уже существующих строк пройдёт в Anti-Join
"""
def __init__(
self,
spark: SparkSession,
history_table: str,
key_col: str,
lookback_days: int = 90,
bloom_fpp: float = 0.01,
bloom_expected_items: Optional[int] = None,
):
self.spark = spark
self.history_table = history_table
self.key_col = key_col
self.lookback_days = lookback_days
self.bloom_fpp = bloom_fpp
self._bloom_filter = None
self._bloom_broadcast = None
self._bloom_expected_items = bloom_expected_items
def build_bloom_filter(self) -> None:
"""
Строит Bloom Filter из исторических данных.
Вызывается один раз при инициализации пайплайна.
Читает только ключевую колонку за последние lookback_days —
минимальный IO при построении.
"""
logger.info(f"Building Bloom Filter from {self.history_table} "
f"(last {self.lookback_days} days)...")
df_keys = (
self.spark.table(self.history_table)
.filter(
F.col("event_date") >= F.date_sub(F.current_date(), self.lookback_days)
)
.select(self.key_col)
.dropna(subset=[self.key_col])
.distinct()
)
# Если expectedNumItems не задан - считаем точно (extra action)
if self._bloom_expected_items is None:
n = df_keys.count()
expected = int(n * 1.3) # +30% запас для роста данных
logger.info(f"Counted {n:,} unique keys, setting expectedNumItems={expected:,}")
else:
expected = self._bloom_expected_items
self._bloom_filter = df_keys.stat.bloomFilter(
colName=self.key_col,
expectedNumItems=expected,
fpp=self.bloom_fpp,
)
# Broadcast: рассылаем фильтр на все экзекьюторы один раз
self._bloom_broadcast = self.spark.sparkContext.broadcast(self._bloom_filter)
logger.info(f"Bloom Filter built and broadcast. FPP={self.bloom_fpp}")
def _make_udf(self):
"""Создаёт UDF, замыкающий broadcast-переменную."""
bf_bc = self._bloom_broadcast
@udf(returnType=BooleanType())
def might_exist_in_history(value) -> bool:
"""True = строка ВОЗМОЖНО уже есть в истории."""
if value is None:
return False # NULL-ключи всегда считаем новыми (вставляем)
return bf_bc.value.mightContain(value)
return might_exist_in_history
def deduplicate(self, df_increment: DataFrame) -> DataFrame:
"""
Выполняет трёхуровневую дедупликацию инкремента.
Returns:
DataFrame с гарантированно новыми строками
(небольшой % False Positives устранён точным Anti-Join)
"""
if self._bloom_filter is None:
raise RuntimeError("Call build_bloom_filter() first!")
key = self.key_col
# ── Уровень 1: Intra-batch dedup ─────────────────────────────────────
# Убираем дубли внутри самого инкремента.
# При дублях берём строку с максимальным event_timestamp
# (или любым другим tiebreaker-критерием).
df_intra_deduped = df_increment.dropDuplicates([key])
logger.info("Level 1: intra-batch dedup done")
# ── Уровень 2: Bloom Filter pre-filter ───────────────────────────────
# Быстро отсекаем строки, которых точно нет в истории.
# Остаются: (реально новые) + (FPR% строк, уже есть в истории).
might_exist_udf = self._make_udf()
# might_exist_udf(key) = True → возможно уже есть → ОТБРОСИТЬ (дубль?)
# might_exist_udf(key) = False → точно новое → ОСТАВИТЬ
df_bf_filtered = df_intra_deduped.filter(
~might_exist_udf(F.col(key))
)
logger.info("Level 2: Bloom Filter pre-filter applied")
# ── Уровень 3: Точный Anti-Join ──────────────────────────────────────
# Устраняем False Positives. На этом этапе df_bf_filtered уже
# в 10–20 раз меньше исходного инкремента → Shuffle минимален.
history_keys = (
self.spark.table(self.history_table)
.filter(
F.col("event_date") >= F.date_sub(F.current_date(), self.lookback_days)
)
.select(key)
.dropna()
)
df_truly_new = df_bf_filtered.join(
history_keys,
on=key,
how="left_anti",
)
logger.info("Level 3: exact anti-join done")
return df_truly_new
# ── Использование ──────────────────────────────────────────────────────────
def run_incremental_ingestion(spark: SparkSession) -> None:
"""Основной ETL-пайплайн с Bloom Filter дедупликацией."""
dedup = BloomFilterDeduplicator(
spark=spark,
history_table="silver.clickstream",
key_col="event_id",
lookback_days=90,
bloom_fpp=0.01,
bloom_expected_items=2_000_000_000, # 2B строк в истории за 90 дней
)
# Строим фильтр один раз при старте (займёт 2–5 минут для 2B строк)
dedup.build_bloom_filter()
# Обрабатываем инкремент
df_incoming = spark.table("bronze.clickstream_incoming")
df_new = dedup.deduplicate(df_incoming)
# Записываем в Silver
df_new.writeTo("silver.clickstream").append()
logger.info(f"Written {df_new.count():,} new rows to silver.clickstream")
5.3. Обработка False Positives в downstream¶
Важно понимать последствия False Positives для различных downstream-потребителей:
# Проверка: насколько велик реальный FPR?
def measure_dedup_accuracy(
df_increment: DataFrame,
df_result: DataFrame,
df_bf_only: DataFrame,
key_col: str,
) -> dict:
"""
Измеряет точность Bloom Filter фильтрации.
Returns:
dict с метриками: фактический FPR, эффективность pre-filter, и т.д.
"""
n_original = df_increment.count()
n_after_bf = df_bf_only.count()
n_final = df_result.count()
n_true_new = n_final # Гарантированно новые (после точного AJ)
n_false_positives = n_after_bf - n_final # Строки, прошедшие BF, но отсечённые AJ
n_filtered_by_bf = n_original - n_after_bf # True Negatives от BF
actual_fpr = n_false_positives / n_after_bf if n_after_bf > 0 else 0
bf_efficiency = n_filtered_by_bf / n_original if n_original > 0 else 0
return {
"original_count": n_original,
"after_bf_count": n_after_bf,
"final_new_count": n_final,
"false_positives": n_false_positives,
"true_negatives_by_bf": n_filtered_by_bf,
"actual_fpr": actual_fpr,
"bf_pre_filter_efficiency": bf_efficiency,
}
Часть 6. False Positives: цена вероятностной оптимизации¶
6.1. Природа False Positives¶
False Positive возникает, когда элемент X не был добавлен в Bloom Filter, но все k хэш-функций указывают на биты, которые были установлены другими элементами. Это называется битовой коллизией.
Чем больше элементов добавлено в фильтр (чем выше заполненность массива), тем больше вероятность случайного совпадения всех k битов - тем выше FPR.
Наглядный пример: фильтр построен на 1B элементах при расчёте на 500M (переполнение в 2x):
# Заполнение фильтра сверх нормы резко увеличивает FPR
# Формула: FPR = (1 - e^(-kn/m))^k
import math
def calculate_fpr(m_bits: int, n_elements: int, k_functions: int) -> float:
"""Вычисляет теоретический FPR Bloom Filter."""
exponent = -k_functions * n_elements / m_bits
fpr = (1 - math.exp(exponent)) ** k_functions
return fpr
# Пример: фильтр рассчитан на 500M элементов с FPR=1%
m = int(500_000_000 * 9.6) # ~4.8 миллиарда бит (~600 МБ)
k = 7 # оптимальное число хэш-функций
print("Реальный FPR при разном заполнении:")
for n_actual in [250_000_000, 500_000_000, 750_000_000, 1_000_000_000]:
fpr = calculate_fpr(m, n_actual, k)
overfill = n_actual / 500_000_000
print(f" n={n_actual/1e6:.0f}M ({overfill:.1f}x расчётного): FPR = {fpr:.3%}")
Ожидаемый вывод:
Реальный FPR при разном заполнении:
n=250M (0.5x расчётного): FPR = 0.001% (лучше расчётного)
n=500M (1.0x расчётного): FPR = 1.000% (соответствует)
n=750M (1.5x расчётного): FPR = 3.681% (деградация)
n=1000M (2.0x расчётного): FPR = 8.865% (серьёзная деградация)
Вывод: всегда задавайте expectedNumItems с запасом 20–30% от реального числа элементов. Это увеличит размер фильтра, но защитит от деградации FPR при росте данных.
6.2. Влияние FPR на производительность downstream¶
False Positive означает, что строка, уже существующая в истории, пройдёт BF pre-filter и попадёт в точный Anti-Join. Это незначительный overhead для Anti-Join, но полная безопасность для данных:
# Количество False Positives при разных FPR и объёмах
# Допустим, инкремент = 10M строк, из них реально новых = 1M (10% новых)
# → В истории есть 9M строк из этих 10M
n_increment = 10_000_000
n_truly_new = 1_000_000 # 10% новых
n_existing = 9_000_000 # 90% уже есть в истории
for fpr in [0.001, 0.01, 0.05, 0.10]:
# BF пропустит: все 1M новых (True Negatives правильно пропустит)
# + FPR% от 9M существующих = False Positives
false_positives = int(n_existing * fpr)
total_after_bf = n_truly_new + false_positives
print(f"FPR={fpr:.1%}: {false_positives:,} False Positives, "
f"{total_after_bf:,} строк идут в Anti-Join "
f"(из {n_increment:,} исходных)")
Ожидаемый вывод:
FPR=0.1%: 9,000 False Positives, 1,009,000 строк идут в Anti-Join (из 10,000,000)
FPR=1.0%: 90,000 False Positives, 1,090,000 строк идут в Anti-Join (из 10,000,000)
FPR=5.0%: 450,000 False Positives, 1,450,000 строк идут в Anti-Join (из 10,000,000)
FPR=10.0%: 900,000 False Positives, 1,900,000 строк идут в Anti-Join (из 10,000,000)
Даже при FPR=10% Anti-Join обрабатывает 1.9M строк вместо 10M - ускорение в 5x. При FPR=1% - 1.09M строк - ускорение почти в 10x.
Часть 7. Bloom Filter в Lakehouse форматах: Delta Lake и Iceberg¶
7.1. Bloom Filter Index в Delta Lake¶
Delta Lake поддерживает Bloom Filter Index - фильтры хранятся внутри Parquet-файлов как часть метаданных. При чтении Spark проверяет фильтр каждого файла и пропускает нерелевантные.
Это file-level skipping - принципиально другой уровень, чем row-level pre-filter в памяти.
# Создание таблицы с Bloom Filter Index в Delta Lake
spark.sql("""
CREATE TABLE silver.clickstream (
event_id STRING NOT NULL,
session_id STRING,
user_id BIGINT,
event_type STRING,
event_date DATE,
event_ts TIMESTAMP,
page_url STRING,
properties MAP<STRING, STRING>
)
USING delta
PARTITIONED BY (event_date)
TBLPROPERTIES (
-- Включаем Bloom Filter для высококардинальных колонок
'delta.bloomFilterIndex.enabled' = 'true',
'delta.bloomFilterIndexOnColumn.event_id.enabled' = 'true',
'delta.bloomFilterIndexOnColumn.event_id.fpp' = '0.01',
'delta.bloomFilterIndexOnColumn.session_id.enabled' = 'true',
'delta.bloomFilterIndexOnColumn.session_id.fpp' = '0.01'
)
""")
# Добавление Bloom Filter Index к существующей таблице Delta
spark.sql("""
ALTER TABLE silver.clickstream
SET TBLPROPERTIES (
'delta.bloomFilterIndex.enabled' = 'true',
'delta.bloomFilterIndexOnColumn.event_id.enabled' = 'true',
'delta.bloomFilterIndexOnColumn.event_id.fpp' = '0.01'
)
""")
# ВАЖНО: Bloom Filter Index применяется только к НОВЫМ записываемым файлам.
# Для существующих файлов нужен OPTIMIZE/REWRITE:
spark.sql("OPTIMIZE silver.clickstream ZORDER BY (event_date, user_id)")
# После OPTIMIZE новые файлы будут содержать BF Index
7.2. Bloom Filter Index в Apache Iceberg¶
Iceberg хранит Bloom Filter как часть Column Statistics на уровне файла. Настройка через TBLPROPERTIES:
# Создание Iceberg таблицы с Bloom Filter
spark.sql("""
CREATE TABLE catalog.silver.clickstream (
event_id STRING NOT NULL,
session_id STRING,
user_id BIGINT,
event_type STRING,
event_date DATE,
event_ts TIMESTAMP
)
USING iceberg
PARTITIONED BY (days(event_ts))
TBLPROPERTIES (
-- Bloom Filter для event_id (высокая кардинальность)
'write.parquet.bloom-filter-enabled.column.event_id' = 'true',
'write.parquet.bloom-filter-fpp.column.event_id' = '0.01',
-- Bloom Filter для session_id
'write.parquet.bloom-filter-enabled.column.session_id' = 'true',
'write.parquet.bloom-filter-fpp.column.session_id' = '0.01',
-- Bloom Filter для user_id
'write.parquet.bloom-filter-enabled.column.user_id' = 'true',
'write.parquet.bloom-filter-fpp.column.user_id' = '0.01'
)
""")
# Проверка работы Bloom Filter через Data Skipping статистику
spark.sql("""
SELECT
operation,
numTargetFilesAdded,
numTargetFilesRemoved,
operationParameters
FROM silver.clickstream.history
ORDER BY version DESC
LIMIT 10
""").show(truncate=False)
7.3. Механика File-Level Skipping¶
Когда Spark выполняет запрос с предикатом WHERE event_id = 'abc123', происходит следующее:
Min/Max statistics - дешёвый первый барьер: если event_id = 'abc123' вне диапазона [Min, Max] файла - пропускаем. Bloom Filter - второй барьер: даже если значение попадает в диапазон Min/Max, BF может сказать «нет этого значения».
Вместе они дают высокую эффективность Data Skipping для точечных lookup-запросов.
7.4. Сравнение: in-memory BF vs file-level BF¶
| Критерий | in-memory stat.bloomFilter() | File-level BF Index (Delta/Iceberg) |
|---|---|---|
| Когда применяется | Во время выполнения запроса (runtime) | При чтении файлов (scan phase) |
| Что отсекает | Строки в DataFrame | Целые Parquet файлы |
| Для чего лучше | Incremental dedup, custom ETL | Point lookups, MERGE ON-условия |
| Требует build step | Да (отдельная Action) | Нет (встроен в write) |
| Актуальность | Только для заданного набора данных | Всегда актуален для файла |
| Стоимость build | O(n) scan исторической таблицы | 0 (автоматически при write) |
Часть 8. Анализ Physical Plan: видим Bloom Filter в действии¶
8.1. Что смотреть в Spark UI¶
После включения Runtime Bloom Filter Spark показывает его в Physical Plan. Разберём ключевые узлы:
# Пример запроса для анализа
df_events = spark.table("silver.clickstream") # 2 млрд строк
df_sessions = spark.table("silver.valid_sessions") # 10 млн строк
df_joined = df_events.join(df_sessions, on="session_id", how="inner")
df_joined.explain(mode="formatted")
Типичный физический план с включённым Runtime Bloom Filter:
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- BroadcastHashJoin [session_id#12], [session_id#45], Inner, BuildRight, false
:- Filter isnotnull(session_id#12)
: +- Filter might_contain(bloomFilter#23, session_id#12) ← BF применяется
: +- FileScan parquet silver.clickstream [...] ← до scan
: PushedFilters: [IsNotNull(session_id)]
+- BroadcastExchange HashedRelationBroadcastMode(...)
+- BloomFilterAggregate(
estimatedNumItems=10000000,
numBits=94906265) ← BF строится
+- FileScan parquet silver.valid_sessions [session_id]
Ключевые узлы для анализа:
BloomFilterAggregate- здесь Spark строит фильтр по маленькой таблицеmight_contain(bloomFilter, ...)- здесь фильтр применяется к большой таблицеPushedFilters- предикаты, переданные вниз к Parquet readerFileScan- если после него стоит BF фильтр, Spark применяет его до передачи данных наверх
8.2. Spark UI: метрики Shuffle Reduction¶
В Spark UI (вкладка SQL → физический план → Statistics) смотрите:
Stage N: Exchange (Shuffle)
Shuffle Write: 12.5 GB ← данные, отправленные на shuffle
Shuffle Read: 11.8 GB ← данные, полученные после shuffle
Records Written: 47,832,091
Сравните с запуском без BF:
Stage N: Exchange (Shuffle) - без BF
Shuffle Write: 187.3 GB ← ~15x больше!
Shuffle Read: 183.1 GB
Records Written: 743,128,432
Bloom Filter сократил shuffle write с 187 ГБ до 12.5 ГБ - в 15 раз.
8.3. Принудительное включение и отключение¶
# Форсируем применение Bloom Filter (если AQE его не выбрал автоматически)
spark.conf.set("spark.sql.runtime.filter.enabled", "true")
spark.conf.set("spark.sql.runtime.filter.bloomFilter.enabled", "true")
# Увеличиваем лимиты (для больших probe side)
spark.conf.set(
"spark.sql.runtime.filter.bloomFilter.maxNumBits",
str(8 * 1024 * 1024 * 1024) # 8 ГБ в битах = 8_589_934_592
)
spark.conf.set(
"spark.sql.runtime.filter.bloomFilter.maxNumItems",
"100000000" # До 100M элементов в одном BF
)
# Отключение для диагностики (сравнение с и без BF)
spark.conf.set("spark.sql.runtime.filter.bloomFilter.enabled", "false")
df_no_bf = df_large.join(df_small, on="key", how="inner")
df_no_bf.explain(mode="formatted") # Увидим обычный SortMergeJoin или BHJ
Часть 9. Performance internals: память и Shuffle Reduction¶
9.1. Почему Bloom Filter так компактен¶
Hash Set из BIGINT значений требует:
BIGINT: 8 байт/элемент- Hash Set overhead (load factor ~0.75): 8 / 0.75 ≈ 10.7 байт/элемент
- Для 1 млрд элементов: ~10.7 ГБ
Bloom Filter при FPR=1%:
- Нужно
9.6 бит/элемент= 1.2 байт/элемент - Для 1 млрд элементов: ~1.2 ГБ
- Экономия: в ~9x меньше памяти
При FPR=0.1%: 14.4 бит/элемент = 1.8 байт/элемент → ~1.8 ГБ для 1B элементов - всё равно в 6x компактнее HashSet.
9.2. Когда Bloom Filter не умещается в Broadcast¶
Если фильтр превышает spark.sql.autoBroadcastJoinThreshold (по умолчанию 10 МБ для BHJ, но для Runtime BF настройка отдельная), Spark не будет делать Broadcast и не применит BF автоматически. Решения:
# Увеличить порог для broadcast BF объекта
spark.conf.set("spark.sql.runtime.filter.bloomFilter.maxNumBits", str(4 * 1024**3 * 8)) # 4 ГБ
# Или построить частичный BF (только по последним N дням)
# вместо всей истории
bf_partial = (
spark.table("silver.events")
.filter(F.col("event_date") >= F.date_sub(F.current_date(), 30)) # 30 дней вместо 90
.stat.bloomFilter("event_id", expectedNumItems=500_000_000, fpp=0.01)
)
Практическое правило: Bloom Filter объект при FPR=1% занимает ~1.2 байт на элемент. Для Broadcast лучше держать его в пределах 200–500 МБ. Это ~170–420 млн элементов. Если истории больше - используйте file-level BF Index в Delta/Iceberg.
9.3. CPU overhead: хэш-вычисления¶
Проверка принадлежности в BF требует k хэш-вычислений на строку. При k=7 и 10 млн строк инкремента - это 70 млн хэш-операций. Каждая очень быстрая (murmur3), но в сумме это несколько секунд CPU.
Сравнение: Anti-Join по 2B историческим строкам - это Shuffle 200+ ГБ + сравнение 2 × 10^18 пар ключей. BF overhead - ничто по сравнению с этим.
# Benchmark: время BF проверки vs время Anti-Join
import time
# BF pre-filter (только вычисление хэшей, без Shuffle)
start = time.time()
count_after_bf = df_increment.filter(~might_exist_udf(F.col("event_id"))).count()
bf_time = time.time() - start
# Полный Anti-Join (с Shuffle)
start = time.time()
count_exact = df_increment.join(history_keys, on="event_id", how="left_anti").count()
exact_time = time.time() - start
print(f"BF pre-filter time: {bf_time:.1f}s (→ {count_after_bf:,} строк)")
print(f"Exact Anti-Join time: {exact_time:.1f}s (→ {count_exact:,} строк)")
print(f"Speedup (BF vs exact): {exact_time / bf_time:.1f}x")
Часть 10. Анти-паттерны и production-проблемы¶
10.1. Низкая кардинальность ключей¶
Bloom Filter эффективен только при высокой кардинальности. На низкой кардинальности он бесполезен или вреден:
# ПЛОХО: BF на колонке с 50 уникальными значениями (низкая кардинальность)
# department_id: ["HR", "IT", "Finance", ..., "Legal"] - 50 значений
# BF из 50 значений: практически всегда говорит "есть"
# → FPR близок к 100% → BF ничего не отсекает
bad_bf = df_small.stat.bloomFilter("department_id", expectedNumItems=50, fpp=0.01)
# Фильтр размером ~1 КБ, но с крошечным эффектом
# ХОРОШО: BF на высококардинальной колонке
good_bf = df_small.stat.bloomFilter("user_id", expectedNumItems=50_000_000, fpp=0.01)
# Фильтр размером ~60 МБ, эффективно отсекает 95%+ строк
Правило: Bloom Filter целесообразен, когда процент совпадающих строк между большим и малым датасетами менее 20–30%. При более высоком overlap проще использовать обычный JOIN или Broadcast Join.
10.2. Несвежий (Stale) Bloom Filter¶
Если вы кэшируете фильтр между запусками, он устаревает: новые данные появляются в истории, но фильтр о них не знает. Это не вызывает ошибок данных (новые данные в истории будут пропущены через BF как потенциально новые + отфильтрованы точным AJ), но снижает эффективность pre-filter.
# Несвежий фильтр: пропускает больше False Positives
# Решение: пересоздавать фильтр при каждом запуске пайплайна
# или хранить TTL и проверять его
from datetime import datetime, timedelta
class CachedBloomFilter:
"""BF с временем жизни (TTL)."""
def __init__(self, ttl_minutes: int = 60):
self._filter = None
self._built_at = None
self._ttl = timedelta(minutes=ttl_minutes)
def is_fresh(self) -> bool:
if self._filter is None or self._built_at is None:
return False
return datetime.now() - self._built_at < self._ttl
def get_or_build(self, spark, table: str, key_col: str) -> any:
if not self.is_fresh():
self._filter = _build_filter(spark, table, key_col)
self._built_at = datetime.now()
return self._filter
10.3. Слишком агрессивный FPP = слишком маленький размер¶
Установка очень маленького expectedNumItems при большом фактическом числе элементов катастрофически увеличивает FPR:
# ОПАСНО: указали 10M, реальных элементов 500M
# → FPR деградирует с 1% до ~30%+
bad_bf = df_large.stat.bloomFilter(
"event_id",
expectedNumItems=10_000_000, # 10M
fpp=0.01
)
# При реальных 500M элементов: FPR ≈ 30–40% (практически бесполезен)
# ПРАВИЛЬНО: с запасом 30%
actual_n = df_large.select("event_id").distinct().count() # Считаем точно
good_bf = df_large.stat.bloomFilter(
"event_id",
expectedNumItems=int(actual_n * 1.3),
fpp=0.01
)
10.4. Использование BF там, где лучше Broadcast Join¶
Если меньшая таблица помещается в память каждого экзекьютора (~100–200 МБ), используйте Broadcast Hash Join напрямую - он точный (без False Positives) и быстрее BF:
# ХОРОШО: маленькая таблица помещается в broadcast
small_table_size_mb = 50 # МБ
if small_table_size_mb < 100:
# Broadcast Hash Join - точный, без FPR, без overhead построения BF
result = df_large.join(F.broadcast(df_small), on="key", how="inner")
else:
# Bloom Filter оправдан
bf = df_small.stat.bloomFilter("key", expectedNumItems=10_000_000, fpp=0.01)
# ... применяем BF
10.5. Unstable hashing после Schema Evolution¶
Если схема таблицы изменилась (добавлена колонка в начало - нестандартная схема), хэш-функции внутри Bloom Filter могут давать другие результаты для тех же значений. В Delta/Iceberg это не проблема (BF встроен в файл и применяется к конкретным значениям колонки). Но при хранении сериализованного BF объекта между версиями кода - потенциальный источник ошибок.
Решение: всегда пересоздавать BF при изменении схемы или использовать file-level BF Index.
10.6. Сводная матрица анти-паттернов¶
Часть 11. End-to-End Production Кейс: ускорение clickstream dedup через Bloom Filter¶
11.1. Постановка задачи¶
Сценарий: крупный интернет-ритейлер. Klikstream (события браузера и приложений) собирается с 200 млн MAU. Каждый час приходит батч - 200 млн событий. Исторический Silver датасет - 50 млрд строк за 2 года.
Требования:
- Инкрементальная дедупликация: новые события добавить в Silver, дубли пропустить
- Время обработки батча: не более 30 минут
- Кластер: 100 экзекьюторов × 32 ГБ RAM
- Гарантии: ни одно новое событие не должно быть потеряно
Проблема: классический Anti-Join по 50 млрд строк - это Shuffle на ~400 ГБ трафика и 3–4 часа времени. Нужно решение.
11.2. Архитектура решения¶
11.3. Полная реализация¶
"""
clickstream_dedup_pipeline.py
Production-кейс: Bloom Filter ускорение дедупликации clickstream
200M строк/час, 50B строк история
"""
from pyspark.sql import SparkSession, functions as F, DataFrame
from datetime import datetime
import logging
import time
logger = logging.getLogger(__name__)
# ── Конфигурация ──────────────────────────────────────────────────────────────
BRONZE_TABLE = "catalog.bronze.clickstream_incoming"
SILVER_TABLE = "catalog.silver.clickstream"
KEY_COL = "event_id"
DATE_COL = "event_date"
LOOKBACK_DAYS = 90 # Насколько далеко смотреть в историю для Anti-Join
BLOOM_FPP = 0.01 # 1% FPR для file-level BF Index
def configure_spark(spark: SparkSession) -> SparkSession:
"""
Настраиваем Spark для максимальной эффективности с BF.
Ключевые настройки:
- AQE: адаптивное выполнение (перепланирование после каждого stage)
- Runtime Filter: автоматическое применение BF в JOIN
- Shuffle partitions: ~400 для 50 ГБ shuffle при 128 МБ/партицию
"""
configs = {
# AQE - критично для Runtime Bloom Filter
"spark.sql.adaptive.enabled": "true",
"spark.sql.adaptive.coalescePartitions.enabled": "true",
"spark.sql.adaptive.skewJoin.enabled": "true",
# Runtime Bloom Filter
"spark.sql.runtime.filter.enabled": "true",
"spark.sql.runtime.filter.bloomFilter.enabled": "true",
"spark.sql.runtime.filter.bloomFilter.maxNumBits": str(8 * 1024**3 * 8),
"spark.sql.runtime.filter.bloomFilter.maxNumItems": "100000000",
# Shuffle
"spark.sql.shuffle.partitions": "400",
# Broadcast threshold
"spark.sql.autoBroadcastJoinThreshold": "104857600", # 100 МБ
}
for k, v in configs.items():
spark.conf.set(k, v)
return spark
# ── Stage 1: Нормализация и intra-batch dedup ─────────────────────────────────
def stage1_normalize_and_intra_dedup(spark: SparkSession) -> DataFrame:
"""
Читает бронзовый батч, нормализует и дедуплицирует внутри батча.
Intra-batch dedup выполняется через dropDuplicates по event_id.
При дублях внутри батча (один источник прислал дважды) берём
строку с максимальным event_ts.
"""
t = time.time()
df_raw = spark.table(BRONZE_TABLE)
# Нормализация
df_norm = (
df_raw
.withColumn(KEY_COL, F.trim(F.col(KEY_COL))) # Убираем пробелы
.withColumn("event_ts", F.to_utc_timestamp( # Приводим к UTC
F.col("event_ts"), "UTC"
))
.withColumn(DATE_COL, F.to_date(F.col("event_ts"))) # Дата из timestamp
.filter(F.col(KEY_COL).isNotNull()) # Убираем NULL ключи
)
# Intra-batch dedup: берём запись с максимальным event_ts при дублях
df_deduped = df_norm.dropDuplicates([KEY_COL])
# Кэшируем - будем использовать несколько раз
df_deduped.cache()
count = df_deduped.count()
logger.info(f"Stage 1: {count:,} unique rows after intra-batch dedup "
f"({time.time() - t:.1f}s)")
return df_deduped
# ── Stage 2 + 3: BF-assisted Anti-Join через Iceberg file-level index ─────────
def stage2_bloom_filter_anti_join(
spark: SparkSession,
df_batch: DataFrame,
) -> DataFrame:
"""
Использует Iceberg Bloom Filter Index для эффективного Anti-Join.
Вместо явного in-memory BF объекта полагаемся на file-level BF Index
в Silver таблице (Iceberg). Spark автоматически применяет его при
выполнении JOIN/Anti-Join, пропуская файлы, в которых точно нет ключей.
Дополнительно: Runtime Bloom Filter от AQE строит in-memory BF
по ключам батча и применяет к Silver при сканировании.
"""
t = time.time()
# Создаём временное представление батча
df_batch.createOrReplaceTempView("batch_events")
# Anti-Join: Catalyst + AQE + Runtime BF + Iceberg file-level BF
# работают совместно:
# 1. AQE строит Runtime BF по batch_events.event_id
# 2. BF транслируется на экзекьюторы, читающие Silver
# 3. Iceberg file-level BF дополнительно пропускает файлы
# 4. Из оставшихся файлов читаем только нужные строки
# 5. Anti-Join завершает точную фильтрацию
silver_keys = (
spark.table(SILVER_TABLE)
.filter(
# Ограничиваем lookback window - не смотрим в "вечность"
F.col(DATE_COL) >= F.date_sub(F.current_date(), LOOKBACK_DAYS)
)
.select(KEY_COL)
)
df_new = df_batch.join(silver_keys, on=KEY_COL, how="left_anti")
# Не вызываем count() здесь - это Action, откладываем до write
logger.info(f"Stage 2+3: Anti-Join plan built ({time.time() - t:.1f}s)")
return df_new
# ── Stage 4: Write ─────────────────────────────────────────────────────────────
def stage4_write_to_silver(df_new: DataFrame) -> None:
"""
Записывает новые строки в Silver (Iceberg).
При записи Iceberg автоматически строит BF Index для event_id
(если TBLPROPERTIES настроены) - файлы сразу доступны для
эффективного будущего Anti-Join.
"""
t = time.time()
# Записываем с партиционированием по дате
df_new.writeTo(SILVER_TABLE).append()
logger.info(f"Stage 4: Write complete ({time.time() - t:.1f}s)")
# ── Мониторинг ────────────────────────────────────────────────────────────────
def collect_metrics(spark: SparkSession, df_batch: DataFrame, df_new: DataFrame) -> dict:
"""
Собирает метрики для мониторинга качества дедупликации.
ВАЖНО: count() - это Action, выполняется отдельными задачами.
В production используйте Spark Listener или metrics API
для сбора без дополнительных Action.
"""
batch_count = df_batch.count()
new_count = df_new.count()
existing_count = batch_count - new_count
# Читаем статистику последнего commit из Iceberg metadata
iceberg_stats = spark.sql(f"""
SELECT
added_data_files_count,
added_records_count,
total_data_files_count,
total_records_count
FROM {SILVER_TABLE}.snapshots
ORDER BY committed_at DESC
LIMIT 1
""").collect()[0]
return {
"batch_size": batch_count,
"new_events": new_count,
"duplicate_events": existing_count,
"duplicate_rate": f"{existing_count / batch_count:.1%}",
"silver_total_files": iceberg_stats["total_data_files_count"],
"silver_total_rows": iceberg_stats["total_records_count"],
}
# ── Main Pipeline ──────────────────────────────────────────────────────────────
def main():
spark = (
SparkSession.builder
.appName("clickstream_dedup_bloom_filter")
.getOrCreate()
)
configure_spark(spark)
pipeline_start = time.time()
logger.info("=== Clickstream Dedup Pipeline START ===")
# Stage 1: Normalize & intra-batch dedup
df_batch = stage1_normalize_and_intra_dedup(spark)
# Stage 2+3: BF-assisted Anti-Join
df_new = stage2_bloom_filter_anti_join(spark, df_batch)
# Stage 4: Write to Silver (triggers lazy evaluation)
stage4_write_to_silver(df_new)
# Collect & log metrics
metrics = collect_metrics(spark, df_batch, df_new)
for k, v in metrics.items():
logger.info(f" {k}: {v}")
total_time = time.time() - pipeline_start
logger.info(f"=== Pipeline DONE in {total_time/60:.1f} min ===")
if __name__ == "__main__":
main()
11.4. Настройка Silver таблицы (Iceberg) для максимального эффекта¶
# Создаём Silver таблицу с полными оптимизациями
spark.sql(f"""
CREATE TABLE {SILVER_TABLE} (
event_id STRING NOT NULL COMMENT 'Уникальный ID события',
session_id STRING COMMENT 'ID сессии',
user_id BIGINT COMMENT 'ID пользователя',
event_type STRING COMMENT 'Тип события (click/view/purchase)',
event_date DATE COMMENT 'Дата события (партиция)',
event_ts TIMESTAMP COMMENT 'Время события в UTC',
page_url STRING COMMENT 'URL страницы',
device_type STRING COMMENT 'desktop/mobile/tablet',
country_code STRING COMMENT 'ISO код страны',
_ingested_at TIMESTAMP COMMENT 'Время ingestion',
_batch_id STRING COMMENT 'ID батча для трейсабилити'
)
USING iceberg
PARTITIONED BY (event_date)
TBLPROPERTIES (
-- Bloom Filter для высококардинальных ключей
'write.parquet.bloom-filter-enabled.column.event_id' = 'true',
'write.parquet.bloom-filter-fpp.column.event_id' = '0.01',
'write.parquet.bloom-filter-enabled.column.session_id' = 'true',
'write.parquet.bloom-filter-fpp.column.session_id' = '0.01',
'write.parquet.bloom-filter-enabled.column.user_id' = 'true',
'write.parquet.bloom-filter-fpp.column.user_id' = '0.01',
-- Целевой размер файла: 512 МБ (баланс IO и кол-ва файлов)
'write.target-file-size-bytes' = '536870912',
-- Автоматический compaction через OPTIMIZE
'write.spark.fanout.enabled' = 'true'
)
""")
11.5. Benchmark: до и после Bloom Filter¶
Реальные метрики для сценария 200M строк/час, 50B история:
| Метрика | Без BF (классический Anti-Join) | С Iceberg BF Index + Runtime BF |
|---|---|---|
| Время пайплайна | 3.5–4 часа | 25–35 минут |
| Shuffle Write | 380–420 ГБ | 40–60 ГБ |
| Shuffle Read | 370–410 ГБ | 38–58 ГБ |
| IO (Parquet read) | ~5 ТБ | ~100–200 ГБ |
| Файлов прочитано из Silver | 100% (~50,000 файлов) | 2–5% (~1,000–2,500 файлов) |
| OOM на экзекьюторах | Часто | Редко |
| Ускорение | 1x | ~7–10x |
Ключевой вклад в ускорение: Iceberg file-level BF сокращает IO с 5 ТБ до 100–200 ГБ. Runtime BF дополнительно уменьшает Shuffle.
Итоги урока¶
Bloom Filter - это не «магическая настройка», а архитектурный инструмент, требующий понимания workload-характеристик:
Что такое Bloom Filter: вероятностная структура данных на основе битового массива и k хэш-функций. Гарантирует отсутствие False Negatives, допускает False Positives с управляемой вероятностью FPR.
Три уровня применения в Spark:
- Runtime Bloom Filter (автоматически, Spark 3.3+) - Catalyst строит in-memory BF и применяет его до Shuffle
- stat.bloomFilter() API (явно) - вы строите фильтр в коде для custom pre-filter логики
- File-level BF Index (Delta/Iceberg) - BF хранится в Parquet-метаданных, пропускает нерелевантные файлы при сканировании
Когда применять: высококардинальные ключи, процент совпадений между таблицами менее 20–30%, Anti-Join или semi-join по большой исторической таблице, incremental dedup pipelines.
Когда не применять: низкая кардинальность (< 1000 уникальных значений), маленькая таблица умещается в broadcast, overlap > 50% (BF почти ничего не отсечёт).
Правильный sizing: всегда задавайте expectedNumItems с запасом 30%, иначе FPR деградирует при росте данных. Для 15 млрд ключей с FPR=1% нужно ~18 ГБ памяти (file-level) - строим BF Index в Iceberg/Delta, не в памяти.
End-to-end результат: для clickstream 200M событий/час на 50B истории - ускорение Anti-Join с 4 часов до 30 минут за счёт сокращения IO в 25x и Shuffle в 7x.