Spark и Flink

Средний уровень

Почему возникают ошибки OutOfMemory и как их на самом деле исправить?

«Просто увеличьте память executor'а» решает проблему примерно в половине случаев и впустую тратит деньги в другой половине. Вот как понять, какой у вас случай.

Читать на:
Реальная нехватка ресурсов Перекошенная партиция Broadcast join пошёл не так Неограниченное состояние
Примерное ранжирование по частоте, с которой это реальная причина OOM — не измерено.

Ошибка 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 — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.