开始使用免费开始使用

在 Spark 连接中使用广播

请记住,Spark 中的表连接会在集群的各个工作节点之间拆分执行。 如果数据不在本地,就需要进行多次 shuffle 操作,这会对性能产生负面影响。 我们将使用 Spark 的 broadcast 操作,让每个节点都拥有一份指定数据的副本。

几点提示:

  • 广播较小的 DataFrame。 DataFrame 越大,传输到工作节点所需的时间越长。
  • 对于很小的 DataFrame,可能更适合不进行广播,让 Spark 自行决定优化。
  • 如果查看查询执行计划,出现 broadcastHashJoin 说明您已成功配置了广播。

DataFrame flights_dfairports_df 已为您提供。

本练习是课程的一部分

使用 PySpark 进行数据清洗

查看课程

练习说明

  • pyspark.sql.functions 导入 broadcast() 方法。
  • 通过将 flights_dfairports_df 进行连接并使用广播,创建新的 DataFrame broadcast_df
  • 显示查询计划,并与原始版本进行对比思考差异。

交互式实操练习

通过完成这段示例代码来试试这个练习。

# 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.____()
编辑并运行代码