Ответ
Перекос данных (data skew) в Apache Spark — это ситуация, когда данные распределены между партициями RDD или DataFrame крайне неравномерно. В результате задачи (tasks), обрабатывающие перекошенные партиции с огромным объемом данных, выполняются намного дольше, чем остальные. Это создает "длинный хвост" выполнения, простаивание ресурсов и резко снижает общую производительность job, особенно на стадиях shuffle (например, при join или groupBy).
Типичный пример и проблема:
// Предположим, 99% событий имеют user_id = 0 (незалогиненный пользователь)
val eventLogs: DataFrame = ... // DataFrame со столбцом [user_id, event]
// При join или groupBy по user_id все записи с user_id=0 попадут в одну партицию
val skewedAgg = eventLogs.groupBy("user_id").count()
// Задача, обрабатывающая партицию для user_id=0, будет монструозной и медленной.
Основные стратегии борьбы с перекосом:
-
Соление ключей (Salting / Key Augmentation): Добавляем случайный префикс или суффикс к ключам, чтобы разбить "тяжелый" ключ на множество мелких.
import org.apache.spark.sql.functions._ val saltNum = 100 // Количество "солей" // Соление для агрегации val saltedDF = eventLogs .withColumn("salted_key", concat(col("user_id"), lit("_"), (rand() * saltNum).cast("int"))) .groupBy("salted_key") .agg(sum("amount").as("partial_sum")) .withColumn("original_key", split(col("salted_key"), "_").getItem(0)) .groupBy("original_key").agg(sum("partial_sum").as("total_sum")) -
Использование Broadcast Join: Если одна из таблиц для join достаточно мала, чтобы поместиться в память всех executor'ов, Spark может разослать ее копию, избежав дорогостоящего shuffle. Это решает проблему перекоса для этого конкретного join.
// smallDF будет broadcasted, shuffle не произойдет val result = largeDF.join(broadcast(smallDF), "key") -
Оптимизация разделения (Repartitioning): Явное увеличение числа партиций для стадии shuffle может помочь распределить нагрузку лучше, хотя при сильном перекосе одного ключа это не всегда эффективно.
val repartitionedDF = df.repartition(200, col("key")) -
Использование Skew Join Hint (Spark SQL 3.0+): Позволяет явно указать ядру Spark на перекошенные ключи и их значения для специальной оптимизации.
-- Указываем, что в таблице 'orders' для ключа 'customer_id' значения (1, 7, 9) являются перекошенными SELECT /*+ SKEW('orders', 'customer_id', (1, 7, 9)) */ * FROM orders JOIN customers ON orders.customer_id = customers.id;
Выбор стратегии зависит от конкретных данных, размера кластера и типа операции.