❗️Вопрос с собеса в BigTech:
«Мы удвоили количество executor ов, но Spark джоба не ускорилась. Почему?»
Это классическая ловушка для тех, кто пытается лечить любую проблему масштабированием 🥸
У тебя есть медленный pipeline. Ты увеличиваешь количество executor ов в 2 раза, перезапускаешь job, а время выполнения остаётся прежним или даже растёт.
Что происходит?
1️⃣Data skew 🥴
Если один ключ содержит 90% данных, все связанные с ним записи могут попасть в одну shuffle-partition. Один task будет работать 20 минут, пока остальные завершатся за секунды.
В таком случае дополнительные executor’ы не помогут: job ждёт самую медленную task 🐌
В Spark UI сравни Max и Median по длительности tasks и объёму Shuffle Read .
Большой разрыв — сигнал проверить перекос данных❗️
2️⃣Узкое место в shuffle🔀
JOIN , GROUP BY и repartition могут вызвать shuffle — перераспределение данных между executor’ами.
Если bottleneck (узкое место) находится в shuffle, добавление машин не убирает саму пересылку, сортировку и запись промежуточных данных на диск.
Сначала нужно проверить план и объём shuffle, а затем рассмотреть broadcast join, предварительную агрегацию, фильтрацию до join и настройку числа shuffle-partitions.
3️⃣Недостаточно задач для параллелизма🟰
Количество executor’ов само по себе не создаёт работу. Если в stage всего несколько крупных tasks, дополнительные executor’ы будут простаивать.
Поэтому нужно смотреть не только на число executor’ов, но и на количество tasks, размер partition и соотношение доступных ядер к параллельной работе.
4️⃣Слишком много мелких файлов💛
Миллион файлов по 10 КБ — это не много вычислений, а много служебных операций: listing, планирование и запуск tasks.
Масштабирование кластера проблему не устранит. Нужно менять стратегию записи и объединять мелкие файлы.
5️⃣Memory pressure, GC и spill🕳
Если tasks обрабатывают слишком большие partition, executor может тратить время на сборку мусора или сбрасывать промежуточные данные на диск ( spill ).
Spill не всегда означает ошибку — это допустимый механизм Spark. Но если он массовый, а GC Time высокий, нужно проверить размер partition, shuffle и конфигурацию памяти.
Как отвечать на собеседовании 🧐
Я бы сказал:
☑️«Сначала открою Spark UI и найду самый долгий stage».
☑️«Сравню Max и Median длительности tasks».
☑️ «Проверю Shuffle Read/Write , spill и GC Time ».
☑️ «Посмотрю, хватает ли tasks для загруженных ядер».
☑️«Только после диагностики буду менять конфигурацию или добавлять ресурсы».
Главная мысль: больше executor’ов ускоряют job только тогда, когда есть достаточно независимой работы. Если причина в skew, shuffle, мелких файлах или memory pressure, масштабирование лишь увеличит стоимость, но не устранит bottleneck.
📚Для изучения:
⏺ Руководство для начинающих по Spark UI
⏺Оптимизируем Shuffle в Spark
Ставь 🔥, если было полезно!
«Мы удвоили количество executor ов, но Spark джоба не ускорилась. Почему?»
Это классическая ловушка для тех, кто пытается лечить любую проблему масштабированием 🥸
У тебя есть медленный pipeline. Ты увеличиваешь количество executor ов в 2 раза, перезапускаешь job, а время выполнения остаётся прежним или даже растёт.
Что происходит?
1️⃣Data skew 🥴
Если один ключ содержит 90% данных, все связанные с ним записи могут попасть в одну shuffle-partition. Один task будет работать 20 минут, пока остальные завершатся за секунды.
В таком случае дополнительные executor’ы не помогут: job ждёт самую медленную task 🐌
В Spark UI сравни Max и Median по длительности tasks и объёму Shuffle Read .
Большой разрыв — сигнал проверить перекос данных❗️
2️⃣Узкое место в shuffle🔀
JOIN , GROUP BY и repartition могут вызвать shuffle — перераспределение данных между executor’ами.
Если bottleneck (узкое место) находится в shuffle, добавление машин не убирает саму пересылку, сортировку и запись промежуточных данных на диск.
Сначала нужно проверить план и объём shuffle, а затем рассмотреть broadcast join, предварительную агрегацию, фильтрацию до join и настройку числа shuffle-partitions.
3️⃣Недостаточно задач для параллелизма🟰
Количество executor’ов само по себе не создаёт работу. Если в stage всего несколько крупных tasks, дополнительные executor’ы будут простаивать.
Поэтому нужно смотреть не только на число executor’ов, но и на количество tasks, размер partition и соотношение доступных ядер к параллельной работе.
4️⃣Слишком много мелких файлов💛
Миллион файлов по 10 КБ — это не много вычислений, а много служебных операций: listing, планирование и запуск tasks.
Масштабирование кластера проблему не устранит. Нужно менять стратегию записи и объединять мелкие файлы.
5️⃣Memory pressure, GC и spill🕳
Если tasks обрабатывают слишком большие partition, executor может тратить время на сборку мусора или сбрасывать промежуточные данные на диск ( spill ).
Spill не всегда означает ошибку — это допустимый механизм Spark. Но если он массовый, а GC Time высокий, нужно проверить размер partition, shuffle и конфигурацию памяти.
Как отвечать на собеседовании 🧐
Я бы сказал:
☑️«Сначала открою Spark UI и найду самый долгий stage».
☑️«Сравню Max и Median длительности tasks».
☑️ «Проверю Shuffle Read/Write , spill и GC Time ».
☑️ «Посмотрю, хватает ли tasks для загруженных ядер».
☑️«Только после диагностики буду менять конфигурацию или добавлять ресурсы».
Главная мысль: больше executor’ов ускоряют job только тогда, когда есть достаточно независимой работы. Если причина в skew, shuffle, мелких файлах или memory pressure, масштабирование лишь увеличит стоимость, но не устранит bottleneck.
📚Для изучения:
⏺ Руководство для начинающих по Spark UI
⏺Оптимизируем Shuffle в Spark
Ставь 🔥, если было полезно!