Spark DataFrame过滤:保留new_id存在于db_id中的行
解决Spark DataFrame过滤问题:保留存在于另一个DataFrame中的行
嘿,这个需求在Spark开发里挺常见的,我给你几个实用的实现方案,你可以根据数据集的大小和性能需求来选择:
先构建示例数据
首先咱们先把你提到的示例DataFrame创建出来,方便后续测试:
import org.apache.spark.sql.SparkSession val spark = SparkSession.builder().appName("FilterExample").master("local[*]").getOrCreate() import spark.implicits._ // 构建DF1,包含new_id列 val df1 = Seq(1, 2, 3, 4).toDF("new_id") // 构建DF2,包含db_id列 val df2 = Seq(1, 4, 5, 6, 10).toDF("db_id")
方法1:内连接(Inner Join)
这是最稳妥的方式,不管两个DataFrame的数据量大小都适用。内连接只会保留两边匹配的行,之后我们只需要保留DF1的new_id列即可:
val filteredDF1 = df1.join(df2, df1("new_id") === df2("db_id"), "inner") .select(df1("new_id")) // 只保留DF1的目标列 .distinct() // 可选:如果DF1/DF2存在重复值,用这个去重
方法2:isin + collect_set(适合小数据集)
如果DF2的数据量很小,可以先把db_id的值收集到Driver端的集合里,再用isin过滤DF1。注意:如果DF2数据量大,这种方法会导致Driver端内存溢出,谨慎使用!
// 从DF2中提取所有db_id值,转成Set val dbIdSet = df2.select("db_id").collect().map(_.getInt(0)).toSet // 过滤DF1,只保留new_id在dbIdSet中的行 val filteredDF1 = df1.filter($"new_id".isin(dbIdSet.toSeq:_*))
方法3:广播小表优化Join性能(适合DF2数据量小)
当DF2是小表时,用broadcast把它广播到所有Executor节点,能大幅减少shuffle操作,提升性能:
import org.apache.spark.sql.functions.broadcast val filteredDF1 = df1.join(broadcast(df2), df1("new_id") === df2("db_id"), "inner") .select(df1("new_id")) .distinct()
验证结果
运行以下代码查看过滤后的结果:
filteredDF1.show()
输出会是你想要的结果:
+------+ |new_id| +------+ | 1| | 4| +------+
内容的提问来源于stack exchange,提问作者Martee
相关产品推荐
相关产品推荐

