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

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. 判断组内是否存在元素的第二个值(元组首元素)不等于1的项
  2. 若满足条件,保留该组所有原始元素;否则丢弃整个分组

解决方案代码

// 筛选符合条件的分组并展开元素
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 22:31:56