Перейти к содержимому
9 / 15

Джойны в Spark: broadcast, sort-merge, shuffle hash — как выбирается стратегия

Стратегий джойна несколько, и выбирает их Catalyst по оценкам размеров. BroadcastHashJoin — маленькая сторона целиком рассылается на все исполнители, shuffle не нужен; включается автоматически, если оценка размера меньше spark.sql.autoBroadcastJoinThreshold (10 МБ по умолчанию). SortMergeJoin — стратегия по умолчанию для двух больших сторон: обе перемешиваются по ключу, сортируются и сливаются; надёжна, но требует двух shuffle. ShuffleHashJoin строит хеш-таблицу вместо сортировки и выигрывает, когда одна сторона заметно меньше другой, но всё же не влезает в broadcast. BroadcastNestedLoopJoin и CartesianProduct появляются при джойне без условия равенства и означают перебор пар — на больших данных это приговор. Повлиять можно подсказками: broadcast(df) или hint("merge"), а начиная с версии 3.0 доступны все четыре подсказки. При включённом AQE решение может быть пересмотрено уже во время выполнения по фактическим размерам.

Джойны в Spark: broadcast, sort-merge, shuffle hash — как выбирается стратегия | JScriptiser