Deduplication: dropDuplicates vs Window row_number - когда и что выбрать

Deduplication: dropDuplicates vs Window row_number - когда и что выбрать

streaming

Дублирование данных - одна из самых коварных проблем в data engineering. В отличие от упавшего pipeline, который виден сразу, дубли часто незаметны до тех пор, пока финансовый директор не замечает, что выручка «внезапно» выросла вдвое. К этому моменту аналитики уже несколько недель принимали решения на основе задвоенных цифр.

Spark предлагает два основных инструмента дедупликации: dropDuplicates() и Window + row_number(). Они внешне похожи - оба убирают «лишние» строки - но под капотом работают принципиально по-разному, потребляют разные объёмы памяти и CPU, и решают разные инженерные задачи. Понимание этой разницы отделяет инженера, который «написал код дедупликации», от инженера, который «спроектировал надёжный dedup-слой».


1. Природа дублей: два принципиально разных типа

1.1 Технические дубликаты

Технические дубликаты - это абсолютно идентичные строки, появившиеся из-за особенностей инфраструктуры:

  • Kafka at-least-once delivery: брокер при reconnect доставляет сообщение повторно
  • Airflow retry: задача упала и перезапустилась, источник отдал те же данные
  • HTTP API retry: при сетевом таймауте клиент повторил запрос, сервер обработал дважды
  • Debezium replication: после рестарта коннектора CDC начинает с последнего checkpoint, включая уже отправленные события
  • Double ingestion: два параллельных запуска пайплайна записали одни и те же файлы

Все поля таких дублей одинаковы, включая event_id, timestamp и payload. Задача проста: убрать лишние копии, оставить одну.

1.2 Бизнес-дубликаты (версионирование)

Бизнес-дубликаты - это несколько версий одной и той же сущности с разными значениями:

order_id=1001 | status=created    | updated_at=10:00:00 | amount=1500.0
order_id=1001 | status=paid       | updated_at=10:05:00 | amount=1500.0
order_id=1001 | status=shipped    | updated_at=11:00:00 | amount=1500.0
order_id=1001 | status=delivered  | updated_at=15:00:00 | amount=1500.0

Это не ошибка системы - это нормальная история изменений заказа через CDC. В Bronze-слое все четыре строки корректны и нужны. В Silver-слое для построения «текущего состояния» нужна только одна - актуальная (delivered). Но какую именно строку выбрать - это бизнес-решение, а не техническое.

1.3 Почему неверная дедупликация хуже её отсутствия

Отсутствие дедупликации даёт явное нарушение данных - выручка задвоена, её видно. Неверная дедупликация даёт молчаливое искажение: метрики выглядят правдоподобно, но содержат ошибку. Несколько реальных сценариев:

  • Выбрали строку status=created вместо status=delivered → выручка корректна, но конверсия воронки занижена
  • При дедупликации по user_id случайно удалили пользователя из когорты → ML-модель обучена на неполных данных
  • После CDC-дедупликации осталась версия с amount=0 (отмена), а не итоговая → выручка = 0 у реальных клиентов

Правило: перед написанием dedup-логики всегда формализуйте бизнес-правило: «правильная запись - это та, которая имеет максимальный LSN среди всех строк с одинаковым order_id».


2. Под капотом: как Spark физически выполняет дедупликацию

2.1 Физический план dropDuplicates

from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window

spark = SparkSession.builder \
    .appName("Deduplication-Deep-Dive") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "hadoop") \
    .config("spark.sql.catalog.lakehouse.warehouse", "/tmp/warehouse") \
    .getOrCreate()

# Тестовые данные: технические дубликаты по event_id
events = spark.createDataFrame([
    ("evt-001", "user_1", "click", "2026-05-23 10:00:00"),
    ("evt-001", "user_1", "click", "2026-05-23 10:00:00"),  # точный дубликат
    ("evt-002", "user_2", "view",  "2026-05-23 10:01:00"),
    ("evt-003", "user_1", "buy",   "2026-05-23 10:05:00"),
    ("evt-003", "user_1", "buy",   "2026-05-23 10:05:00"),  # точный дубликат
], ["event_id", "user_id", "event_type", "event_ts"])

# Применяем dropDuplicates и смотрим план выполнения
dedup_df = events.dropDuplicates(["event_id"])
dedup_df.explain(mode="formatted")

Вывод explain покажет физический план:

== Physical Plan ==
* HashAggregate (aggr expressions: [])     ← агрегация для нахождения уникальных строк
  * Exchange hashpartitioning(event_id)    ← SHUFFLE по хэшу event_id
    * HashAggregate (partial)              ← локальная предагрегация на executor
      * Scan

