Что такое перекос данных (data skew) в Apache Spark?

«Что такое перекос данных (data skew) в Apache Spark?» — вопрос из категории Apache Spark, который задают на 33% собеседований Data Инженер. Ниже — развёрнутый ответ с разбором ключевых моментов.

Ответ

Перекос данных (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, будет монструозной и медленной.

Основные стратегии борьбы с перекосом:

  1. Соление ключей (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"))
  2. Использование Broadcast Join: Если одна из таблиц для join достаточно мала, чтобы поместиться в память всех executor'ов, Spark может разослать ее копию, избежав дорогостоящего shuffle. Это решает проблему перекоса для этого конкретного join.

    // smallDF будет broadcasted, shuffle не произойдет
    val result = largeDF.join(broadcast(smallDF), "key")
  3. Оптимизация разделения (Repartitioning): Явное увеличение числа партиций для стадии shuffle может помочь распределить нагрузку лучше, хотя при сильном перекосе одного ключа это не всегда эффективно.

    val repartitionedDF = df.repartition(200, col("key"))
  4. Использование 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;

Выбор стратегии зависит от конкретных данных, размера кластера и типа операции.