Spark

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

Как обнаружить и исправить перекос данных (data skew)?

Перекос невидим в агрегированных метриках и очевиден в метриках по отдельным задачам — нужно просто смотреть в правильное место.

Читать на:
Задача 1 Задача 2 Задача 3 Задача 4 (перекошенная) Задача 5
Перекошенный стейдж, иллюстрация — одна задача выполняет намного больше работы, чем остальные.

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

Как это действительно увидеть

Средняя и даже p95 длительность задач могут выглядеть нормально, даже когда перекос — реальная проблема: он часто затрагивает достаточно малую долю задач, чтобы процентильная статистика его скрывала. Что не скрывает: соотношение между максимальной и медианной длительностью задачи внутри одного стейджа. У здорового shuffle-стейджа это соотношение заметно ниже 3x. Соотношение выше 10x — это перекос, а не шум.

Исправление 1: «соление» перекошенного ключа

Для join'а или агрегации, перекошенных по конкретному значению ключа (один ID клиента, один регион, один статус-флаг доминирует по объёму), «соление» (salting) распределяет строки этого ключа по нескольким синтетическим под-ключам перед shuffle, а затем объединяет результаты после. Это самое точечное решение, когда известно, какой ключ — проблема, но оно требует изменения кода по обеим сторонам join'а.

Исправление 2: Adaptive Query Execution (AQE)

Если вы на Spark 3.0+, spark.sql.adaptive.enabled=true (включено по умолчанию в последних версиях) включает автоматическую обработку перекошенных join'ов — Spark обнаруживает перегруженную партицию во время выполнения и разбивает её перед join'ом, без изменения кода. Это первое, что стоит проверить перед ручным «солением»: убедиться, что AQE действительно включён, а порог обработки перекоса не выставлен настолько высоко, что ваш случай его не задействует.

Исправление 3: broadcast меньшей стороны вместо этого

Если перекос на самом деле — это дисбаланс размеров join'а, а не перекос внутри столбца, и одна из таблиц достаточно мала для broadcast — переход на broadcast join полностью обходит shuffle (и перекос вместе с ним). Применимо только если меньшая таблица действительно укладывается в лимиты памяти для broadcast join; см. статью про OOM выше о том, что происходит, когда не укладывается.

Где это видно в opti-pipe: разброс длительности задач из вашего event log Spark — один из сигналов, которые движок правил читает напрямую: когда соотношение максимум/медиана в shuffle-стейдже превышает порог перекоса, это отмечается как конкретная рекомендация, а не то, что вам пришлось бы заметить самостоятельно в Spark UI.

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

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