Что происходит шаг за шагом:

  1. Scan: каждый executor читает свою часть данных локально
  2. HashAggregate (partial): локально на каждом executor строится хэш-таблица по event_id. Если строка с таким event_id уже видена - она отбрасывается. Это reduces shuffle volume
  3. Exchange (Shuffle): строки перераспределяются между executors так, чтобы все строки с одинаковым event_id попали на один executor (hash partitioning)
  4. HashAggregate (final): финальная дедупликация - каждый executor видит все строки со «своими» event_id и оставляет по одной

Ключевое: dropDuplicates реализован через Hash Aggregation - ту же механику, что и GROUP BY. Это быстро, но не предоставляет никакого контроля над тем, какая именно строка остаётся.

2.2 Физический план Window + row_number

# Бизнес-дедупликация: берём строку с максимальным LSN
orders_cdc = spark.createDataFrame([
    ("order_1001", 101, "2026-05-23 10:00:00", "created",  1500.0),
    ("order_1001", 103, "2026-05-23 10:02:00", "paid",     1500.0),  # out-of-order
    ("order_1001", 102, "2026-05-23 10:01:00", "shipped",  1500.0),
    ("order_1002", 104, "2026-05-23 10:05:00", "created",   750.0),
], ["order_id", "lsn", "event_ts", "status", "amount"])

window_spec = Window.partitionBy("order_id").orderBy(F.col("lsn").desc())
dedup_window = orders_cdc \
    .withColumn("rn", F.row_number().over(window_spec)) \
    .filter(F.col("rn") == 1)

dedup_window.explain(mode="formatted")

Физический план Window функции значительно сложнее:

== Physical Plan ==
* Filter (rn = 1)
  * Window (row_number, partitionBy=[order_id], orderBy=[lsn DESC])
    * Exchange hashpartitioning(order_id)   ← SHUFFLE по partitionBy
      * Sort (order_id, lsn DESC)           ← СОРТИРОВКА перед window
        * Scan

Что добавляется по сравнению с dropDuplicates:

  1. Shuffle по partitionBy: аналогично dropDuplicates, но причина другая - нужно, чтобы все строки одного order_id попали на один executor для корректного вычисления ранга
  2. Sort: локальная сортировка внутри каждого executor по (order_id, lsn DESC). Это дополнительная O(n log n) операция, которой нет в dropDuplicates
  3. Window computation: вычисление row_number для каждой строки в отсортированной партиции
  4. Filter: финальная фильтрация rn = 1

2.3 Сравнение физических планов: диаграмма

Ключевое отличие: dropDuplicates делает Shuffle + Hash Aggregation. Window row_number делает Shuffle + Sort + Window Computation. Сортировка - это дополнительная нагрузка, но именно она обеспечивает детерминизм: мы знаем, какая строка останется.


3. dropDuplicates: когда, как и чего избегать

3.1 Базовое использование

# Пример 1: Полная дедупликация - все поля должны совпадать
all_fields_dedup = events.distinct()
# Эквивалентно: events.dropDuplicates()
# Использовать когда: точная копия строки, включая все метаданные

# Пример 2: Дедупликация по бизнес-ключу
# Использовать когда: нужно убрать retries по UUID события
key_dedup = events.dropDuplicates(["event_id"])
# ВАЖНО: из группы дублей останется ПРОИЗВОЛЬНАЯ строка
# Для технических ретраев это OK: все дубли идентичны

# Пример 3: Составной ключ
composite_dedup = events.dropDuplicates(["user_id", "event_type", "event_ts"])
# Использовать когда: нет единого UUID, ключ составной

3.2 Недетерминизм dropDuplicates: практическая демонстрация

# Ситуация: у одного event_id разные значения в поле processing_node
# (разные executor-ы обработали один и тот же ивент)
ambiguous_events = spark.createDataFrame([
    ("evt-001", "user_1", "click", "executor-01"),
    ("evt-001", "user_1", "click", "executor-02"),  # то же событие, другой executor
], ["event_id", "user_id", "event_type", "processing_node"])

# dropDuplicates("event_id") оставит произвольную строку
# При каждом запуске может выбрать executor-01 или executor-02
result = ambiguous_events.dropDuplicates(["event_id"])
result.show()
# Вывод: либо executor-01, либо executor-02 - непредсказуемо!

# Это допустимо если:
# 1. processing_node не влияет на бизнес-логику (технический атрибут)
# 2. Нам важно только убрать дубль, а не выбрать конкретный

# Это НЕДОПУСТИМО если:
# 1. Разные строки имеют разные значения в бизнес-полях
# 2. Нужен детерминированный результат для audit trail

3.3 Оптимизация: repartition перед dropDuplicates

Если данные изначально распределены неравномерно по ключу дедупликации, Shuffle в dropDuplicates создаёт дополнительную нагрузку. Помочь может предварительный repartition:

