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

Scala中如何基于列对匹配过滤不同Spark DataFrame?

Spark DataFrame 基于列对的过滤实现方法

你可以通过以下两种方式实现基于列对的过滤需求:

方法一:使用半连接(Semi Join)

这是适配大数据场景的最优方案,半连接仅保留左表(第一个DataFrame)中与右表(第二个DataFrame)满足列对匹配条件的行,且不会引入右表的额外列,完全契合你的需求。

示例代码(Python)

from pyspark.sql import SparkSession

# 初始化Spark会话
spark = SparkSession.builder.appName("ColumnPairFilter").getOrCreate()

# 构造第一个DataFrame
df1_data = [("aaa", "111"), ("bbb", "222"), ("ccc", "333")]
df1 = spark.createDataFrame(df1_data, ["value_1", "value_2"])

# 构造第二个DataFrame
df2_data = [("aaa", "222"), ("aaa", "111"), ("bbb", "111")]
df2 = spark.createDataFrame(df2_data, ["value_3", "value_4"])

# 基于列对匹配执行半连接
filtered_df = df1.join(df2, (df1.value_1 == df2.value_3) & (df1.value_2 == df2.value_4), "semi")

# 查看结果
filtered_df.show()

执行后输出结果:

+-------+-------+
|value_1|value_2|
+-------+-------+
|    aaa|    111|
+-------+-------+

方法二:基于列对集合过滤

如果第二个DataFrame的数据量较小,可以先将其中的列对收集为本地集合,再用isin进行过滤。注意这种方法不适合大数据量场景,因为collect会将数据拉取到Driver节点,可能引发内存溢出问题。

示例代码(Python)

# 收集第二个DataFrame中的列对为元组列表
pair_list = df2.select("value_3", "value_4").rdd.map(lambda x: (x[0], x[1])).collect()

# 过滤第一个DataFrame
filtered_df = df1.filter((df1.value_1, df1.value_2).isin(pair_list))

filtered_df.show()

输出结果与方法一一致。

内容的提问来源于stack exchange,提问作者elena

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 16:55:09