Hash-ключ для составных PK: sha2/xxhash64 - dedup по 10+ колонкам без shuffle penalty

Почему составные PK из 10+ колонок убивают производительность, как SHA-2 и xxhash64 решают проблему, Birthday Paradox и защита от коллизий, end-to-end кейс дедупликации 5 млрд строк

streaming

Введение: проблема составных ключей в 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 происходит следующее:

  1. Сериализация ключа - каждая строка сериализуется в байты для передачи по сети. Ключ из 15 колонок сериализуется дольше, занимает больше памяти и требует больше CPU.
  2. Передача данных по сети - между экзекьюторами перемещаются данные. Чем больше колонок в ключе, тем больше байт передаётся.
  3. Десериализация и сравнение - Spark должен сравнить строки по всем 15 колонкам, что в 15 раз медленнее, чем по одной.
  4. 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). Причины:

  1. Детерминированность между системами - разные СУБД и языки программирования дают одинаковый SHA-256 для одних данных.
  2. Аудит и compliance - в финансовых и медицинских системах требуется криптографическая стойкость хэша.
  3. Декомпозиция ключей - SHA-256 достаточно длинный, чтобы гарантировать практическое отсутствие коллизий даже при экзабайтных объёмах.
  4. Интероперабельность - 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.