# Без оптимизации: Spark перераспределяет данные при dropDuplicates
basic_dedup = large_events_df.dropDuplicates(["event_id"])

# С оптимизацией: если данные уже партиционированы по event_id,
# Shuffle в dropDuplicates становится no-op (данные уже там где надо)
optimized_dedup = (
    large_events_df
    .repartition(200, F.col("event_id"))  # явное партиционирование
    .dropDuplicates(["event_id"])          # Shuffle теперь минимален
)

# Ещё эффективнее для incremental pipelines:
# если Bronze таблица bucketed по event_id,
# Spark может читать данные без перераспределения

3.4 Когда dropDuplicates - правильный выбор

Сценарий Подходит dropDuplicates? Почему
Kafka at-least-once retries Дубли идентичны, выбор строки не важен
Bronze ingestion cleanup Технический уровень, нет бизнес-логики
Дедупликация event_id из CDC UUID одинаков, строки идентичны
Выбор актуального статуса заказа Нужен порядок - кто «последний»
CDC с разными версиями строки Строки различаются, нужен LSN
SCD Type 2 dedup Требует valid_from/valid_to логику

4. Window row_number: детерминированная бизнес-дедупликация

4.1 row_number, rank и dense_rank: разница

# Тестовые данные: три события для одного заказа с одинаковым timestamp у двух
cdc_data = spark.createDataFrame([
    ("order_1001", 101, "2026-05-23 10:00:00", "created"),
    ("order_1001", 103, "2026-05-23 10:02:00", "paid"),
    ("order_1001", 103, "2026-05-23 10:02:00", "paid_duplicate"),  # tie по lsn!
    ("order_1002", 104, "2026-05-23 10:05:00", "created"),
], ["order_id", "lsn", "event_ts", "status"])

window = Window.partitionBy("order_id").orderBy(F.col("lsn").desc())

comparison = cdc_data \
    .withColumn("row_number",   F.row_number().over(window)) \
    .withColumn("rank",         F.rank().over(window)) \
    .withColumn("dense_rank",   F.dense_rank().over(window))

comparison.show()
# order_id   | lsn | status          | row_number | rank | dense_rank
# order_1001 | 103 | paid            |     1      |  1   |     1
# order_1001 | 103 | paid_duplicate  |     2      |  1   |     1       ← разница!
# order_1001 | 101 | created         |     3      |  3   |     2
# order_1002 | 104 | created         |     1      |  1   |     1

Три функции ведут себя по-разному при tie (одинаковых значениях в ORDER BY):

  • row_number(): присваивает уникальный ранг каждой строке даже при tie. Среди строк с одинаковым lsn=103 одна получит rn=1, другая rn=2. Порядок при tie недетерминирован - это важно понимать!
  • rank(): при tie обе строки получают одинаковый ранг 1. Следующий ранг пропускается (3, а не 2).
  • dense_rank(): при tie обе строки получают ранг 1, следующий ранг 2 (без пропуска).

Для дедупликации используйте row_number(): он гарантирует ровно одну строку с rn=1. При rank() или dense_rank() может остаться несколько строк с rank=1 при tie.

4.2 Детерминизм при tie: добавляем secondary sort

Если в данных возможны tie по основному ключу сортировки, нужен secondary sort - второй критерий для однозначного упорядочивания:

# Проблема: два события с одинаковым lsn
# row_number() выберет произвольно - это НЕДетерминировано!

# Решение: добавляем secondary sort по полю, которое гарантирует уникальность
# Например, event_id (UUID) всегда уникален
window_deterministic = Window.partitionBy("order_id").orderBy(
    F.col("lsn").desc(),       # primary: LSN убыванием (последнее событие)
    F.col("event_id").desc()   # secondary: UUID как tie-breaker
)

deterministic_dedup = cdc_data \
    .withColumn("rn", F.row_number().over(window_deterministic)) \
    .filter(F.col("rn") == 1) \
    .drop("rn")

# Теперь при любом количестве запусков - один и тот же результат

Добавление event_id (UUID) как вторичного критерия сортировки превращает row_number из «произвольного» в «детерминированный». UUID уникален, поэтому среди любых двух строк с одинаковым LSN всегда одна «победит» - всегда одна и та же при повторном запуске.

4.3 Паттерны бизнес-дедупликации

# Паттерн 1: "Latest Event Wins" - берём строку с максимальным LSN
# Применение: CDC-поток, обновление профиля клиента
latest_by_lsn = (
    cdc_df
    .withColumn("rn",
        F.row_number().over(
            Window.partitionBy("entity_id")
                  .orderBy(F.col("lsn").desc(), F.col("event_id").desc())
        )
    )
    .filter(F.col("rn") == 1)
    .drop("rn")
)

