Spark 1.6如何基于特定列执行Except操作过滤数据?
在Spark 1.6中移除匹配特定键的行:替代EXCEPT的方案
在Spark 1.6里,确实没办法直接针对特定列执行EXCEPT操作——因为EXCEPT严格要求两个DataFrame的列数、列名和数据类型完全一致,它是对整个行的去重排除,而不是针对特定键组合。不过你完全可以用更合适的方案来实现你的需求,这里有两种常用的方法,优先推荐第一种:
方法1:使用Left Anti Join(最推荐)
Spark 1.6正式引入了left_anti连接类型,这正是为你的场景量身打造的:它会保留左表(dfA)中所有在右表(dfB)的指定键组合里不存在的行,而且不会引入右表的任何额外列,性能也非常高效。
代码示例
Scala版本
// 先从dfB中提取需要匹配的键列,再执行左反连接 val filteredDF = dfA.join( dfB.select("year", "cid"), // 只保留用于匹配的列 Seq("year", "cid"), // 指定连接的键 "left_anti" // 连接类型:左反连接 )
Python版本
filtered_df = dfA.join( dfB.select("year", "cid"), ["year", "cid"], "left_anti" )
这个方法的优势在于:
- 无需修改原DataFrame的结构,直接保留dfA的所有列
- 是Spark优化器原生支持的操作,执行效率远高于其他手动过滤的方法
- 适合大数据量场景,不会把数据拉到Driver端导致内存问题
方法2:收集键组合后用isin过滤(仅适合小数据集)
如果dfB的数据量很小(比如只有几千行),你也可以先把dfB中的(year, cid)组合收集到Driver端,再用isin来过滤dfA。但注意这个方法不适合大数据量,因为collect()会把dfB的所有键数据加载到Driver内存中,可能导致OOM。
代码示例
Scala版本
// 提取并去重dfB的键组合,收集到Driver端 val keyPairs = dfB.select("year", "cid").distinct().collect() // 过滤dfA中不在键组合里的行 val filteredDF = dfA.filter( !((col("year"), col("cid")).isin(keyPairs: _*)) )
Python版本
# 提取键组合并转换为元组列表 key_pairs = [tuple(row) for row in dfB.select("year", "cid").distinct().collect()] # 过滤dfA filtered_df = dfA.filter( ~((dfA["year"], dfA["cid"]).isin(key_pairs)) )
总结
优先选择Left Anti Join,它是处理这类“按键排除行”场景的标准高效方案;只有当dfB数据量极小的时候,才考虑用isin的方法。
内容的提问来源于stack exchange,提问作者RefiPeretz
相关产品推荐
相关产品推荐

