Spark
Средний уровеньКогда broadcast join в Spark реально помогает, а когда — вредит?
Broadcast join — это ставка на то, что одна сторона join'а достаточно мала, чтобы дёшево скопировать её везде. Когда ставка верна — это самый быстрый join в Spark. Когда нет — самый быстрый способ уронить driver по OOM.
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 — и
получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.