# Паттерн 2: "Most Complete Record" - берём строку с наибольшим числом непустых полей
# Применение: разные источники присылают частичные записи
most_complete = (
    messy_df
    .withColumn("completeness_score",
        sum([F.when(F.col(c).isNotNull(), 1).otherwise(0)
             for c in messy_df.columns])
    )
    .withColumn("rn",
        F.row_number().over(
            Window.partitionBy("entity_id")
                  .orderBy(F.col("completeness_score").desc(),
                           F.col("event_ts").desc())
        )
    )
    .filter(F.col("rn") == 1)
    .drop("rn", "completeness_score")
)

# Паттерн 3: "Priority Source Wins" - данные из CRM важнее данных из API
# Применение: несколько источников дают разные версии одной сущности
source_priority = {
    "crm_system": 1,   # высший приоритет
    "api_v2":     2,
    "legacy_db":  3,   # низший приоритет
}

priority_dedup = (
    multi_source_df
    .withColumn("source_priority",
        F.when(F.col("source") == "crm_system", F.lit(1))
         .when(F.col("source") == "api_v2",     F.lit(2))
         .otherwise(F.lit(3))
    )
    .withColumn("rn",
        F.row_number().over(
            Window.partitionBy("customer_id")
                  .orderBy(F.col("source_priority").asc(),  # 1 = победитель
                           F.col("event_ts").desc())
        )
    )
    .filter(F.col("rn") == 1)
    .drop("rn", "source_priority")
)

5. Производительность: когда каждый подход тормозит

5.1 Проблема Data Skew в дедупликации

Data Skew - это ситуация, когда один partition key встречается несравнимо чаще других. Для дедупликации это критично: весь массив строк с «горячим» ключом уходит на один executor.

Как Skew проявляется для каждого подхода:

  • dropDuplicates: executor с горячим ключом строит огромную hash-таблицу. При 50M строк одного user_id - почти гарантированный OOM или спилл на диск.
  • Window row_number: ещё хуже - executor должен не только хранить все 50M строк в памяти, но и сортировать их. Сортировка 50M строк - это секунды или минуты на одном узле, пока остальные ждут.

5.2 Решение Skew для dropDuplicates: ранняя фильтрация

# Стратегия 1: Filter before dedup - убираем известные «ботовые» ключи
clean_df = raw_events \
    .filter(F.col("user_id") != "bot_user") \
    .dropDuplicates(["event_id"])

# Стратегия 2: Отдельная обработка «горячих» ключей
hot_keys = ["user_BOT", "user_CRAWLER"]

normal_events = raw_events.filter(~F.col("user_id").isin(hot_keys))
hot_events    = raw_events.filter(F.col("user_id").isin(hot_keys))

# Нормальные события: обычный dropDuplicates
normal_dedup = normal_events.dropDuplicates(["event_id"])

# Горячие ключи: агрегируем вместо дедупликации (считаем только count)
hot_aggregated = (
    hot_events
    .groupBy("user_id", "event_type", F.to_date("event_ts").alias("event_date"))
    .agg(F.countDistinct("event_id").alias("event_count"))
)

# Объединяем обработанные наборы
result = normal_dedup.unionByName(hot_aggregated, allowMissingColumns=True)

5.3 Решение Skew для Window: Salting + двухфазная дедупликация

# Window с горячим ключом тормозит, потому что все строки ключа на одном executor
# Решение: Salting - добавляем случайный суффикс к ключу для распределения

import random

SALT_BUCKETS = 10  # Делим горячий ключ на 10 частей

# Фаза 1: Дедупликация внутри "соли" (параллельно на 10 executor-ах)
phase1 = (
    large_cdc_df
    # Для нормальных ключей: salt=0 (дедупликация без разделения)
    # Для горячих ключей: salt=0..9 (разделяем на 10 частей)
    .withColumn("salt",
        F.when(
            F.col("user_id").isin(hot_keys),
            (F.crc32(F.col("event_id")) % SALT_BUCKETS).cast("string")
        ).otherwise(F.lit("0"))
    )
    .withColumn("salted_key",
        F.concat_ws("__", F.col("user_id"), F.col("salt")))
    .withColumn("rn",
        F.row_number().over(
            Window.partitionBy("salted_key")
                  .orderBy(F.col("lsn").desc())
        )
    )
    .filter(F.col("rn") == 1)
    .drop("rn", "salt", "salted_key")
)

# Фаза 2: Финальная дедупликация по реальному ключу
# После фазы 1 данных стало в 10 раз меньше - теперь Window работает быстро
phase2 = (
    phase1
    .withColumn("rn",
        F.row_number().over(
            Window.partitionBy("user_id")
                  .orderBy(F.col("lsn").desc())
        )
    )
    .filter(F.col("rn") == 1)
    .drop("rn")
)

