如何使用Spark基于另一张表的同名字段过滤当前表的记录
问题解答
初始写法报错原因
Spark 的 where/filter 算子仅支持引用当前算子所属 DataFrame 的列、常量值或广播变量,不能直接引用其他未关联 DataFrame 的列,因此你最初的写法不符合语法规则会报错。
现有实现的有效性
你当前使用的内连接写法逻辑是正确的,只是默认会返回两张表的所有字段,所以会出现重复的关联字段,除了手动调用 drop() 删除冗余字段外,还有两种更优的替代方案:
更优方案
方案1:左半连接(通用最优,推荐)
左半连接(LEFT SEMI JOIN)的设计目标就是仅返回左表中能和右表匹配到的记录,不会返回右表的任何字段,完全匹配你的过滤需求,同时性能比普通内连接更好:匹配到第一条符合条件的右表记录就会终止当前行的匹配,不需要处理右表后续数据,也不会产生冗余字段。
Scala DSL 写法(关联字段同名时更简洁)
val joinColumns = Seq("servizio_rap", "filiale_rap", "codice_rap") val filtredTablehb004_p = fixTablehb004_p.join( filtredTable109_p, joinColumns, "left_semi" )
SQL 写法
val filtredTablehb004_p = sparkSession.sql( """SELECT a.* FROM fixTablehb004_p a LEFT SEMI JOIN filtredTable109_p b ON a.servizio_rap = b.servizio_rap AND a.filiale_rap = b.filiale_rap AND a.codice_rap = b.codice_rap """.stripMargin)
方案2:广播变量过滤(仅适用于右表数据量小的场景)
如果 filtredTable109_p 的数据量远小于 Spark 默认的广播阈值(默认1000MB),可以把右表的关联键广播到所有执行节点,避免Shuffle操作,性能更高:
import org.apache.spark.sql.functions.struct // 收集右表的关联键组合并广播 val rightKeys = filtredTable109_p .select(struct("servizio_rap", "filiale_rap", "codice_rap").as("key")) .collect() .map(_.getAs[Row]("key")) val broadcastKeys = sparkSession.sparkContext.broadcast(rightKeys) // 直接过滤左表 val filtredTablehb004_p = fixTablehb004_p.filter( struct("servizio_rap", "filiale_rap", "codice_rap").isin(broadcastKeys.value:_*) )
内容的提问来源于stack exchange,提问作者Miko
相关产品推荐
相关产品推荐

