Использование broadcasting при объединении таблиц в Spark
Помните, что объединение таблиц в Spark распределяется между рабочими узлами кластера. Если данные не находятся локально, требуются операции перемешивания (shuffle), которые могут негативно сказаться на производительности. Вместо этого мы воспользуемся операцией broadcast в Spark, чтобы передать копию указанных данных каждому узлу.
Несколько советов:
- Передавайте через broadcast меньший DataFrame. Чем больше DataFrame, тем больше времени потребуется на его передачу рабочим узлам.
- Для небольших DataFrame broadcasting может быть излишним — лучше дать Spark самому выбрать оптимальную стратегию.
- Если в плане выполнения запроса вы видите broadcastHashJoin, значит, broadcasting настроен успешно.
DataFrame'ы flights_df и airports_df уже доступны для работы.
Это упражнение является частью курса
Очистка данных с помощью PySpark
Инструкции к упражнению
- Импортируйте метод
broadcast()изpyspark.sql.functions. - Создайте новый DataFrame
broadcast_df, объединивflights_dfсairports_dfс использованием broadcasting. - Выведите план запроса и обратите внимание на отличия от исходного варианта.
Интерактивное практическое упражнение
Попробуйте выполнить это упражнение, дополнив этот пример кода.
# Import the broadcast method from pyspark.sql.functions
from ____ import ____
# Join the flights_df and airports_df DataFrames using broadcasting
broadcast_df = flights_df.____(____(airports_df), \
flights_df["Destination Airport"] == airports_df["IATA"] )
# Show the query plan and compare against the original
broadcast_df.____()