Как работает двухфазный подход: в фазе 1 горячий ключ user_BOT разделяется на 10 «солёных» ключей user_BOT__0 ... user_BOT__9. Каждый executor обрабатывает ~5M строк вместо 50M. После фазы 1 для каждой «соли» осталась одна строка - итого максимум 10 строк для горячего ключа. Фаза 2 делает финальную дедупликацию этих 10 строк - это мгновенно.

5.4 Сравнение производительности: матрица выбора


6. Дедупликация в Structured Streaming

6.1 Почему streaming dedup сложнее batch

В batch-обработке все данные известны заранее - мы просто выбираем лучшую строку из набора. В streaming системе поток бесконечен: события продолжают поступать, и мы не знаем, придёт ли ещё один дубликат через 5 секунд или через 5 минут.

Поэтому streaming deduplication требует state management: Spark должен помнить, какие event_id уже были обработаны, чтобы отбросить повторные доставки.

# Stateful deduplication в Structured Streaming
streaming_df = (
    spark.readStream
         .format("kafka")
         .option("kafka.bootstrap.servers", "kafka:9092")
         .option("subscribe", "order-events")
         .load()
         .selectExpr("CAST(value AS STRING) AS json_payload",
                     "timestamp AS kafka_ts")
         # ... десериализация JSON ...
         .withColumn("event_ts", F.col("event_ts").cast("timestamp"))
)

# withWatermark: говорим Spark, как долго ждать опоздавшие события
# "10 minutes" = события, опоздавшие более чем на 10 минут после watermark,
#                 будут отброшены без обработки
# Это также ограничивает размер state: Spark забывает event_id старше watermark
deduplicated_stream = (
    streaming_df
    .withWatermark("event_ts", "10 minutes")
    .dropDuplicates(["event_id"])  # stateful dedup: state живёт в рамках watermark
)

query = (
    deduplicated_stream
    .writeStream
    .format("iceberg")
    .outputMode("append")
    .option("checkpointLocation", "s3://checkpoints/dedup-stream/")
    .toTable("lakehouse.bronze.order_events_dedup")
)

6.2 Проблема Infinite State без Watermark

Без watermark Spark хранит в state store все event_id, которые когда-либо были обработаны. После нескольких дней работы стрима state store переполняется - OOM на executor.

С watermark Spark периодически «забывает» старые event_id. Watermark = max(event_ts) - delay. Когда watermark проходит время события, его event_id удаляется из state. Это ограничивает размер state временным окном.

Компромисс watermark: более широкое окно ("1 hour") даёт лучшую защиту от поздних дублей, но потребляет больше памяти. Более узкое окно ("1 minute") экономит память, но дубли, пришедшие через 2 минуты, не будут отловлены.

6.3 Дедупликация при Replay

При повторном запуске streaming pipeline с checkpoint Spark продолжит с сохранённого offset. Но state store при этом восстанавливается из checkpoint: Spark «помнит», какие event_id уже были обработаны в рамках текущего watermark-окна.

Если произошёл полный сброс checkpoint (например, при смене логики), state потерян. В этом случае для защиты от дублей при replay нужен дополнительный слой - например, MERGE INTO в sink-таблице.


7. groupBy + first/last: третий подход

7.1 Агрегационная дедупликация

Помимо dropDuplicates и Window, существует третий подход - агрегация через groupBy:

# groupBy + first: быстрая дедупликация с выбором произвольной строки
# (аналогично dropDuplicates по производительности, но с контролем над агрегатами)
groupby_dedup = (
    orders_cdc
    .groupBy("order_id")
    .agg(
        F.max("lsn").alias("max_lsn"),               # берём максимальный LSN
        F.first("status").alias("status"),           # произвольный статус (!)
        F.max("amount").alias("amount"),             # максимальная сумма
        F.max("event_ts").alias("latest_event_ts"),  # последнее время
    )
)
# ПРОБЛЕМА: first("status") берёт произвольный статус, не тот что соответствует max_lsn!
# Это АНТИПАТТЕРН для CDC-дедупликации

# Правильный вариант через struct + max:
# Упаковываем нужные поля в struct, берём struct с максимальным LSN
correct_groupby = (
    orders_cdc
    .withColumn("latest_record",
        F.struct(
            F.col("lsn"),
            F.col("status"),
            F.col("amount"),
            F.col("event_ts"),
        )
    )
    .groupBy("order_id")
    .agg(F.max("latest_record").alias("rec"))
    .select(
        "order_id",
        F.col("rec.lsn").alias("lsn"),
        F.col("rec.status").alias("status"),
        F.col("rec.amount").alias("amount"),
    )
)
# max("struct") в Spark сравнивает struct-ы лексикографически по первому полю
# Если первое поле - lsn (BIGINT), то max выберет строку с максимальным LSN
# Это детерминированно и быстрее Window (нет Sort!)

