Ответ
Broadcast join — это стратегия оптимизации операции соединения (join) в Spark, при которой небольшая таблица (DataFrame) полностью копируется на все рабочие узлы (executors) кластера. Это позволяет избежать дорогостоящей перетасовки данных (shuffle) для большой таблицы, так как соединение теперь может выполняться локально на каждом узле.
Когда использовать?
Когда одна из таблиц значительно меньше другой (обычно меньше 10 МБ, что определяется параметром spark.sql.autoBroadcastJoinThreshold).
Пример на PySpark:
from pyspark.sql import SparkSession
from pyspark.sql.functions import broadcast
spark = SparkSession.builder.appName("BroadcastJoinExample").getOrCreate()
# Небольшой справочник (например, таблица стран)
dim_country = spark.createDataFrame([
(1, "USA"),
(2, "UK"),
(3, "Germany")
], ["country_id", "country_name"])
# Большая фактовая таблица (например, логи транзакций)
fact_transactions = spark.createDataFrame([
(100, 1, 150.0),
(101, 2, 89.99),
(102, 3, 200.5),
(103, 1, 55.0)
], ["transaction_id", "country_id", "amount"])
# Явное указание на broadcast join с помощью функции broadcast()
result_df = fact_transactions.join(broadcast(dim_country), on="country_id")
result_df.show()
# +-------------+----------+------+-------------+
# | transaction_id| country_id| amount| country_name|
# +-------------+----------+------+-------------+
# | 100| 1| 150.0| USA|
# | 103| 1| 55.0| USA|
# | 101| 2| 89.99| UK|
# | 102| 3| 200.5| Germany|
# +-------------+----------+------+-------------+
Ключевые моменты:
- Автоматика: Spark Catalyst Optimizer может автоматически выбрать broadcast join, если размер маленькой таблицы меньше порога.
- Принудительное вещание: Функция
broadcast()используется для явного указания, даже если таблица чуть больше порога (но вы уверены, что это эффективно). - Предостережение: Попытка broadcast очень большой таблицы может привести к нехватке памяти (OutOfMemoryError) на исполнителях, так как её копия будет храниться в памяти каждого узла.