Spark RDD分组后过滤需求:groupByKey后按指定条件筛选
Spark RDD groupByKey 后按组条件筛选实现方案
场景1:筛选组内存在第二个元素不等于1的键对应元素
原始输入RDD代码
// 创建RDD val rdd = sc.makeRDD(List(("a", (1, "m")), ("b", (1, "m")), ("a", (1, "n")), ("b", (2, "n")), ("c", (1, "m")), ("c", (5, "m")), ("d", (1, "m")), ("d", (1, "n")))) val groupRDD = rdd.groupByKey()
实现思路
基于groupByKey分组结果,对每个分组完成两步操作:
- 判断组内是否存在元素的第二个值(元组首元素)不等于1的项
- 若满足条件,保留该组所有原始元素;否则丢弃整个分组
解决方案代码
// 筛选符合条件的分组并展开元素 val resultRDD = groupRDD.flatMap { case (key, iter) => val hasNonOne = iter.exists(_._1 != 1) if (hasNonOne) { iter.map((key, _)) } else { Iterator.empty } } // 输出结果验证 resultRDD.collect().foreach(println)
执行结果
("b", (1, "m")), ("b", (2, "n")), ("c", (1, "m")), ("c", (5, "m"))
场景2:筛选组内并非所有第二个元素均为"x"的键对应元素
输入示例RDD
val rdd = sc.makeRDD(List(("a",("x","m")), ("a",("x","n")), ("b",("x","m")), ("b",("y","n")), ("c",("x","m")), ("c",("z","m")), ("d",("x","m")), ("d",("x","n")))) val groupRDD = rdd.groupByKey()
实现思路
基于groupByKey分组结果,判断每个分组是否存在元素的第二个值(元组首元素)不等于"x"(等价于并非所有元素的该值都是"x"),满足条件则保留该组所有原始元素。
解决方案代码
// 筛选符合条件的分组并展开元素 val resultRDD = groupRDD.flatMap { case (key, iter) => val hasNonX = iter.exists(_._1 != "x") if (hasNonX) { iter.map((key, _)) } else { Iterator.empty } } // 输出结果验证 resultRDD.collect().foreach(println)
执行结果
("b",("x","m")), ("b",("y","n")), ("c",("x","m")), ("c",("z","m"))
内容的提问来源于stack exchange,提问作者1580923067
相关产品推荐
相关产品推荐