F.max(F.struct(...)) - менее известная техника. Spark сравнивает struct-ы лексикографически: сначала по первому полю, при равенстве - по второму. Если первым полем в struct указан lsn, то max выберет строку с максимальным LSN - это эквивалентно ROW_NUMBER ORDER BY lsn DESC, но без фазы сортировки.

7.2 Сравнение трёх подходов

Критерий dropDuplicates Window row_number groupBy + max(struct)
Детерминизм ❌ Произвольная строка ✅ По ORDER BY ✅ По первому полю struct
CPU (Sort) Нет сортировки Сортировка O(n log n) Нет сортировки
Memory Hash-таблица Hash + Sort buffer Hash-таблица
Гибкость логики Минимальная Максимальная Средняя
Skew-устойчивость Средняя Плохая (без Salting) Средняя
Streaming ✅ dropDuplicates+WM ❌ Не поддерживается ❌ Не поддерживается
Когда использовать Техн. дубли Бизнес-дедупликация Быстрая dedup с выбором поля

8. Антипаттерны дедупликации

8.1 Дедупликация без детерминированного ordering

# АНТИПАТТЕРН: нет secondary sort - tie приводит к недетерминированному результату
bad_window = Window.partitionBy("order_id").orderBy(F.col("event_ts").desc())
# Если два события имеют одинаковый event_ts - row_number выберет произвольно
# При перезапуске может выбрать другую строку → разные результаты = ломаный pipeline

# ПРАВИЛЬНО: всегда добавляйте tie-breaker
good_window = Window.partitionBy("order_id").orderBy(
    F.col("event_ts").desc(),
    F.col("event_id").desc()   # UUID - всегда уникален, гарантирует детерминизм
)

8.2 current_timestamp() как ordering key

# АНТИПАТТЕРН: сортировка по времени обработки вместо времени события
very_bad_window = Window.partitionBy("order_id").orderBy(
    F.current_timestamp().desc()  # ЭТО ВСЕГДА ОДИНАКОВОЕ ЗНАЧЕНИЕ для всего батча!
)
# current_timestamp() внутри окна возвращает одно значение для всего micro-batch
# Это не даёт никакой полезной сортировки между строками батча!

# ПРАВИЛЬНО: используйте время из источника (LSN или event_ts из WAL)
correct_window = Window.partitionBy("order_id").orderBy(
    F.col("lsn").desc()  # LSN из CDC источника - монотонный, корректный
)

8.3 Дедупликация после JOIN порождает новые дубли

# Типичная ловушка: JOIN создаёт дубли, которых не было до него

orders_df   = spark.table("silver.orders")         # 1000 строк
customers_df = spark.table("silver.customers")     # может иметь дубли по customer_id!

# Если customers_df содержит 3 строки для customer_A:
# JOIN (orders × customers) для customer_A даст 3 строки вместо 1!
joined = orders_df.join(customers_df, on="customer_id", how="left")

# Результат: joined содержит БОЛЬШЕ строк, чем orders_df
# Дедупликация ПОСЛЕ join не поможет - строки заказа теперь различаются

# ПРАВИЛЬНО: дедупликация ПЕРЕД join
clean_customers = customers_df.dropDuplicates(["customer_id"])
joined_clean = orders_df.join(clean_customers, on="customer_id", how="left")

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

8.4 Игнорирование Null-ключей

# АНТИПАТТЕРН: в данных есть null customer_id
# dropDuplicates оставит все null-строки (null != null в SQL семантике)
events_with_nulls = spark.createDataFrame([
    (None,     "click"),  # null customer_id
    (None,     "buy"),    # тоже null - dropDuplicates ОСТАВИТ ОБЕ!
    ("user_1", "click"),
], ["customer_id", "event_type"])

result = events_with_nulls.dropDuplicates(["customer_id"])
result.show()
# customer_id | event_type
# null        | click        ← остался
# null        | buy          ← тоже остался! null != null
# user_1      | click

# ПРАВИЛЬНО: обрабатываем null-ключи отдельно
non_null = events_with_nulls.filter(F.col("customer_id").isNotNull()) \
                              .dropDuplicates(["customer_id"])
null_events = events_with_nulls.filter(F.col("customer_id").isNull()) \
                                 .withColumn("customer_id", F.lit("UNKNOWN"))

result_clean = non_null.unionByName(null_events)

9. End-to-End кейс: production dedup pipeline

9.1 Сценарий

Компания получает события заказов из Kafka (at-least-once). Нужно построить pipeline: Bronze (убираем технические дубли) → Silver (берём актуальный статус) → Gold (выручка без double-counting).

import datetime
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window

