Spark и Flink
Средний уровеньПочему возникают ошибки OutOfMemory и как их на самом деле исправить?
«Просто увеличьте память executor'а» решает проблему примерно в половине случаев и впустую тратит деньги в другой половине. Вот как понять, какой у вас случай.
Ошибка OOM говорит вам, что память закончилась. Она не говорит почему — а исправление полностью зависит от причины.
1. Реальная нехватка ресурсов
Простой случай: размер самой большой партиции/окна/состояния действительно не помещается в выделенную память. Это реальная проблема, и решение — правда больше памяти (или меньший размер партиции — см. статью про shuffle-партиции выше).
Сигнал: использование heap стабильно растёт на протяжении жизни задачи/job'а и падает вблизи настроенного лимита, без резкого скачка.
2. Перекошенная партиция, а не заниженное среднее
Если 95% ваших партиций используют 2 ГБ, а одна — 40 ГБ, увеличение памяти для всех ради выживания одного выброса дорого и часто всё равно недостаточно. Настоящее решение — устранить сам перекос: «посолить» (salt) перекошенный ключ join'а или перепартиционировать по более равномерному ключу.
Сигнал: OOM происходит только у небольшого числа задач, а не у большинства, и они коррелируют с конкретным диапазоном ключей или партицией.
3. Broadcast join пошёл не так
Порог broadcast join в Spark (spark.sql.autoBroadcastJoinThreshold) определяет, копируется ли меньшая таблица целиком на каждый executor. Если эта «меньшая» таблица выросла за пределы того, что могут вместить ваши executor'ы — или оценка размера в Spark неверна (часто бывает после фильтра или join'а выше по цепочке, который меняет кардинальность без обновления статистики) — каждый executor падает с OOM, пытаясь удержать broadcast, который на самом деле уже не маленький.
Сигнал: OOM происходит очень рано в стейдже, до реальной работы shuffle, и отключение broadcast join для этого запроса делает его успешным (медленнее, но успешным).
4. Неограниченное состояние в потоковой задаче (специфично для Flink)
Во Flink OOM часто вызван не одним большим батчем — а состоянием, которое не должно было расти бесконечно, но делает именно это: окно, которое никогда не закрывается из-за неверно настроенного watermark, или keyed state без TTL, который бесконечно накапливает ключи. Это не проявится в быстром тесте; это проявится через дни или недели в продакшене, когда heap медленно растёт и не восстанавливается.
Сигнал: размер чекпоинта стабильно растёт на протяжении дней/недель без соответствующего роста объёма входных данных.
Прежде чем увеличивать память: сначала проверьте время GC (Spark) или тренд длительности чекпоинтов (Flink). Оба — сильные индикаторы того, с какой из четырёх причин вы столкнулись, и оба — цифры, которые движок правил opti-pipe уже читает из вашего event log или истории чекпоинтов, а не то, что вам нужно вручную выискивать в Spark/Flink UI.
Посмотрите, как это выглядит на вашем собственном пайплайне.
Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.