You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 04:53:13