Способы дедупликации в Spark

Я тут наткнулась на статью с провокационным названием Stop Using dropDuplicates()! и не смогла пройти мимо нее.

Честно говоря, в большинстве случаев я использую либо dropDuplicates(), если мне необходимо удалить дубли по всем или выбранным столбцам в датафрейме, либо groupBy + count()/agg().

1️⃣ В статье утверждается, что стандартная функция dropDuplicates() в PySpark вызывает глобальный шаффл, что ведет к резкому падению производительности на больших объёмах данных. Это может стать серьезной проблемой при дедупликации миллиардных датасетов, т.к. перегруженные партиции создают перекосы данных, приводя к out-of-memory ошибкам и сбоям воркеров.

Чтобы гарантировать удаление всех дубликатов по всему датафрейму, dropDuplicates() / distinct() делает полный шаффл данных. Spark должен сгруппировать все потенциально дублирующиеся строки вместе на одной партиции для их сравнения и удаления. Это перемешивание происходит по ключам, определенным столбцами в dropDuplicates(subset=[…]) или по всем столбцам для пустого вызова.

2️⃣ Автор рекомендуют вместо dropDuplicates() применять оконные функции, а также предварительно анализировать распределение ключей и разумно репартиционировать данные.

В чем разница оконных функций и dropDuplicates()? Ключевая идея автора в том, что при использовании оконной функции (row_number() и последующей фильтрации по rn=1) вы явно контролируете, какой дубликат (например, с наибольшей/наименьшей датой или ID) будет сохранен. Кроме того, перемешивание происходит только по столбцам, указанным в partitionBy. Если данные уже были разделены (например, с помощью repartition()) по тем же ключам, шаффл может быть пропущен.

3️⃣ В случае перекоса данных - автор предлагает применить salting ключей, т.е. добавления небольшого случайного префикса или суффикса к ключу перед агрегацией/дедупликацией, а затем на втором этапе десолтировать (удалить соль) и агрегировать уже по исходному ключу.

Что в итоге?

Я бы не сказала, что эта статья - универсальная рекомендация, если хочется оптимизировать дедупликацию. Все очень сильно зависит от вашей задачи и данных. Я придерживаюсь следующего мнения:

🔸на небольшом объеме данных можно не парится и применять dropDuplicates().

🔸если нужно очистить все строки по всем полям на большом объеме, то оконная функция - сомнительное решение, она также приведет к глобальному шаффлу. Действительно полезно в случае, когда нам нужен контроль над строками.

🔸оконная функция по ключевым полям и dropDuplicates()с указанными полями будут работать практически идентично.

🔸если соль распределена неверно, то салтинг данных может привести к формированию некорректных ключей, и как следствие, к некорректной дедупликации. Требует аккуратности.

🔸groupBy() + agg(first()): для случаев, когда нужно взять первое значение из группы по определенному столбцу будет эффективнее оконной функции.

А как вы удаляете дубликаты?

©️что-то на инженерном