为何结构完全一致的Spark DataFrame执行join操作会得到不同结果
问题背景
更新:该问题的根源是一个bug,已于Spark 3.2.0版本修复。
两次运行输入的DataFrame结构完全相同,但输出结果存在差异,仅第二次运行返回了预期结果df6,已知可以通过给DataFrame设置别名得到正确结果。
问题
Spark生成df3的底层运行机制是什么?join的on子句中明确写了df1.c1 == df2.c2,但Spark显然没有正确识别指定的两个DataFrame,其底层逻辑是什么?如何预判这类行为?
第一次运行(df3结果错误)
data = [ (1, 'bad', 'A'), (4, 'ok', None)] df1 = spark.createDataFrame(data, ['ID', 'Status', 'c1']) df1 = df1.withColumn('c2', F.lit('A')) df1.show() #+---+------+----+---+ #| ID|Status| c1| c2| #+---+------+----+---+ #| 1| bad| A| A| #| 4| ok|null| A| #+---+------+----+---+ df2 = df1.filter((F.col('Status') == 'ok')) df2.show() #+---+------+----+---+ #| ID|Status| c1| c2| #+---+------+----+---+ #| 4| ok|null| A| #+---+------+----+---+ df3 = df2.join(df1, (df1.c1 == df2.c2), 'full') df3.show() #+----+------+----+----+----+------+----+----+ #| ID|Status| c1| c2| ID|Status| c1| c2| #+----+------+----+----+----+------+----+----+ #| 4| ok|null| A|null| null|null|null| #|null| null|null|null| 1| bad| A| A| #|null| null|null|null| 4| ok|null| A| #+----+------+----+----+----+------+----+----+
第二次运行(df6结果正确)
data = [ (1, 'bad', 'A', 'A'), (4, 'ok', None, 'A')] df4 = spark.createDataFrame(data, ['ID', 'Status', 'c1', 'c2']) df4.show() #+---+------+----+---+ #| ID|Status| c1| c2| #+---+------+----+---+ #| 1| bad| A| A| #| 4| ok|null| A| #+---+------+----+---+ df5 = spark.createDataFrame(data, ['ID', 'Status', 'c1', 'c2']).filter((F.col('Status') == 'ok')) df5.show() #+---+------+----+---+ #| ID|Status| c1| c2| #+---+------+----+---+ #| 4| ok|null| A| #+---+------+----+---+ df6 = df5.join(df4, (df4.c1 == df5.c2), 'full') df6.show() #+----+------+----+----+---+------+----+---+ #| ID|Status| c1| c2| ID|Status| c1| c2| #+----+------+----+----+---+------+----+---+ #|null| null|null|null| 4| ok|null| A| #| 4| ok|null| A| 1| bad| A| A| #+----+------+----+----+---+------+----+---+
执行计划差异
df3.explain() == Physical Plan == BroadcastNestedLoopJoin BuildRight, FullOuter, (c1#23335 = A) :- *(1) Project [ID#23333L, Status#23334, c1#23335, A AS c2#23339] : +- *(1) Filter (isnotnull(Status#23334) AND (Status#23334 = ok)) : +- *(1) Scan ExistingRDD[ID#23333L,Status#23334,c1#23335] +- BroadcastExchange IdentityBroadcastMode, [id=#9250] +- *(2) Project [ID#23379L, Status#23380, c1#23381, A AS c2#23378] +- *(2) Scan ExistingRDD[ID#23379L,Status#23380,c1#23381] df6.explain() == Physical Plan == SortMergeJoin [c2#23459], [c1#23433], FullOuter :- *(2) Sort [c2#23459 ASC NULLS FIRST], false, 0 : +- Exchange hashpartitioning(c2#23459, 200), ENSURE_REQUIREMENTS, [id=#9347] : +- *(1) Filter (isnotnull(Status#23457) AND (Status#23457 = ok)) : +- *(1) Scan ExistingRDD[ID#23456L,Status#23457,c1#23458,c2#23459] +- *(4) Sort [c1#23433 ASC NULLS FIRST], false, 0 +- Exchange hashpartitioning(c1#23433, 200), ENSURE_REQUIREMENTS, [id=#9352] +- *(3) Scan ExistingRDD[ID#23431L,Status#23432,c1#23433,c2#23434]
两次运行的物理计划存在差异,内部使用了不同的join实现(BroadcastNestedLoopJoin和SortMergeJoin),但这一点本身无法解释结果差异,因为不同的内部join实现返回的结果应该一致。
问题解答
底层逻辑
这个问题是Spark 3.2.0之前版本的已知bug,核心出在join条件的解析逻辑上:
- 第一次运行中,
df2是df1经过filter衍生出的DataFrame,c2列是通过F.lit('A')生成的常量列。老版本Spark在解析df1.c1 == df2.c2这个条件时,错误地将衍生出来的常量列df2.c2直接替换为了字面量A,同时没有正确区分该列所属的DataFrame,最终实际生效的join条件变成了c1 = A,和预期的两表列匹配逻辑完全不同,才会出现错误的全连接结果。 - 第二次运行中,
df4和df5是两个完全独立创建的DataFrame,没有衍生关系,列的唯一标识ID完全独立,Spark不会做错误的常量替换,join条件被正确解析为两表的列匹配,因此结果符合预期。
执行计划里的join类型差异只是表象,本质是join条件解析错误导致的逻辑差异。
规避和预判方法
- 版本低于3.2.0的Spark环境下,只要join的两个DataFrame存在衍生关系,一律给两个DataFrame设置别名,通过别名指定join列,写法参考:
df3 = df2.alias("df2").join(df1.alias("df1"), F.col("df1.c1") == F.col("df2.c2"), "full")
- 写join条件时优先使用带别名的
F.col写法,减少列解析的歧义。 - 遇到join结果不符合预期时,优先调用
explain()方法查看实际生效的join条件,只要实际条件和编写逻辑不一致,就可以快速定位到解析问题。
内容的提问来源于stack exchange,提问作者ZygD
相关产品推荐
相关产品推荐

