Когда Spark работает медленно (спиллы, шаффлы, перекосы данных)?

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

Ответ

Производительность Spark падает из-за нескольких типовых проблем:

  1. Перекос данных (Data Skew): Неравномерное распределение данных по партициям приводит к тому, что одни задачи выполняются долго, а другие простаивают.

    • Решение: Использовать соление ключей (salting) или увеличить количество партиций.
      // Добавление случайного префикса к ключу для борьбы с перекосом
      val saltedDf = df.withColumn("salted_key", concat(col("key"), lit("_"), (rand() * 100).cast("int")))
  2. Чрезмерные шаффлы: Операции join, groupBy и orderBy вызывают перемешивание данных между узлами — это дорогая операция ввода-вывода.

    • Решение: Использовать broadcast join для маленьких датафреймов, увеличить spark.sql.shuffle.partitions и применять partitionBy при записи.
  3. Спиллы на диск: Когда данные не помещаются в оперативную память исполнителя (executor), Spark записывает их на диск, что резко замедляет работу.

    • Решение: Увеличить spark.executor.memory, настроить spark.memory.fraction и использовать persist(StorageLevel.MEMORY_AND_DISK) для промежуточных датафреймов.
  4. Неоптимальные трансформации: Использование collect() на больших датасетах или пользовательских функций (UDF) там, где можно обойтись встроенными.

    • Решение: Всегда фильтровать данные как можно раньше и минимизировать передачу данных между драйвером и исполнителями.

Для диагностики я в первую очередь анализирую Spark UI, смотрю на время выполнения стадий (Stages), объем данных при шаффле (Shuffle Read/Write) и наличие спиллов (Spill).