spark = SparkSession.builder \
    .appName("Dedup-E2E-Pipeline") \
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions") \
    .config("spark.sql.catalog.lakehouse", "org.apache.iceberg.spark.SparkCatalog") \
    .config("spark.sql.catalog.lakehouse.type", "hadoop") \
    .config("spark.sql.catalog.lakehouse.warehouse", "/tmp/warehouse") \
    .getOrCreate()

# ─── СИМУЛЯЦИЯ KAFKA INPUT ────────────────────────────────────────────────────
# Входные данные из Kafka: технические дубли + бизнес-версии одного заказа
raw_kafka_events = spark.createDataFrame([
    # Технические дубли: evt-001 доставлен дважды (Kafka retry)
    ("evt-001", "order_1001", 101, "2026-05-23 10:00:00", "created",  1500.0),
    ("evt-001", "order_1001", 101, "2026-05-23 10:00:00", "created",  1500.0),  # дубль!
    # Бизнес-версии: заказ меняет статус (CDC UPDATE)
    ("evt-002", "order_1001", 102, "2026-05-23 10:05:00", "paid",     1500.0),
    ("evt-003", "order_1001", 103, "2026-05-23 11:00:00", "shipped",  1500.0),
    # Out-of-order: evt-004 пришло после evt-003, но LSN меньше (network delay)
    ("evt-004", "order_1001", 104, "2026-05-23 12:00:00", "delivered", 1500.0),
    # Другой заказ
    ("evt-005", "order_1002", 105, "2026-05-23 10:30:00", "created",   750.0),
    ("evt-006", "order_1002", 106, "2026-05-23 10:35:00", "cancelled",   0.0),
    # Технический дубль для order_1002
    ("evt-005", "order_1002", 105, "2026-05-23 10:30:00", "created",   750.0),  # дубль!
], ["event_id", "order_id", "lsn", "event_ts", "status", "amount"])

print(f"Входных событий из Kafka: {raw_kafka_events.count()}")
# 8 строк (включая 2 технических дубликата)

# ─── ШАГ 1: BRONZE - ТЕХНИЧЕСКАЯ ДЕДУПЛИКАЦИЯ ───────────────────────────────
# Bronze убирает технические дубликаты по event_id
# Не нужен detminism: все дубли идентичны
print("\n=== BRONZE: Техническая дедупликация по event_id ===")

bronze_df = (
    raw_kafka_events
    .dropDuplicates(["event_id"])  # убираем Kafka at-least-once дубли
    .withColumn("_ingested_at", F.current_timestamp())
    .withColumn("ingestion_date", F.current_date())
)

print(f"После Bronze dedup: {bronze_df.count()} строк")
# 6 строк (убрали 2 технических дубля)

spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.bronze.order_events (
        event_id      STRING,
        order_id      STRING,
        lsn           BIGINT,
        event_ts      STRING,
        status        STRING,
        amount        DOUBLE,
        _ingested_at  TIMESTAMP,
        ingestion_date DATE
    )
    USING iceberg
    PARTITIONED BY (ingestion_date)
""")

bronze_df.writeTo("lakehouse.bronze.order_events").append()

print("Bronze таблица:")
spark.table("lakehouse.bronze.order_events") \
     .select("event_id", "order_id", "lsn", "status") \
     .orderBy("order_id", "lsn") \
     .show()
# Видим: для order_1001 - 4 строки (created/paid/shipped/delivered)
# Для order_1002 - 2 строки (created/cancelled)
# Технические дубли убраны

# ─── ШАГ 2: SILVER - БИЗНЕС-ДЕДУПЛИКАЦИЯ ────────────────────────────────────
# Silver берёт только актуальный статус (строку с максимальным LSN)
print("\n=== SILVER: Бизнес-дедупликация через ROW_NUMBER + LSN ===")

window_by_lsn = Window.partitionBy("order_id").orderBy(
    F.col("lsn").desc(),
    F.col("event_id").desc()   # tie-breaker для детерминизма
)

silver_df = (
    spark.table("lakehouse.bronze.order_events")
    .withColumn("event_ts", F.to_timestamp("event_ts"))
    .withColumn("amount",   F.col("amount").cast("decimal(10,2)"))
    .withColumn("rn",       F.row_number().over(window_by_lsn))
    .filter(F.col("rn") == 1)   # только актуальный статус
    .drop("rn", "_ingested_at", "ingestion_date")
    .withColumn("_silver_processed_at", F.current_timestamp())
)

print(f"После Silver dedup: {silver_df.count()} строк")
# 2 строки - по одной на каждый заказ

silver_df.select("order_id", "lsn", "status", "amount").show()
# order_1001 | lsn=104 | delivered | 1500.0   ← последний статус
# order_1002 | lsn=106 | cancelled | 0.0      ← последний статус

spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.silver.orders_latest (
        event_id              STRING,
        order_id              STRING,
        lsn                   BIGINT,
        event_ts              TIMESTAMP,
        status                STRING,
        amount                DECIMAL(10,2),
        _silver_processed_at  TIMESTAMP
    )
    USING iceberg
    PARTITIONED BY (month(event_ts))
""")

