Что делать, если твой Spark-кластер из 200 воркеров работает как один 🐌
Запускаешь джобу на 200 воркеров, а вся нагрузка едет на одном узле, пока остальные 199 простаивают🥺
Это data skew - неравномерное распределение данных по партициям. Spark делит данные по ключу (обычно при join или groupBy), и если один ключ встречается в разы чаще остальных - весь воркер, который его обрабатывает, становится бутылочным горлышком. Время выполнения джобы определяет не средний воркер, а самый перегруженный.
Заметить просто: открой Spark UI - один или несколько тасков выполняются заметно дольше остальных, иногда доходит до OOM (out of memory) на конкретном экзекьюторе.
Что с этим делать 🤔
Adaptive Query Execution - начни отсюда ▶️
В Spark 3.x можно включить adaptive query execution со skew join handling прямо в конфиге. Spark сам на лету разобьёт перегруженные партиции. Закрывает большинство случаев без единой строчки кода.
Broadcast join, если одна таблица маленькая📊
Можно разослать маленький датасет на все экзекьюторы и обойтись без шаффла полностью. Часто проще и быстрее, чем что-либо солить.
Salting, если ничего не помогло🧂
Добавляешь к перегруженному ключу случайное число - соль. Один огромный партишн превращается в несколько поменьше. Для join большую таблицу солишь, а маленькую размножаешь под каждое значение соли. Это последний инструмент, не первый — усложняет код, и если есть способ попроще, лучше им и обойтись.
🔖полезно почитать дополнительно
Если было полезно, ставь 🔥
Запускаешь джобу на 200 воркеров, а вся нагрузка едет на одном узле, пока остальные 199 простаивают🥺
Это data skew - неравномерное распределение данных по партициям. Spark делит данные по ключу (обычно при join или groupBy), и если один ключ встречается в разы чаще остальных - весь воркер, который его обрабатывает, становится бутылочным горлышком. Время выполнения джобы определяет не средний воркер, а самый перегруженный.
Заметить просто: открой Spark UI - один или несколько тасков выполняются заметно дольше остальных, иногда доходит до OOM (out of memory) на конкретном экзекьюторе.
Что с этим делать 🤔
Adaptive Query Execution - начни отсюда ▶️
В Spark 3.x можно включить adaptive query execution со skew join handling прямо в конфиге. Spark сам на лету разобьёт перегруженные партиции. Закрывает большинство случаев без единой строчки кода.
Broadcast join, если одна таблица маленькая📊
Можно разослать маленький датасет на все экзекьюторы и обойтись без шаффла полностью. Часто проще и быстрее, чем что-либо солить.
Salting, если ничего не помогло🧂
Добавляешь к перегруженному ключу случайное число - соль. Один огромный партишн превращается в несколько поменьше. Для join большую таблицу солишь, а маленькую размножаешь под каждое значение соли. Это последний инструмент, не первый — усложняет код, и если есть способ попроще, лучше им и обойтись.
🔖полезно почитать дополнительно
Если было полезно, ставь 🔥