You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 07:49:59