Spark

Senior

Спекулятивное выполнение в Spark: страховка или скрытый множитель расходов?

spark.speculation перезапускает задачу, работающую заметно медленнее остальных, на теории, что дело в узле, на котором она выполняется, а не в самой задаче. Эта теория иногда неверна.

Читать на:
Временная проблема узла Перекос данных (похоже) Реально медленная логика
Чем обычно оказывается паттерн медленной задачи — иллюстративно, не измерено.

При включённой спекуляции Spark следит за задачами, работающими заметно медленнее медианы в своём стейдже, и запускает дублирующую попытку в другом месте, оставляя ту, что завершится первой.

Случай, где это работает как надо

Одна задача, застрявшая на уровне в 10 раз выше медианы, в то время как все остальные задачи стейджа завершаются нормально — учебный случай: сбойный узел, шумный сосед на общей инфраструктуре, временный сбой диска или сети. Спекуляция запускает дубликат, дубликат завершается на здоровом узле, и задача восстанавливается без чьего-либо вмешательства — именно та страховка, для которой она и задумана.

Случай, где становится хуже

Перекос данных даёт тот же внешний симптом — одна задача заметно медленнее остальных — по совершенно другой причине: этой задаче просто досталось больше данных для обработки. Спекуляция запускает дубликат задачи, которая изначально не могла завершиться быстро ни на каком узле, так что вы платите за две медленные попытки вместо одной, без улучшения реального времени завершения. Перекошенная задача всё равно завершается последней в любом случае.

Как различить их до того, как включать переключатель

Проверяйте объём входных данных, а не только длительность, у медленной задачи по сравнению с соседями в том же стейдже. Примерно равный объём входных данных с одной сильно более медленной задачей указывает на инфраструктурную проблему, с которой спекуляция реально поможет. Медленная задача с заметно большим объёмом входных данных (из "Input Metrics" или байтов shuffle read в event log) — это перекос, и фикс здесь — репартиционирование или «соление» ключа join'а, а не спекуляция, которая просто удваивает стоимость задачи, изначально обречённой быть самой долгой.

Эмпирическое правило: спекуляция — разумное значение по умолчанию для инфраструктурного шума, но не замена реальному поиску и исправлению перекоса — а на задаче, где перекос и есть настоящая проблема, она тихо увеличивает стоимость, ничего не давая для времени выполнения.

Посмотрите, как это выглядит на вашем собственном пайплайне.

Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.