Deduplication: dropDuplicates vs Window row_number - когда и что выбрать
Deduplication: dropDuplicates vs Window row_number - когда и что выбрать
Дублирование данных - одна из самых коварных проблем в 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
Что происходит шаг за шагом:
- Scan: каждый executor читает свою часть данных локально
- HashAggregate (partial): локально на каждом executor строится хэш-таблица по
event_id. Если строка с такимevent_idуже видена - она отбрасывается. Это reduces shuffle volume - Exchange (Shuffle): строки перераспределяются между executors так, чтобы все строки с одинаковым
event_idпопали на один executor (hash partitioning) - 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:
- Shuffle по
partitionBy: аналогично dropDuplicates, но причина другая - нужно, чтобы все строки одногоorder_idпопали на один executor для корректного вычисления ранга - Sort: локальная сортировка внутри каждого executor по
(order_id, lsn DESC). Это дополнительная O(n log n) операция, которой нет в dropDuplicates - Window computation: вычисление
row_numberдля каждой строки в отсортированной партиции - 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 не занимается дедупликацией - она должна быть решена раньше.