PySpark中含OR子句的连接条件下如何避免Broadcast Nested Loop Join
问题现象与原因分析
我已经通过配置spark.sql.autoBroadcastJoinThreshold=-1关闭了自动广播,但遇到以下差异现象:
场景1:单字段等值连接
使用单字段作为连接条件时,执行计划为SortMergeJoin,无广播操作:
data1 = [("123","abc")] schema =(StructType([ StructField("id1",StringType(),True), StructField("id2",StringType(),True)])) data_df1 = spark.createDataFrame(data=data1,schema=schema) data2 = [("123","abc")] schema =(StructType([ StructField("id1",StringType(),True), StructField("id2",StringType(),True)])) data_df2 = spark.createDataFrame(data=data1,schema=schema) join_condition_1 = [(data_df1.id1 == data_df2.id1) ] data_df1.join(data_df2,join_condition_1, "left").explain()
执行计划:
== Physical Plan == SortMergeJoin [id1#308], [id1#312], LeftOuter :- *(1) Sort [id1#308 ASC NULLS FIRST], false, 0 : +- Exchange(coordinator id: 1272911761) hashpartitioning(id1#308, 2000), coordinator[target post-shuffle partition size: 67108864] : +- Scan ExistingRDD[id1#308,id2#309] +- *(3) Sort [id1#312 ASC NULLS FIRST], false, 0 +- Exchange(coordinator id: 1272911761) hashpartitioning(id1#312, 2000), coordinator[target post-shuffle partition size: 67108864] +- *(2) Filter isnotnull(id1#312) +- Scan ExistingRDD[id1#312,id2#313]
场景2:OR条件连接
当连接条件包含OR子句时,执行计划变为BroadcastNestedLoopJoin,触发广播:
join_condition_2 = [(data_df1.id1 == data_df2.id1) | (data_df1.id2 == data_df2.id2) ] data_df1.join(data_df2,join_condition_2, "left").explain()
执行计划:
== Physical Plan == BroadcastNestedLoopJoin BuildRight, LeftOuter, ((id1#308 = id1#312) || (id2#309 = id2#313)) :- Scan ExistingRDD[id1#308,id2#309] +- BroadcastExchange IdentityBroadcastMode +- Scan ExistingRDD[id1#312,id2#313]
核心原因
Spark的SortMergeJoin仅支持基于等值条件的AND组合连接,这类连接可以通过哈希分区将匹配的数据分发到同一分区,再通过排序合并高效完成连接。但当连接条件包含OR时,一条数据可能满足任意一个OR条件,无法通过分区策略将匹配数据聚合到同一分区,因此SortMergeJoin无法处理这类场景。
此时Spark会选择NestedLoopJoin作为备选策略,而NestedLoopJoin有两种实现:
ShuffledNestedLoopJoin:对两张表分区后在分区内做嵌套循环;BroadcastNestedLoopJoin:将小表广播到所有Executor,在每个Executor上做嵌套循环。
即使关闭了autoBroadcastJoinThreshold=-1,Spark在OR条件场景下仍会优先选择BroadcastNestedLoopJoin——因为对于无法分区匹配的场景,广播小表的开销比全表 shuffle 更低,这是当前场景下的最优选择。
Spark 2.4.4中避免触发广播的解决方案
要绕过OR条件导致的BroadcastNestedLoopJoin,可将OR连接拆分为两个独立的等值Join,再通过Union合并结果,让每个子Join都能使用SortMergeJoin,从而避免广播。
示例代码:
# 第一个等值Join:基于id1匹配 join_id1 = data_df1.join(data_df2, data_df1.id1 == data_df2.id1, "left") # 第二个等值Join:基于id2匹配,且排除已被id1匹配的记录(避免重复) # 先筛选出未被id1匹配的主表数据 unmatched_by_id1 = data_df1.join(data_df2, data_df1.id1 == data_df2.id1, "left_anti") # 再基于id2进行左连接 join_id2 = unmatched_by_id1.join(data_df2, unmatched_by_id1.id2 == data_df2.id2, "left") # 合并两个结果 final_result = join_id1.union(join_id2) final_result.explain()
方案说明
- 使用
left_anti过滤已被第一个Join匹配的记录,确保最终结果无重复; - 拆分后的每个子Join都是单字段等值条件,Spark会自动选择
SortMergeJoin,不会触发广播; - 该方案需执行两次Join和一次Union,但这是Spark 2.4.4中处理OR连接条件、避免广播的可行方案。
内容的提问来源于stack exchange,提问作者Ayan Biswas
相关产品推荐
相关产品推荐

