Spark
JuniorПочему моя Spark-задача вдруг стала медленнее без изменений в коде?
В девяти случаях из десяти дело не в коде. Дело в данных, кластере или конфигурации, которая незаметно перестала соответствовать одному из них.
Вы не трогали задачу. DAG идентичен. Последний деплой был три недели назад. И тем не менее вчерашний ночной прогон занял 40 минут вместо 12. Это один из самых частых тикетов в data engineering, и он почти всегда сводится к одной из трёх причин.
1. Данные выросли, а конфигурация — нет
Число shuffle-партиций, память executor'а и порог broadcast join обычно задаются один раз, на этапе первоначальной разработки, исходя из объёма данных на тот момент. Через полгода та же таблица выросла в 4 раза, а spark.sql.shuffle.partitions=200 остался прежним — теперь он создаёт партиции слишком большие, чтобы комфортно помещаться в память, что заставляет данные сбрасываться на диск при каждом shuffle-этапе. В коде ничего не менялось; изменилось предположение, заложенное в конфигурации.
Как проверить: сравните число входных строк и объём shuffle read/write между недавним медленным прогоном и более старым быстрым, во вкладке Stages в Spark UI. Резкий скачок здесь, с примерно пропорциональным замедлением, прямо указывает на эту причину.
2. Join, который раньше был сбалансирован, перестал им быть
Перекос данных не заявляет о себе — он просто означает, что горстка партиций выполняет в 10-100 раз больше работы, чем остальные, так что общее время определяется самой медленной задачей, а не средней. Распределение ключей, которое было примерно равномерным при запуске, могло сместиться по мере изменения паттернов использования (один клиент, один регион, один тип события начинает доминировать по объёму).
Как проверить: в детальном виде стейджа в Spark UI отсортируйте задачи по длительности. Стейдж, где медиана — 4 секунды, а максимум — 6 минут, это перекос, без всяких сомнений.
3. Кластер — уже не тот кластер, на котором вы тестировали
Автомасштабируемые кластеры, пулы spot-инстансов и общие мультитенантные кластеры могут незаметно выделять вашей задаче меньше executor'ов или более медленные, особенно при конкуренции с другими задачами. Это проявляется как замедление по времени выполнения при идентичном CPU-времени на задачу — задача выполняет тот же объём работы, просто дольше ждёт ресурсов для этого.
Как проверить: сравните число реально выделенных (не запрошенных) executor'ов между прогонами и посмотрите на задержку планировщика задач, а не только на длительность задач.
Самый быстрый способ отличить эти причины друг от друга: нужны две вещи рядом — ваша текущая конфигурация и реальные метрики прогона, который действительно произошёл. Именно на этом строится движок правил opti-pipe: он читает реальное время выполнения задач, объём shuffle и использование памяти из вашего event log и сверяет это с настроенным числом shuffle-партиций, памятью executor'а и числом инстансов — вместо того чтобы вы сами разглядывали Spark UI и гадали, какая из трёх причин перед вами.
Посмотрите, как это выглядит на вашем собственном пайплайне.
Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.