Ответ
Производительность Spark падает из-за нескольких типовых проблем:
-
Перекос данных (Data Skew): Неравномерное распределение данных по партициям приводит к тому, что одни задачи выполняются долго, а другие простаивают.
- Решение: Использовать соление ключей (salting) или увеличить количество партиций.
// Добавление случайного префикса к ключу для борьбы с перекосом val saltedDf = df.withColumn("salted_key", concat(col("key"), lit("_"), (rand() * 100).cast("int")))
- Решение: Использовать соление ключей (salting) или увеличить количество партиций.
-
Чрезмерные шаффлы: Операции
join,groupByиorderByвызывают перемешивание данных между узлами — это дорогая операция ввода-вывода.- Решение: Использовать
broadcast joinдля маленьких датафреймов, увеличитьspark.sql.shuffle.partitionsи применятьpartitionByпри записи.
- Решение: Использовать
-
Спиллы на диск: Когда данные не помещаются в оперативную память исполнителя (executor), Spark записывает их на диск, что резко замедляет работу.
- Решение: Увеличить
spark.executor.memory, настроитьspark.memory.fractionи использоватьpersist(StorageLevel.MEMORY_AND_DISK)для промежуточных датафреймов.
- Решение: Увеличить
-
Неоптимальные трансформации: Использование
collect()на больших датасетах или пользовательских функций (UDF) там, где можно обойтись встроенными.- Решение: Всегда фильтровать данные как можно раньше и минимизировать передачу данных между драйвером и исполнителями.
Для диагностики я в первую очередь анализирую Spark UI, смотрю на время выполнения стадий (Stages), объем данных при шаффле (Shuffle Read/Write) и наличие спиллов (Spill).