Spark

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

Когда broadcast join в Spark реально помогает, а когда — вредит?

Broadcast join — это ставка на то, что одна сторона join'а достаточно мала, чтобы дёшево скопировать её везде. Когда ставка верна — это самый быстрый join в Spark. Когда нет — самый быстрый способ уронить driver по OOM.

Читать на:
Корректный broadcast Устаревшая оценка размера Фильтр не продавлен
Примерно как часто каждый сценарий объясняет неожиданность с broadcast join — иллюстративно, не измерено.

spark.sql.autoBroadcastJoinThreshold (по умолчанию 10 МБ) говорит планировщику: если оценённый размер таблицы меньше этого порога, пропустить shuffle и вместо этого отправить полную копию на каждый executor.

Что на самом деле контролирует порог

Ниже порога планировщик заменяет shuffle join (обе стороны репартиционируются и перемешиваются по сети) на broadcast join (маленькая сторона собирается на driver, затем передаётся каждому executor'у как hash-таблица). Отсутствие shuffle для большой стороны означает отсутствие shuffle-read стейджа, отсутствие связанного с shuffle перекоса и обычно заметно более быстрый join — именно поэтому планировщик делает это автоматически, как только оценивает таблицу как достаточно маленькую.

Отказ: broadcast там, где должен был остаться shuffle

Оценка, которой пользуется планировщик, берётся из статистики таблицы, которая устаревает — таблица, которая была 8 МБ на момент последнего расчёта статистики, а сейчас весит 800 МБ, всё равно попадёт под broadcast, если статистику никто не обновил. Результат: каждый executor пытается одновременно держать в памяти хеш-таблицу на 800 МБ, и часто именно driver (который сначала собирает данные для broadcast) падает по OOM раньше любого executor'а.

Как понять это по event log

Ищите стейдж с очень малым числом задач (стейдж сбора broadcast-данных фактически выполняется на одной задаче) с необычно высоким пиковым потреблением памяти, сразу за которым следует падение на уровне стейджа или OOM драйвера в логах приложения, а не отказ отдельной задачи. OOM при shuffle join, наоборот, выглядит как падение множества задач в стейдже нормального размера — совсем другая картина.

Если не уверены, происходит ли broadcast вообще: проверьте физический план на BroadcastHashJoin против SortMergeJoin — это единственное место, где Spark прямо говорит, какую стратегию выбрал, вместо того чтобы догадываться по таймингу.

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

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