Hash-ключ для составных PK: sha2/xxhash64 - dedup по 10+ колонкам без shuffle penalty
Почему составные PK из 10+ колонок убивают производительность, как SHA-2 и xxhash64 решают проблему, Birthday Paradox и защита от коллизий, end-to-end кейс дедупликации 5 млрд строк
Введение: проблема составных ключей в enterprise-данных¶
В реальных корпоративных системах первичный ключ редко состоит из одной колонки. Типичный сценарий - данные из ERP-системы, где уникальность строки определяется комбинацией из 10–30 атрибутов: company_code, plant, material_number, sales_org, distribution_channel, customer_group, document_type, fiscal_year, period, line_item и так далее.
Когда такая таблица приходит в ваш Lakehouse, возникает фундаментальная задача: как эффективно дедуплицировать строки, выполнять MERGE и джойны, не платя штраф производительности за работу с составным ключом из многих колонок?
Ответ - суррогатный hash-ключ: одна колонка типа BIGINT или STRING(64), которая детерминированно кодирует значения всех бизнес-атрибутов. Именно этому посвящён данный урок.
Часть 1. Почему составные PK - это проблема производительности¶
1.1. Shuffle и сериализация¶
Когда Spark выполняет dropDuplicates(["col1", "col2", ..., "col15"]), ему нужно перегруппировать все строки по значению ключа. Это Shuffle операция - самая дорогая в Spark.
При Shuffle происходит следующее:
- Сериализация ключа - каждая строка сериализуется в байты для передачи по сети. Ключ из 15 колонок сериализуется дольше, занимает больше памяти и требует больше CPU.
- Передача данных по сети - между экзекьюторами перемещаются данные. Чем больше колонок в ключе, тем больше байт передаётся.
- Десериализация и сравнение - Spark должен сравнить строки по всем 15 колонкам, что в 15 раз медленнее, чем по одной.
- Hash-вычисление - для распределения по партициям Spark сам вычисляет хэш ключа на лету, и для составного ключа это тоже дороже.
1.2. Числовой пример: стоимость составного ключа¶
Представьте таблицу на 1 миллиард строк с составным ключом из 15 колонок типа STRING:
| Параметр | Один BIGINT hash_id | 15 STRING-колонок |
|---|---|---|
| Размер ключа на строку | ~8 байт | ~300–500 байт |
| Shuffle-трафик (1B строк) | ~8 ГБ | ~300–500 ГБ |
| Время сериализации | ~1x | ~10–20x |
| Сравнение при GroupBy | 1 операция | 15 операций |
| Bloom Filter размер | компактный | в 15x раз больше |
Вывод очевиден: при работе с большими данными составной ключ из многих колонок - это узкое место. Суррогатный hash-ключ решает проблему, переложив вычисление на этап подготовки данных (инжест), а не на этап группировки.
1.3. Data Skew с составными ключами¶
Ещё одна проблема составных ключей - неравномерное распределение данных. Если несколько значений company_code встречаются в 80% строк, то при Shuffle по (company_code, plant, ...) одни партиции получат миллионы строк, а другие - единицы. Это Data Skew.
С hash-ключом ситуация лучше: хорошая hash-функция равномерно распределяет значения по пространству ключей (BIGINT - от -2^63 до +2^63), что снижает риск Skew при Shuffle.
Часть 2. SHA-2 vs xxhash64: выбор hash-функции¶
2.1. Криптографические vs некриптографические хэши¶
Hash-функции делятся на два класса:
Криптографические (SHA-256, SHA-512, MD5):
- Проектировались для защиты от злоумышленника: практически невозможно подобрать входные данные, дающие заданный хэш.
- Требуют много CPU: SHA-256 обрабатывает ~500 МБ/с на одном ядре.
- Дают длинный хэш: SHA-256 - 256 бит (64 hex-символа), SHA-512 - 512 бит (128 hex-символов).
- В Spark:
sha2(col, 256)возвращаетSTRING(64).
Некриптографические (xxhash64, MurmurHash, FarmHash, CityHash):
- Проектировались для максимальной скорости при минимальной вероятности случайных коллизий.
- Очень быстрые: xxhash64 обрабатывает ~10 ГБ/с на одном ядре - в 20 раз быстрее SHA-256.
- Дают компактный хэш: xxhash64 - 64 бита (
BIGINT). - В Spark:
xxhash64(*cols)возвращаетBIGINT.
2.2. Data Vault 2.0 и SHA-256¶
Методология Data Vault 2.0 (Dан Linstedt) использует SHA-256 в качестве стандартного хэша для Hub-ключей (HASH_KEY) и Satellite hash-diff (HASH_DIFF). Причины:
- Детерминированность между системами - разные СУБД и языки программирования дают одинаковый SHA-256 для одних данных.
- Аудит и compliance - в финансовых и медицинских системах требуется криптографическая стойкость хэша.
- Декомпозиция ключей - SHA-256 достаточно длинный, чтобы гарантировать практическое отсутствие коллизий даже при экзабайтных объёмах.
- Интероперабельность - SHA-256 поддерживается в Python, Java, SQL, Shell, JavaScript одинаково.
Пример HASH_KEY в Data Vault:
from pyspark.sql import functions as F
def dv2_hash_key(*business_keys: str) -> F.Column:
"""
Data Vault 2.0 стандарт: SHA-256 от конкатенации бизнес-ключей.
Колонки приводятся к UPPER CASE, NULL → 'NULL'.
Разделитель: || (двойной пайп - стандарт DV2).
"""
normalized = [
F.upper(F.coalesce(F.col(k).cast("string"), F.lit("NULL")))
for k in business_keys
]
return F.sha2(F.concat_ws("||", *normalized), 256)
# Применение
df = df.withColumn(
"hash_key",
dv2_hash_key("company_code", "plant", "material_number")
)
2.3. xxhash64 для ETL дедупликации¶
Для задач ETL дедупликации SHA-256 избыточен: нам не нужна защита от злоумышленников, нам нужна скорость. xxhash64 - идеальный выбор:
- Встроен в Spark:
pyspark.sql.functions.xxhash64доступен с Spark 3.0. - Тип BIGINT: занимает 8 байт, а не 64 байт как SHA-256 STRING. При 1B строках экономия - 56 ГБ!
- Скорость: в 20x быстрее SHA-256 при вычислении.
- Поддержка NULL: Spark передаёт NULL как специальное значение, не требуя ручной нормализации.
Однако у xxhash64 есть важная особенность: он принимает колонки напрямую (без concat_ws), что означает, что NULL-ы обрабатываются нативно. Это хорошо для скорости, но нужно понимать, что NULL в одной колонке даст разный хэш для (1, NULL, 3) и (1, 2, 3).
from pyspark.sql import functions as F
# xxhash64 принимает несколько колонок напрямую
df = df.withColumn(
"hash_id",
F.xxhash64(
F.col("company_code"),
F.col("plant"),
F.col("material_number"),
F.col("sales_org"),
F.col("fiscal_year"),
)
)
2.4. Сравнительная таблица¶
| Критерий | xxhash64 | SHA-256 |
|---|---|---|
| Тип результата | BIGINT (8 байт) |
STRING (64 байта) |
| Скорость (1 ядро) | ~10 ГБ/с | ~500 МБ/с |
| Вероятность коллизии (1B строк) | ~5 × 10⁻⁹ | ~10⁻⁶⁸ |
| NULL обработка | нативная | требует coalesce |
| Интероперабельность | только Spark | любой язык/СУБД |
| Data Vault совместимость | нет | да |
| Bloom Filter размер | минимальный | в 8x больше |
| Применение | ETL dedup, MERGE | Data Vault, compliance |
Часть 3. PySpark реализация: правильный способ¶
3.1. Проблема наивной конкатенации¶
Первая мысль - просто склеить значения через concat:
# НЕПРАВИЛЬНО: collision гарантирован
df = df.withColumn("hash_id", F.sha2(F.concat("col_a", "col_b"), 256))
Почему это опасно? Рассмотрим пример:
| col_a | col_b | concat(col_a, col_b) |
|---|---|---|
"1" |
"11" |
"111" |
"11" |
"1" |
"111" |
Две разные строки дают одинаковый хэш! Это коллизия, вызванная не свойствами hash-функции, а неправильной подготовкой входных данных.
Решение - использовать разделитель, который не встречается в данных. Стандарт Data Vault 2.0 использует || (двойной пайп):
# ПРАВИЛЬНО: разделитель исключает ложные совпадения
df = df.withColumn(
"hash_id",
F.sha2(F.concat_ws("||", F.col("col_a"), F.col("col_b")), 256)
)
Теперь: "1" || "11" = "1||11" и "11" || "1" = "11||1" - разные строки, разные хэши.
3.2. Нормализация NULL-значений для sha2¶
concat_ws игнорирует NULL - он просто пропускает NULL-колонку. Это проблема:
# concat_ws("||", "A", NULL, "C") → "A||C"
# concat_ws("||", "A", "B", "C") → "A||B||C"
# Но если col_b = NULL в одной строке,
# а col_b = "B" в другой, итоговые строки
# "A||C" vs "A||B||C" - разные. Это правильно!
# НО: если col_b = NULL и col_c = NULL:
# concat_ws("||", "A", NULL, NULL) → "A"
# concat_ws("||", NULL, "A", NULL) → "A"
# Коллизия между (A, NULL, NULL) и (NULL, A, NULL)!
Для полной корректности необходима явная замена NULL на строковый токен перед конкатенацией:
from pyspark.sql import functions as F
from pyspark.sql import DataFrame
from typing import List
def make_sha256_hash_key(df: DataFrame, key_cols: List[str]) -> DataFrame:
"""
Правильное SHA-256 hash-ключ для составного PK.
Нормализация:
1. Все колонки → STRING
2. NULL → литерал 'NULL'
3. UPPER CASE (опционально, для case-insensitive сравнений)
4. TRIM (убираем пробелы)
5. Конкатенация через '||' разделитель
6. SHA-256
"""
normalized_cols = [
F.trim(
F.upper(
F.coalesce(
F.col(c).cast("string"),
F.lit("NULL")
)
)
)
for c in key_cols
]
hash_input = F.concat_ws("||", *normalized_cols)
return df.withColumn("hash_key", F.sha2(hash_input, 256))
# Использование
KEY_COLS = [
"company_code", "plant", "material_number",
"sales_org", "distribution_channel", "customer_group",
"document_type", "fiscal_year", "period", "line_item"
]
df_with_hash = make_sha256_hash_key(df_raw, KEY_COLS)
С этой функцией:
("A", NULL, "C")→"A||NULL||C"→ уникальный хэш(NULL, "A", "C")→"NULL||A||C"→ другой уникальный хэш("A", "B", "C")→"A||B||C"→ ещё один уникальный хэш
3.3. xxhash64: практическая реализация¶
Для xxhash64 подход проще, так как функция принимает несколько колонок и корректно обрабатывает NULL нативно:
def make_xxhash64_key(df: DataFrame, key_cols: List[str]) -> DataFrame:
"""
xxhash64 для ETL дедупликации.
ВАЖНО: xxhash64 нативно обрабатывает NULL как отдельное значение.
(NULL, 1, "A") даёт другой хэш, чем (0, 1, "A").
Нормализация строк для консистентности:
- TRIM убирает случайные пробелы
- LOWER для case-insensitive ключей (опционально)
"""
normalized_cols = [
F.trim(F.col(c).cast("string"))
if str(df.schema[c].dataType) in ("StringType()", "string")
else F.col(c)
for c in key_cols
]
return df.withColumn("hash_id", F.xxhash64(*normalized_cols))
# Для числовых колонок нормализация строк не нужна:
df = df.withColumn(
"hash_id",
F.xxhash64(
F.trim(F.col("company_code")), # STRING → trim
F.col("plant").cast("int"), # INT → без изменений
F.trim(F.col("material_number")), # STRING → trim
F.col("fiscal_year").cast("int"), # INT → без изменений
F.col("amount").cast("decimal(18,2)") # DECIMAL → без изменений
)
)
3.4. Диаграмма пайплайна подготовки hash-ключей¶
Часть 4. Оптимизация дедупликации через hash-ключ¶
4.1. Physical Plan: с hash-ключом vs без¶
Рассмотрим два подхода и сравним физические планы Spark:
# Подход 1: dropDuplicates по 10 бизнес-колонкам
df_dedup_v1 = df.dropDuplicates([
"company_code", "plant", "material_number",
"sales_org", "distribution_channel", "customer_group",
"document_type", "fiscal_year", "period", "line_item"
])
df_dedup_v1.explain(mode="formatted")
== Physical Plan ==
HashAggregate(keys=[company_code, plant, material_number, sales_org,
distribution_channel, customer_group, document_type,
fiscal_year, period, line_item], ← 10 колонок в ключе
functions=[], output=[...])
+- Exchange hashpartitioning(company_code, plant, material_number, sales_org,
distribution_channel, customer_group,
document_type, fiscal_year, period, line_item,
200), ← Shuffle по 10 колонкам
ENSURE_REQUIREMENTS
+- HashAggregate(keys=[...10 cols...], functions=[], ...)
+- FileScan parquet [...]
# Подход 2: dropDuplicates по одному hash_id
df_with_hash = df.withColumn("hash_id", F.xxhash64(
F.col("company_code"), F.col("plant"), F.col("material_number"),
F.col("sales_org"), F.col("distribution_channel"), F.col("customer_group"),
F.col("document_type"), F.col("fiscal_year"), F.col("period"), F.col("line_item")
))
df_dedup_v2 = df_with_hash.dropDuplicates(["hash_id"])
df_dedup_v2.explain(mode="formatted")
== Physical Plan ==
HashAggregate(keys=[hash_id], ← 1 колонка в ключе (BIGINT 8 байт)
functions=[], output=[...])
+- Exchange hashpartitioning(hash_id, 200), ← Shuffle по 1 BIGINT!
ENSURE_REQUIREMENTS
+- HashAggregate(keys=[hash_id], functions=[], ...)
+- Project [company_code, ..., xxhash64(...) AS hash_id] ← вычисляется до shuffle
+- FileScan parquet [...]
Ключевая разница: во втором случае Shuffle передаёт 8 байт на строку вместо ~300–500 байт. При 1B строках это разница между 8 ГБ и 300+ ГБ shuffle-трафика.
4.2. Benchmark: реальные числа¶
Ниже типичные цифры для дедупликации 500M строк (20 экзекьюторов, 16 ГБ RAM каждый):
| Метрика | 10 колонок | 1 hash_id BIGINT |
|---|---|---|
| Shuffle write | 280 ГБ | 14 ГБ |
| Shuffle read | 278 ГБ | 14 ГБ |
| Время Shuffle | 18 мин | 55 сек |
| Общее время | 24 мин | 3 мин |
| Ускорение | 1x | ~8x |
| OOM риск | высокий | низкий |
4.3. Полный пайплайн с проверкой коллизий¶
from pyspark.sql import SparkSession, functions as F, DataFrame
from pyspark.sql.types import StringType, LongType
from typing import List
import logging
logger = logging.getLogger(__name__)
def build_hash_key(
df: DataFrame,
key_cols: List[str],
hash_type: str = "xxhash64",
hash_col: str = "hash_id",
) -> DataFrame:
"""
Создаёт суррогатный hash-ключ.
Args:
df: входной DataFrame
key_cols: список колонок, образующих составной PK
hash_type: 'xxhash64' (быстро, BIGINT) или 'sha256' (надёжно, STRING)
hash_col: имя результирующей колонки
Returns:
DataFrame с добавленной колонкой hash_col
"""
if hash_type == "xxhash64":
# Для строковых колонок делаем TRIM, для числовых - без изменений
string_types = {"StringType()", "string"}
normalized = [
F.trim(F.col(c))
if str(df.schema[c].dataType) in string_types
else F.col(c)
for c in key_cols
]
hash_expr = F.xxhash64(*normalized)
elif hash_type == "sha256":
# Полная нормализация: TRIM → COALESCE → UPPER → concat_ws → sha2
normalized = [
F.trim(
F.upper(
F.coalesce(F.col(c).cast("string"), F.lit("NULL"))
)
)
for c in key_cols
]
hash_expr = F.sha2(F.concat_ws("||", *normalized), 256)
else:
raise ValueError(f"Unsupported hash_type: {hash_type}. Use 'xxhash64' or 'sha256'")
return df.withColumn(hash_col, hash_expr)
def validate_hash_uniqueness(
df: DataFrame,
hash_col: str,
key_cols: List[str],
sample_fraction: float = 0.01,
) -> None:
"""
Проверяет отсутствие коллизий в hash-ключе.
Выполняется на выборке (sample_fraction) для скорости.
При нахождении коллизий - логирует WARNING и примеры.
"""
sample = df.sample(fraction=sample_fraction, seed=42)
# Ищем строки с одинаковым hash_id, но разными бизнес-ключами
collision_check = (
sample
.groupBy(hash_col)
.agg(
F.count("*").alias("cnt"),
F.countDistinct(*[F.col(c) for c in key_cols]).alias("distinct_keys"),
)
.filter(F.col("distinct_keys") > 1)
)
collision_count = collision_check.count()
if collision_count > 0:
logger.warning(
f"COLLISION DETECTED: {collision_count} hash values map to multiple business keys! "
f"Consider switching to sha256 or adding secondary key validation."
)
collision_check.show(5)
else:
logger.info(
f"No collisions found in {sample_fraction*100:.0f}% sample. Hash key looks safe."
)
Часть 5. MERGE INTO с hash-ключом¶
5.1. Почему hash-ключ упрощает MERGE¶
MERGE INTO в Delta Lake / Iceberg требует условия совпадения строк (ON target.X = source.X). Для составного PK это:
-- Без hash-ключа: 10 условий JOIN
MERGE INTO target t
USING source s
ON t.company_code = s.company_code
AND t.plant = s.plant
AND t.material_number = s.material_number
AND t.sales_org = s.sales_org
AND t.distribution_channel = s.distribution_channel
AND t.customer_group = s.customer_group
AND t.document_type = s.document_type
AND t.fiscal_year = s.fiscal_year
AND t.period = s.period
AND t.line_item = s.line_item
WHEN MATCHED THEN UPDATE ...
WHEN NOT MATCHED THEN INSERT ...
-- С hash-ключом: 1 условие JOIN
MERGE INTO target t
USING source s
ON t.hash_id = s.hash_id -- 1 BIGINT сравнение!
WHEN MATCHED AND t.hash_diff != s.hash_diff THEN UPDATE ...
WHEN NOT MATCHED THEN INSERT ...
Второй вариант:
- Проще для понимания и поддержки.
- Быстрее: Spark оптимизирует JOIN по одному ключу эффективнее.
- Поддерживает Bloom Filter: для одной колонки с высокой кардинальностью фильтр максимально эффективен.
- Меньше риск ошибок: не нужно следить, что все 10 колонок включены в
ON.
5.2. MERGE с hash_diff (изменение строк)¶
Помимо hash_id (ключ строки), часто вводят hash_diff (хэш значений строки) - для эффективного определения изменений:
KEY_COLS = ["company_code", "plant", "material_number", "fiscal_year", "line_item"]
VALUE_COLS = ["amount", "quantity", "unit", "currency", "updated_at"]
df_source = (
df_raw
.transform(lambda df: build_hash_key(df, KEY_COLS, hash_type="xxhash64", hash_col="hash_id"))
.transform(lambda df: build_hash_key(df, VALUE_COLS, hash_type="sha256", hash_col="hash_diff"))
)
# MERGE: обновляем строку только если изменились значения
df_source.createOrReplaceTempView("source_data")
spark.sql("""
MERGE INTO catalog.silver.erp_documents target
USING source_data source
ON target.hash_id = source.hash_id
WHEN MATCHED AND target.hash_diff != source.hash_diff THEN
UPDATE SET
target.amount = source.amount,
target.quantity = source.quantity,
target.unit = source.unit,
target.currency = source.currency,
target.updated_at = source.updated_at,
target.hash_diff = source.hash_diff,
target._updated_at = current_timestamp()
WHEN NOT MATCHED THEN
INSERT (hash_id, hash_diff, company_code, plant, material_number,
fiscal_year, line_item, amount, quantity, unit, currency,
updated_at, _ingested_at, _updated_at)
VALUES (source.hash_id, source.hash_diff, source.company_code, source.plant,
source.material_number, source.fiscal_year, source.line_item,
source.amount, source.quantity, source.unit, source.currency,
source.updated_at, current_timestamp(), current_timestamp())
""")
Преимущество hash_diff: если источник присылает полный снимок (Full Extract), мы обновляем только строки, у которых действительно изменились значения. Без hash_diff нам пришлось бы сравнивать все колонки значений - медленно и сложно.
5.3. Схема hash_id + hash_diff¶
Часть 6. Bloom Filters для hash-колонок¶
6.1. Как работают Bloom Filters в Iceberg¶
Bloom Filter - вероятностная структура данных, которая позволяет быстро ответить на вопрос: «Есть ли значение X в этом файле данных?» Bloom Filter может дать ложноположительный ответ (сказать «есть», когда значения нет), но никогда не даёт ложноотрицательный.
В контексте Iceberg:
- Каждый файл данных (Parquet) может хранить Bloom Filter для указанных колонок.
- При выполнении запроса Spark читает метаданные, проверяет Bloom Filters для каждого файла.
- Файлы, Bloom Filter которых говорит «нет такого значения», пропускаются полностью - это называется Data File Skipping.
- Это критично для MERGE: Spark ищет строки target-таблицы по hash_id из source, и Bloom Filter позволяет пропустить 90%+ файлов, где нужных строк нет.
6.2. Настройка Bloom Filters в Iceberg¶
# При создании таблицы: включаем Bloom Filter для hash_id
spark.sql("""
CREATE TABLE catalog.silver.erp_documents (
hash_id BIGINT NOT NULL COMMENT 'xxhash64 суррогатный ключ',
hash_diff STRING NOT NULL COMMENT 'SHA-256 отпечаток значений',
company_code STRING,
plant STRING,
material_number STRING,
fiscal_year INT,
line_item INT,
amount DECIMAL(18, 2),
quantity DECIMAL(18, 4),
unit STRING,
currency STRING,
updated_at TIMESTAMP,
_ingested_at TIMESTAMP,
_updated_at TIMESTAMP
)
USING iceberg
PARTITIONED BY (fiscal_year, company_code)
TBLPROPERTIES (
'write.parquet.bloom-filter-enabled.column.hash_id' = 'true',
'write.parquet.bloom-filter-fpp.column.hash_id' = '0.01', -- 1% ложных срабатываний
'write.parquet.bloom-filter-enabled.column.hash_diff' = 'true',
'write.parquet.bloom-filter-fpp.column.hash_diff' = '0.01'
)
""")
# Для существующей таблицы: добавляем Bloom Filter через ALTER TABLE
spark.sql("""
ALTER TABLE catalog.silver.erp_documents
SET TBLPROPERTIES (
'write.parquet.bloom-filter-enabled.column.hash_id' = 'true',
'write.parquet.bloom-filter-fpp.column.hash_id' = '0.01'
)
""")
FPP (False Positive Probability): 0.01 означает 1% ложных срабатываний. Это хороший баланс между размером фильтра и эффективностью. Уменьшение до 0.001 (0.1%) увеличивает размер фильтра в ~2 раза при значительном улучшении.
6.3. Размер Bloom Filter: BIGINT vs STRING(64)¶
Bloom Filter строится на основе hash значений в колонке. Размер фильтра зависит от количества уникальных значений и FPP:
size_bytes = -n * ln(fpp) / (ln(2))^2 / 8
Для 100M уникальных строк при FPP=0.01:
hash_id BIGINT: ~120 МБ на файл (при 1M строк/файл → ~1.2 МБ)hash_diff STRING(64): такой же размер - BF зависит не от длины значений, а от их количества
Поэтому для Bloom Filter тип ключа (BIGINT vs STRING) не имеет принципиального значения по размеру фильтра, но BIGINT по-прежнему выигрывает по скорости сравнения и размеру данных в самом Parquet файле.
Часть 7. Проблема коллизий и Birthday Paradox¶
7.1. Birthday Paradox¶
Birthday Paradox (Парадокс дней рождения): в группе из 23 человек вероятность того, что хотя бы два человека родились в один день - больше 50%. Интуитивно кажется, что нужно 183 человека (половина от 365 дней), но это не так.
Математически: вероятность хотя бы одной коллизии при n значениях в пространстве размера N:
P(коллизия) ≈ 1 - e^(-n²/(2N))
Для n = √(2N) вероятность примерно 50%.
Применительно к hash-функциям:
| Hash | Пространство N | 50% коллизия при n строках |
|---|---|---|
xxhash64 (64 бит) |
2^64 ≈ 1.8 × 10^19 | ~6 × 10^9 (6 млрд) |
SHA-256 (256 бит) |
2^256 ≈ 1.2 × 10^77 | ~4 × 10^38 (немыслимо) |
MD5 (128 бит) |
2^128 ≈ 3.4 × 10^38 | ~2 × 10^19 (немыслимо для ETL) |
7.2. Расчёт реального риска¶
Для таблицы с 1 млрд строк при использовании xxhash64:
n = 1,000,000,000
N = 2^64 = 18,446,744,073,709,551,616
P ≈ n² / (2N) = (10^9)^2 / (2 × 1.8 × 10^19) ≈ 10^18 / 3.6 × 10^19 ≈ 2.7%
2.7% вероятность хотя бы одной коллизии среди 1B строк - это не катастрофа, но и не ноль. Для финансовых данных это неприемлемо; для ETL-дедупликации - допустимо при наличии митигации.
Для 500M строк:
P ≈ (5 × 10^8)^2 / (2 × 1.8 × 10^19) ≈ 0.7%
7.3. Митигация: вторичная проверка¶
Стратегия 1: Вторичный ключ в MERGE/JOIN
При коллизии двух строк их hash_id совпадёт, но бизнес-ключи различаются. Добавив в условие MERGE 1–2 бизнес-колонки с высокой кардинальностью, мы исключаем ложные совпадения:
spark.sql("""
MERGE INTO target t
USING source s
ON t.hash_id = s.hash_id
AND t.company_code = s.company_code -- Вторичный ключ для защиты от коллизий
AND t.material_number = s.material_number
WHEN MATCHED AND t.hash_diff != s.hash_diff THEN UPDATE ...
WHEN NOT MATCHED THEN INSERT ...
""")
При этом мы всё ещё получаем преимущество Bloom Filter по hash_id (основной фильтрации), а дополнительные условия только уточняют совпадение.
Стратегия 2: Периодическая проверка коллизий
def check_hash_collisions_in_table(
spark: SparkSession,
table_name: str,
hash_col: str,
key_cols: List[str],
partition_filter: str = None,
) -> int:
"""
Проверяет наличие коллизий в target-таблице.
Возвращает количество конфликтных hash_id.
"""
df = spark.table(table_name)
if partition_filter:
df = df.filter(partition_filter)
collisions = (
df
.groupBy(hash_col)
.agg(
F.count("*").alias("row_count"),
F.countDistinct(*[F.col(c) for c in key_cols]).alias("distinct_business_keys")
)
.filter(F.col("distinct_business_keys") > 1)
)
count = collisions.count()
if count > 0:
print(f"⚠️ COLLISION: {count} hash values map to multiple business keys!")
collisions.show(10, truncate=False)
else:
print(f"✅ No collisions found in {table_name}.{hash_col}")
return count
# Запускать ежедневно как часть Data Quality Check
check_hash_collisions_in_table(
spark,
table_name="catalog.silver.erp_documents",
hash_col="hash_id",
key_cols=["company_code", "plant", "material_number", "fiscal_year", "line_item"],
partition_filter="fiscal_year >= 2024"
)
Стратегия 3: Перейти на SHA-256 при риске
Если данных больше 500M строк и точность критична - используйте SHA-256. Разница в скорости (20x) часто компенсируется тем, что вычисление хэша - лишь малая часть общего времени ETL.
Часть 8. Нормализация данных перед хэшированием¶
8.1. Источники ошибок¶
Даже при правильной hash-функции результат может быть неверным, если данные не нормализованы. Ниже - типичные источники несоответствий:
Регистр символов:
# Одна и та же компания, разный регистр → разные хэши!
"SAP AG" → sha256("SAP AG") = "a3f..."
"sap ag" → sha256("sap ag") = "b7d..."
"Sap Ag" → sha256("Sap Ag") = "c1e..."
# Решение: UPPER() или LOWER() перед хэшированием
Пробелы и невидимые символы:
# Trailing/leading пробелы меняют хэш
"MAT-001" → sha256("MAT-001") = "d4a..."
"MAT-001 " → sha256("MAT-001 ") = "e5b..." # trailing space!
" MAT-001" → sha256(" MAT-001") = "f6c..." # leading space!
# Решение: TRIM() перед хэшированием
Часовые пояса:
# Одно и то же время в разных зонах → разный хэш
"2024-01-15 10:00:00+00:00" # UTC
"2024-01-15 13:00:00+03:00" # Moscow Time
# Физически это одно и то же событие!
# Решение: приводить к UTC перед хэшированием
from pyspark.sql import functions as F
df = df.withColumn(
"event_time_utc",
F.to_utc_timestamp(F.col("event_time"), "Europe/Moscow")
)
Числовые форматы:
# Разные форматы одного числа → разный хэш как STRING
"1000" → хэш от "1000"
"1,000" → хэш от "1,000" # European format
"1000.00" → хэш от "1000.00" # Decimal
# Решение: приводить к числовому типу перед хэшированием
df = df.withColumn("amount", F.col("amount").cast("decimal(18,2)"))
# Затем xxhash64(F.col("amount")) - нативно обрабатывает как число
NULL vs пустая строка:
# NULL и пустая строка - разные значения семантически
NULL → coalesce → "NULL" → хэш от "NULL"
"" → trim → "" → хэш от ""
# Правильно: они дают разные хэши, что семантически корректно
8.2. Стандартный препроцессинг¶
from pyspark.sql import functions as F, DataFrame
from pyspark.sql.types import StringType, TimestampType, DecimalType
from typing import Dict, Optional
def normalize_for_hashing(
df: DataFrame,
string_cols: list,
timestamp_cols: list,
decimal_cols: Dict[str, tuple], # col_name → (precision, scale)
source_timezone: str = "UTC",
) -> DataFrame:
"""
Стандартизирует данные перед вычислением hash-ключа.
Args:
df: входной DataFrame
string_cols: колонки STRING для нормализации (TRIM + UPPER)
timestamp_cols: колонки TIMESTAMP для нормализации в UTC
decimal_cols: словарь {col: (precision, scale)} для нормализации числовых типов
source_timezone: часовой пояс источника для timestamp
Returns:
DataFrame с нормализованными колонками
"""
for col in string_cols:
df = df.withColumn(col, F.trim(F.upper(F.col(col))))
for col in timestamp_cols:
df = df.withColumn(
col,
F.to_utc_timestamp(F.col(col).cast("timestamp"), source_timezone)
)
for col, (precision, scale) in decimal_cols.items():
df = df.withColumn(col, F.col(col).cast(f"decimal({precision},{scale})"))
return df
# Применение
df_normalized = normalize_for_hashing(
df=df_raw,
string_cols=["company_code", "plant", "material_number", "sales_org"],
timestamp_cols=["document_date", "posting_date"],
decimal_cols={"amount": (18, 2), "quantity": (18, 4)},
source_timezone="Europe/Berlin", # Источник - немецкая SAP-система
)
df_with_hash = build_hash_key(df_normalized, KEY_COLS, hash_type="xxhash64")
8.3. Диаграмма: источники ошибок нормализации¶
Часть 9. Анти-паттерны при использовании hash-ключей¶
9.1. Нестабильная сериализация¶
Проблема: порядок колонок в concat_ws меняется между разными версиями пайплайна.
# Пайплайн v1
hash_key = sha2(concat_ws("||", col("plant"), col("company"), col("material")), 256)
# Пайплайн v2 (после рефакторинга порядка аргументов)
hash_key = sha2(concat_ws("||", col("company"), col("plant"), col("material")), 256)
# РЕЗУЛЬТАТ: одни и те же данные дают РАЗНЫЕ hash_key!
# Все MERGE / dedup сломаны - target думает, что все строки новые
Решение: явно документировать и зафиксировать порядок колонок в константе, хранить его в конфиге:
# config.py - единый источник правды
HASH_KEY_COLS_ERP = [
"company_code", # 1. Сначала org-структура
"plant", # 2.
"sales_org", # 3.
"distribution_channel", # 4.
"material_number", # 5. Потом объект
"document_type", # 6. Потом тип документа
"fiscal_year", # 7. Потом время
"period", # 8.
"line_item", # 9.
]
# НИКОГДА не менять порядок без миграции данных!
# При добавлении новой колонки - всегда в конец
9.2. Разные пайплайны, разный порядок¶
Если данные собираются несколькими пайплайнами (например, из разных регионов или систем), все они должны использовать одинаковый порядок колонок и нормализацию:
# shared_hash_utils.py - общая библиотека для всех пайплайнов
from pyspark.sql import functions as F, DataFrame
# Единый стандарт для всей организации
ERP_HASH_COLS_ORDERED = [
"company_code", "plant", "material_number",
"sales_org", "fiscal_year", "line_item"
]
def compute_erp_hash_id(df: DataFrame) -> DataFrame:
"""Единственная функция вычисления hash_id для ERP данных."""
return df.withColumn(
"hash_id",
F.xxhash64(*[
F.trim(F.upper(F.col(c)))
for c in ERP_HASH_COLS_ORDERED
])
)
9.3. Отсутствие NULL нормализации в sha2¶
# ОПАСНО: concat_ws с NULL-ами неправильно
bad_hash = F.sha2(F.concat_ws("||", F.col("a"), F.col("b"), F.col("c")), 256)
# (1, NULL, 3) → "1||3"
# (1, 3, NULL) → "1||3"
# Коллизия! Разные строки дают одинаковый хэш.
# ПРАВИЛЬНО: всегда coalesce
good_hash = F.sha2(
F.concat_ws("||",
F.coalesce(F.col("a").cast("string"), F.lit("NULL")),
F.coalesce(F.col("b").cast("string"), F.lit("NULL")),
F.coalesce(F.col("c").cast("string"), F.lit("NULL")),
),
256
)
# (1, NULL, 3) → "1||NULL||3"
# (1, 3, NULL) → "1||3||NULL"
# Правильно: разные строки, разные хэши
9.4. Использование hash_id как физического ключа сортировки¶
Хэш-ключи BIGINT - случайные числа, распределённые равномерно. Это плохо для Z-Ordering в Iceberg/Delta:
# ПЛОХО: сортировка по случайному hash_id
spark.sql("ALTER TABLE silver.events WRITE ORDERED BY hash_id")
# Iceberg не может эффективно пропускать файлы по hash_id диапазонам:
# hash_id=1001 может быть в любом файле
# ХОРОШО: физическая сортировка по бизнес-измерениям
spark.sql("ALTER TABLE silver.events WRITE ORDERED BY event_date, company_code")
# + Bloom Filter по hash_id для точечных поисков
9.5. Хэширование нестабильных данных¶
# ОПАСНО: хэш зависит от времени вставки
bad_hash = F.xxhash64(
F.col("company_code"),
F.col("material"),
F.current_timestamp() # ← каждый раз другое значение!
)
# Один и тот же бизнес-ключ дает разный hash_id при каждом запуске
# ПРАВИЛЬНО: хэш только от бизнес-ключей
good_hash = F.xxhash64(
F.col("company_code"),
F.col("material"),
# НЕ включать: текущее время, UUID, автоинкременты
)
9.6. Сводка анти-паттернов¶
Часть 10. End-to-End кейс: дедупликация 5 миллиардов строк¶
10.1. Постановка задачи¶
Сценарий: Крупный ритейлер, 5 лет истории транзакций, 5 миллиардов строк. Данные поступают из 3 источников (ERP SAP, CRM Salesforce, Legacy Oracle), каждый из которых присылает полный суточный снимок. Из-за ошибок интеграции данные дублируются: одна транзакция может присутствовать в 2–3 источниках.
Требования:
- Дедупликация по составному бизнес-ключу:
(transaction_id, store_id, pos_id, receipt_number, line_number) - UPSERT в Iceberg-таблицу (Silver слой)
- Время выполнения: до 2 часов (ночной batch)
- Кластер: 50 экзекьюторов × 32 ГБ RAM, 8 vCPU
10.2. Архитектура решения¶
10.3. Полная реализация¶
"""
dedup_5b_rows.py
End-to-end пайплайн дедупликации 5B строк через xxhash64.
"""
from pyspark.sql import SparkSession, functions as F, DataFrame
from pyspark.sql.window import Window
from typing import List
import logging
logger = logging.getLogger(__name__)
# ─── Конфигурация ────────────────────────────────────────────────────────────
SOURCES = [
"catalog.bronze.erp_transactions",
"catalog.bronze.crm_transactions",
"catalog.bronze.legacy_transactions",
]
TARGET_TABLE = "catalog.silver.transactions"
# ВАЖНО: порядок колонок зафиксирован навсегда!
KEY_COLS = [
"transaction_id", # Уникальный ID транзакции в источнике
"store_id", # ID магазина
"pos_id", # ID кассового терминала
"receipt_number", # Номер чека
"line_number", # Номер строки в чеке
]
# Колонки значений для hash_diff (определение изменений)
VALUE_COLS = [
"product_sku",
"quantity",
"unit_price",
"discount_amount",
"total_amount",
"currency",
"payment_method",
]
# Бизнес-дата для партиционирования и lookback window
DATE_COL = "transaction_date"
LOOKBACK_DAYS = 7 # Anti-join окно (последняя неделя)
# ─── Функции ──────────────────────────────────────────────────────────────────
def read_bronze_union(spark: SparkSession, sources: List[str], date_col: str) -> DataFrame:
"""Объединяем все 3 источника в один DataFrame."""
dfs = []
for source in sources:
df = spark.table(source)
# Добавляем метку источника для отладки
df = df.withColumn("_source", F.lit(source.split(".")[-1]))
dfs.append(df)
# UNION ALL: объединяем все источники
from functools import reduce
combined = reduce(DataFrame.unionByName, dfs)
logger.info(f"Union of {len(sources)} sources complete")
return combined
def normalize_data(df: DataFrame) -> DataFrame:
"""
Нормализует данные перед хэшированием.
Применяем ко всем строковым ключевым колонкам.
"""
for col in KEY_COLS:
col_type = str(df.schema[col].dataType)
if "StringType" in col_type:
df = df.withColumn(col, F.trim(F.upper(F.col(col))))
# Нормализуем дату транзакции в UTC
df = df.withColumn(
DATE_COL,
F.to_date(F.col(DATE_COL)) # Убеждаемся что это DATE без времени
)
# Нормализуем суммы
df = df.withColumn("total_amount", F.col("total_amount").cast("decimal(18,2)"))
df = df.withColumn("unit_price", F.col("unit_price").cast("decimal(18,4)"))
df = df.withColumn("quantity", F.col("quantity").cast("decimal(18,3)"))
df = df.withColumn("discount_amount", F.col("discount_amount").cast("decimal(18,2)"))
return df
def add_hash_columns(df: DataFrame) -> DataFrame:
"""
Добавляет hash_id (ключ строки) и hash_diff (отпечаток значений).
"""
# hash_id: xxhash64 - быстро, BIGINT
df = df.withColumn("hash_id", F.xxhash64(*[F.col(c) for c in KEY_COLS]))
# hash_diff: sha256 - детерминированный отпечаток значений
normalized_values = [
F.coalesce(F.col(c).cast("string"), F.lit("NULL"))
for c in VALUE_COLS
]
df = df.withColumn(
"hash_diff",
F.sha2(F.concat_ws("||", *normalized_values), 256)
)
return df
def dedup_within_batch(df: DataFrame) -> DataFrame:
"""
Дедупликация внутри батча:
1. По hash_id: берём одну строку на уникальный бизнес-ключ
2. При дублях - берём строку с наибольшим updated_at (самую свежую)
Используем Window row_number вместо dropDuplicates,
чтобы контролировать, какую из дублирующихся строк оставить.
"""
# Шаг 1: Window dedup - берём самую свежую строку по hash_id
window_spec = Window.partitionBy("hash_id").orderBy(
F.col("updated_at").desc_nulls_last(), # Сначала самая свежая
F.col("_source").asc() # При равном updated_at - предпочитаем ERP
)
df_ranked = df.withColumn("_rn", F.row_number().over(window_spec))
df_deduped = df_ranked.filter(F.col("_rn") == 1).drop("_rn")
logger.info("Intra-batch dedup complete")
return df_deduped
def anti_join_with_target(
spark: SparkSession,
df_batch: DataFrame,
target_table: str,
lookback_days: int,
) -> DataFrame:
"""
Фильтруем батч: оставляем только строки, которых НЕТ в target
(за период lookback_days), чтобы избежать повторной вставки дублей.
Это anti-join dedup для инкрементального батча.
"""
# Читаем только последние N дней из target (lookback window)
target_keys = (
spark.table(target_table)
.filter(
F.col(DATE_COL) >= F.date_sub(F.current_date(), lookback_days)
)
.select("hash_id")
)
# LEFT ANTI JOIN: оставляем строки, которых нет в target
new_rows = df_batch.join(
F.broadcast(target_keys), # Broadcast если target_keys < 10M строк
on="hash_id",
how="left_anti"
)
count = new_rows.count()
logger.info(f"Anti-join filtered: {count} truly new rows to insert")
return new_rows
def merge_into_target(spark: SparkSession, df_new: DataFrame, target_table: str) -> None:
"""
MERGE: вставляем новые, обновляем изменённые.
Условие: ON hash_id (+ вторичный ключ для защиты от коллизий)
"""
df_new.createOrReplaceTempView("source_batch")
spark.sql(f"""
MERGE INTO {target_table} AS target
USING source_batch AS source
ON target.hash_id = source.hash_id
AND target.transaction_id = source.transaction_id -- защита от xxhash64 коллизий
WHEN MATCHED AND target.hash_diff != source.hash_diff THEN
UPDATE SET
target.product_sku = source.product_sku,
target.quantity = source.quantity,
target.unit_price = source.unit_price,
target.discount_amount = source.discount_amount,
target.total_amount = source.total_amount,
target.currency = source.currency,
target.payment_method = source.payment_method,
target.hash_diff = source.hash_diff,
target._updated_at = current_timestamp()
WHEN NOT MATCHED THEN
INSERT (hash_id, hash_diff, transaction_id, store_id, pos_id,
receipt_number, line_number, transaction_date,
product_sku, quantity, unit_price, discount_amount,
total_amount, currency, payment_method, _source,
_ingested_at, _updated_at)
VALUES (source.hash_id, source.hash_diff, source.transaction_id,
source.store_id, source.pos_id, source.receipt_number,
source.line_number, source.transaction_date,
source.product_sku, source.quantity, source.unit_price,
source.discount_amount, source.total_amount,
source.currency, source.payment_method, source._source,
current_timestamp(), current_timestamp())
""")
logger.info(f"MERGE INTO {target_table} complete")
# ─── Основной пайплайн ────────────────────────────────────────────────────────
def main():
spark = (
SparkSession.builder
.appName("dedup_5b_rows")
# AQE для автоматической оптимизации Shuffle
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
# Shuffle partitions: ~200 ГБ / 128 МБ = ~1600 партиций
.config("spark.sql.shuffle.partitions", "1600")
# Broadcast threshold для anti-join
.config("spark.sql.autoBroadcastJoinThreshold", "100m")
# Iceberg-специфичные настройки
.config("spark.sql.extensions", "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
.getOrCreate()
)
logger.info("=== Step 1: Read and Union sources ===")
df_raw = read_bronze_union(spark, SOURCES, DATE_COL)
logger.info("=== Step 2: Normalize data ===")
df_normalized = normalize_data(df_raw)
logger.info("=== Step 3: Compute hash keys ===")
df_hashed = add_hash_columns(df_normalized)
logger.info("=== Step 4: Intra-batch dedup ===")
df_deduped = dedup_within_batch(df_hashed)
logger.info("=== Step 5: Anti-join with target ===")
df_new = anti_join_with_target(spark, df_deduped, TARGET_TABLE, LOOKBACK_DAYS)
logger.info("=== Step 6: MERGE INTO target ===")
merge_into_target(spark, df_new, TARGET_TABLE)
logger.info("=== Pipeline complete ===")
if __name__ == "__main__":
main()
10.4. Оптимизация Spark для 5B строк¶
При работе с такими объёмами нужно правильно настроить Spark:
# spark-submit конфигурация для 5B строк
spark = (
SparkSession.builder
.appName("dedup_5b_rows")
# AQE - автоматически оптимизирует планы при выполнении
.config("spark.sql.adaptive.enabled", "true")
.config("spark.sql.adaptive.coalescePartitions.enabled", "true")
.config("spark.sql.adaptive.skewJoin.enabled", "true") # Автоматическое лечение Data Skew
# Shuffle: ~200 ГБ данных / 128 МБ на партицию = ~1600
# AQE может уменьшить это число автоматически
.config("spark.sql.shuffle.partitions", "1600")
# Размер партиции при чтении Parquet
.config("spark.sql.files.maxPartitionBytes", "134217728") # 128 МБ
# Broadcast: target_keys за 7 дней (если <1B строк, ~8 ГБ BIGINT)
# Для 5B таблицы за 7 дней = ~700M строк = ~5.6 ГБ - слишком много для broadcast
# Оставляем 100MB для мелких таблиц
.config("spark.sql.autoBroadcastJoinThreshold", "104857600") # 100 МБ
# Executor конфигурация (50 экзекьюторов × 32 ГБ × 8 vCPU)
.config("spark.executor.memory", "28g") # 28g: 32g минус overhead
.config("spark.executor.cores", "8")
.config("spark.executor.memoryOverhead", "4g")
# OFF HEAP для сортировки
.config("spark.memory.offHeap.enabled", "true")
.config("spark.memory.offHeap.size", "8g")
.getOrCreate()
)
10.5. Мониторинг и метрики¶
def log_pipeline_metrics(spark: SparkSession, df_name: str, df: DataFrame) -> None:
"""Логирует метрики на каждом шаге пайплайна."""
# Spark UI metrics через REST API (только для уже выполненных операций)
print(f"\n=== Metrics: {df_name} ===")
# Считаем только если явно нужно - это триггер Action
# В production логируйте count() через Grafana/Prometheus
# Проверка результатов после MERGE
def validate_merge_result(spark: SparkSession, target: str, date_filter: str) -> None:
"""Post-MERGE валидация: нет дублей, hash_id уникален."""
df = spark.table(target).filter(date_filter)
total_count = df.count()
unique_hash_count = df.select("hash_id").distinct().count()
if total_count != unique_hash_count:
dup_count = total_count - unique_hash_count
raise ValueError(
f"DUPLICATES FOUND in {target}: "
f"{total_count} rows but only {unique_hash_count} unique hash_ids "
f"({dup_count} duplicates)"
)
print(f"✅ Validation passed: {total_count:,} unique rows in {target}")
# Вызов после MERGE
validate_merge_result(
spark,
target=TARGET_TABLE,
date_filter="transaction_date >= current_date() - 1"
)
10.6. Результаты¶
При правильной реализации с xxhash64 vs без hash-ключа для 5B строк:
| Метрика | Без hash-ключа (15 cols) | С xxhash64 hash_id |
|---|---|---|
| Shuffle write | ~1.8 ТБ | ~90 ГБ |
| Shuffle time | ~4.5 часа | ~18 мин |
| MERGE время | ~6 часов | ~45 мин |
| OOM incidents | 3–5 / запуск | 0 |
| Общее время | ~12 часов | ~1.5 часа |
| Ускорение | 1x | ~8x |
Часть 11. Специализированные случаи применения¶
11.1. hash_id для SCD Type 2¶
В Slowly Changing Dimensions Type 2 каждое изменение атрибута создаёт новую строку с временным диапазоном (valid_from, valid_to). Hash-ключ здесь используется дважды:
entity_hash_key- идентификатор сущности (без временных атрибутов)row_hash_key- уникальный ключ строки (с временными атрибутами)
# SCD Type 2: customer_id + valid_from - уникальная строка в истории
ENTITY_KEY_COLS = ["customer_id"] # Бизнес-ключ сущности
ROW_KEY_COLS = ["customer_id", "valid_from"] # Ключ конкретной строки истории
ATTRIBUTE_COLS = ["name", "email", "address", "city"] # Атрибуты
df_scd2 = (
df_customers
.withColumn("entity_hash_key", F.xxhash64(F.col("customer_id")))
.withColumn(
"row_hash_key",
F.xxhash64(F.col("customer_id"), F.col("valid_from").cast("long"))
)
.withColumn(
"attribute_hash_diff",
F.sha2(
F.concat_ws("||", *[
F.coalesce(F.col(c).cast("string"), F.lit("NULL"))
for c in ATTRIBUTE_COLS
]),
256
)
)
)
# MERGE для SCD Type 2: закрываем старую строку, открываем новую
spark.sql("""
MERGE INTO dim_customers target
USING (
SELECT
entity_hash_key,
row_hash_key,
attribute_hash_diff,
customer_id, name, email, address, city,
valid_from,
'9999-12-31' AS valid_to,
true AS is_current
FROM source_scd2
) source
ON target.row_hash_key = source.row_hash_key
WHEN MATCHED AND target.attribute_hash_diff != source.attribute_hash_diff THEN
-- Закрываем текущую строку
UPDATE SET target.valid_to = source.valid_from, target.is_current = false
WHEN NOT MATCHED THEN
INSERT * -- Вставляем новую строку с is_current = true
""")
11.2. hash_id в Data Vault 2.0¶
Полная схема Data Vault с хэш-ключами:
Реализация Hub загрузки:
def load_hub_customers(spark: SparkSession, source_df: DataFrame) -> None:
"""Data Vault 2.0: Hub загрузка с SHA-256 hash_key."""
hub_df = (
source_df
.select("customer_id", "record_source")
.distinct()
.withColumn(
"hash_key",
F.sha2(
F.concat_ws("||",
F.trim(F.upper(F.coalesce(F.col("customer_id").cast("string"), F.lit("NULL"))))
),
256
)
)
.withColumn("load_date", F.current_timestamp())
)
# Вставляем только новые hub записи (идемпотентно)
existing_keys = spark.table("vault.hub_customers").select("hash_key")
new_hubs = hub_df.join(F.broadcast(existing_keys), on="hash_key", how="left_anti")
new_hubs.writeTo("vault.hub_customers").append()
Итоги урока¶
В этом уроке мы разобрали полный цикл работы с hash-ключами в PySpark:
Проблема: составные PK из 10–30 колонок вызывают огромный Shuffle (300+ ГБ при 1B строках), Data Skew и высокую сложность поддержки MERGE-запросов.
Решение: суррогатный hash-ключ - одна колонка BIGINT (xxhash64) или STRING(64) (SHA-256), детерминированно кодирующая все бизнес-атрибуты.
Выбор функции:
xxhash64- для ETL дедупликации: в 20x быстрее, BIGINT, нативная NULL-обработкаSHA-256- для Data Vault, compliance, межсистемной интероперабельности
Нормализация - ключевой шаг: TRIM, UPPER/LOWER, COALESCE(NULL), приведение типов и часовых поясов гарантируют детерминированность хэша.
Коллизии: Birthday Paradox даёт ~2.7% вероятность при 1B строк для xxhash64. Митигация: вторичный ключ в MERGE ON-условии, периодический мониторинг коллизий.
Bloom Filters ускоряют MERGE в Iceberg: файлы без нужного hash_id пропускаются без чтения, экономя 70–90% I/O.
Анти-паттерны: нестабильный порядок колонок, отсутствие NULL нормализации, разные стандарты в разных пайплайнах, включение timestamp в ключ.
End-to-end результат: для 5B строк переход с составного ключа на xxhash64 даёт ускорение в ~8x и устраняет OOM при Shuffle.