在 Spark 连接中使用广播
请记住,Spark 中的表连接会在集群的各个工作节点之间拆分执行。 如果数据不在本地,就需要进行多次 shuffle 操作,这会对性能产生负面影响。 我们将使用 Spark 的 broadcast 操作,让每个节点都拥有一份指定数据的副本。
几点提示:
- 广播较小的 DataFrame。 DataFrame 越大,传输到工作节点所需的时间越长。
- 对于很小的 DataFrame,可能更适合不进行广播,让 Spark 自行决定优化。
- 如果查看查询执行计划,出现 broadcastHashJoin 说明您已成功配置了广播。
DataFrame flights_df 和 airports_df 已为您提供。
本练习是课程的一部分
使用 PySpark 进行数据清洗
练习说明
- 从
pyspark.sql.functions导入broadcast()方法。 - 通过将
flights_df与airports_df进行连接并使用广播,创建新的 DataFramebroadcast_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.____()