Способы дедупликации в 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()): для случаев, когда нужно взять первое значение из группы по определенному столбцу будет эффективнее оконной функции.
А как вы удаляете дубликаты?
©️что-то на инженерном