НачатьНачать бесплатно

Использование 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.____()
Редактировать и запускать код