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
相关产品推荐
相关产品推荐

