Spark获取基于单列的DataFrame非交集行的高效实现问询
高效实现Spark DataFrame非交集行提取(基于单列匹配)
需求场景
给定两个DataFrame:
- df1(参考表):
| _id | name |
|---|---|
| 12 | abc |
| 56 | def |
| 90 | jkl |
- df2(待处理表):
| _id | user_name | event_name |
|---|---|---|
| 12 | abc | some |
| 34 | xyz | left |
| 56 | def | right |
| 78 | ghi | middle |
| 90 | jkl | front |
需要从df2中提取**_id不在df1中的行**,最终结果如下:
| _id | user_name | event_name |
|---|---|---|
| 34 | xyz | left |
| 78 | ghi | middle |
已尝试方案的问题
- 基于collect+isin过滤:
val idsList = df1.select("_id").dropDuplicates().map(_.getInt(0)).collect.toList val dffDf = df2.filter(not($"_id".isin(idsList: _*)))
问题:collect会将df1的_id数据拉取到Driver端,数据量大时会导致Driver内存压力陡增、执行耗时过长。
- Left Anti Join:
df2.join(df1, df2("_id") === df1("_id"),"leftanti")
问题:抛出解析迭代次数超限错误:
org.apache.spark.sql.catalyst.errors.package$TreeNodeException: Max iterations (100) reached for batch Resolution, please set 'spark.sql.analyzer.maxIterations' to a larger value
可行解决方案
方案1:优化Left Anti Join写法,简化解析逻辑
先对df1做预处理,只保留需要的_id列并去重,然后用Seq("_id")指定连接列(避免列名冲突导致的解析逻辑复杂化):
// 预处理df1,仅保留去重后的_id列 val df1Filter = df1.select("_id").distinct() // 使用Seq指定连接列,让Spark解析逻辑更清晰 val resultDf = df2.join(df1Filter, Seq("_id"), "leftanti")
这种写法减少了不必要的列参与join,降低了Spark解析器的负担,大概率能避开迭代次数超限的问题。
方案2:用Left Outer Join替代Left Anti Join
如果方案1仍报错,可以用左外连接后过滤空值的方式实现相同逻辑,不同的解析路径可能绕过原问题:
val df1Filter = df1.select("_id").distinct() val resultDf = df2.join(df1Filter, Seq("_id"), "leftouter") // 过滤df1中无匹配的行 .filter(df1Filter("_id").isNull) // 移除df1带来的冗余列 .drop(df1Filter("_id"))
方案3:结合广播优化(适用于df1数据量较小的场景)
如果df1的数据量不大,使用broadcast广播小表,既减少shuffle提升性能,也可能改变执行计划的解析流程:
import org.apache.spark.sql.functions.broadcast val df1Filter = df1.select("_id").distinct() val resultDf = df2.join(broadcast(df1Filter), Seq("_id"), "leftanti")
方案4:临时调整解析器迭代次数(终极workaround)
如果以上方案都无效,可以临时调大Spark解析器的迭代次数参数,绕过报错:
// 在执行join前设置参数 spark.conf.set("spark.sql.analyzer.maxIterations", 200) val df1Filter = df1.select("_id").distinct() val resultDf = df2.join(df1Filter, Seq("_id"), "leftanti")
注意:这只是临时 workaround,若数据量持续增大可能仍会出现问题,优先推荐前三种方案。
内容的提问来源于stack exchange,提问作者Ashit_Kumar
相关产品推荐
相关产品推荐