silver_df.createOrReplaceTempView("incoming_silver")
spark.sql("""
    MERGE INTO lakehouse.silver.orders_latest t
    USING incoming_silver s ON t.order_id = s.order_id
    WHEN MATCHED AND t.lsn < s.lsn THEN UPDATE SET *
    WHEN NOT MATCHED THEN INSERT *
""")

# ─── ШАГ 3: GOLD - ВЫРУЧКА БЕЗ DOUBLE-COUNTING ──────────────────────────────
print("\n=== GOLD: Выручка по confirmed-заказам без дублей ===")

# Считаем только confirmed/delivered заказы (не cancelled)
gold_df = (
    spark.table("lakehouse.silver.orders_latest")
    .filter(F.col("status").isin("confirmed", "delivered", "paid"))
    .groupBy(F.to_date("event_ts").alias("order_date"))
    .agg(
        F.sum("amount").alias("total_revenue"),
        F.count("order_id").alias("orders_count"),
    )
    .withColumn("_updated_at", F.current_timestamp())
)

spark.sql("""
    CREATE TABLE IF NOT EXISTS lakehouse.gold.daily_revenue_dedup (
        order_date    DATE,
        total_revenue DECIMAL(15,2),
        orders_count  BIGINT,
        _updated_at   TIMESTAMP
    )
    USING iceberg
    PARTITIONED BY (order_date)
""")

gold_df.writeTo("lakehouse.gold.daily_revenue_dedup").overwritePartitions()

print("Gold - итоговая выручка:")
spark.table("lakehouse.gold.daily_revenue_dedup").show()
# order_date  | total_revenue | orders_count
# 2026-05-23  | 1500.00       | 1            ← только delivered заказ
# (cancelled order_1002 не считается)

# ─── ПРОВЕРКА ИДЕМПОТЕНТНОСТИ ─────────────────────────────────────────────────
print("\n=== ПРОВЕРКА: Повторный запуск не меняет результат ===")
# Запускаем Bronze → Silver → Gold снова с теми же данными
# ОЖИДАЕМЫЙ РЕЗУЛЬТАТ: те же 1500.00, те же 2 строки в Silver

# ... (повторные вызовы дадут тот же результат)
print("Gold после повторного запуска:")
spark.table("lakehouse.gold.daily_revenue_dedup").show()
# Тот же результат - pipeline идемпотентен ✓

9.2 Объяснение результата

Входных событий из Kafka: 8
(включает 2 технических дубликата)

=== BRONZE ===
После Bronze dedup: 6 строк
(dropDuplicates убрал evt-001 и evt-005 дубли)

=== SILVER ===
После Silver dedup: 2 строки
(ROW_NUMBER по LSN выбрал:
  - order_1001: lsn=104 delivered (последний из 4 версий)
  - order_1002: lsn=106 cancelled (последний из 2 версий))

=== GOLD ===
total_revenue = 1500.00, orders_count = 1
(только delivered order_1001; cancelled order_1002 исключён)

Ключевые решения в pipeline:

  • Bronze: dropDuplicates(["event_id"]) - быстро, без ordering, достаточно для технических дублей
  • Silver: ROW_NUMBER + ORDER BY lsn.desc(), event_id.desc() - детерминированный выбор актуального статуса
  • Silver sink: MERGE INTO ... WHEN MATCHED AND t.lsn < s.lsn - идемпотентная запись, обновляет только если LSN вырос
  • Gold: filter(status.isin("confirmed", "delivered", "paid")) - бизнес-правило исключения cancelled

10. Выбор подхода: финальная шпаргалка

dropDuplicates - правильный выбор когда:

  • Дубликаты технические (полностью идентичные строки)
  • Не важно, какая строка останется - бизнес-смысл одинаков
  • Нужна высокая производительность без sort-overhead
  • Streaming с watermark: единственный встроенный вариант

Window row_number - правильный выбор когда:

  • Нужна детерминированная бизнес-логика: «последний по LSN», «с максимальным score»
  • CDC pipeline: несколько версий одной сущности
  • Enterprise pipeline: audit trail, воспроизводимость результата
  • Сложные tie-breaking правила

groupBy + max(struct) - альтернатива когда:

  • Нужен детерминизм, но без sort-overhead от Window
  • Достаточно выбора «максимального» по одному числовому полю
  • Данные уже частично агрегированы

Главный принцип: чем ближе к Gold-слою, тем важнее детерминизм. Bronze может использовать dropDuplicates. Silver - всегда Window + deterministic ordering. Gold не занимается дедупликацией - она должна быть решена раньше.