PySpark中DataFrame与其自过滤结果执行left_anti join失效问题
问题解答
底层原因
这个问题的核心是Spark Catalyst优化器的列引用歧义问题:
- 从
left过滤得到的right_fail和原left是同源派生关系,二者的同名列共享完全相同的表达式ID(ExprId,Spark内部用来唯一标识列的属性)。 - 你在join条件中写的
right_side.lower_street_number这类引用,Spark解析时会因为ExprId重复,错误地把右表的列引用解析为左表的列。 - 对应你的案例:左表第一条需要被过滤的行的
lower_street_number本身是null,被错误代入join条件后right_side.lower_street_number.isNotNull()就变成了false,导致这条行没有被join匹配上,最终left_anti join会保留它,出现结果不符合预期的问题。 - 而你重新构造的
right_success是独立的DataFrame,所有列的ExprId和left的列没有重叠,所以Spark可以正确区分左右表的列引用,join正常生效。
可行解决方案
不需要重新构造DataFrame,以下三个方案都可以解决问题:
方案1:给右表的列重新生成别名(保留原有列名)
仅需要给右表的所有列做一次别名重定义(仍使用原有列名),就能生成新的ExprId,消除歧义:
from pyspark.sql.functions import col right_fix = left.filter("lower_street_number IS NOT NULL")\ .select([col(c).alias(c) for c in left.columns]) result = join_files(left, right_fix) result.count() # 返回正确结果1
方案2:join时给左右表加明确别名,条件中用别名引用列
修改join逻辑,强制明确区分左右表的列,从根源避免解析歧义:
from pyspark.sql.functions import col def join_files(left_side, right_side): left_alias = left_side.alias("l") right_alias = right_side.alias("r") join_condition = [ ( (col("r.lower_street_number").isNotNull()) & (col("r.upper_street_number").isNotNull()) & (col("r.lower_street_number") <= col("l.street_number")) & (col("r.upper_street_number") >= col("l.street_number")) ) ] return left_alias.join(right_alias, join_condition, "left_anti")
方案3:给右表执行checkpoint(适合复杂派生场景)
如果右表是经过多层转换得到的,可以用checkpoint切断和原左表的谱系,自然消除ExprId重复:
spark.sparkContext.setCheckpointDir("./tmp_checkpoint") # 需先指定checkpoint存储路径 right_fix = left.filter("lower_street_number IS NOT NULL").checkpoint()
内容的提问来源于stack exchange,提问作者Matthew Hinea
相关产品推荐
相关产品推荐

