Spark RDD如何结合Filter与Row函数实现数据过滤?
解决RDD中结合Filter与Row转换的问题
嘿,我完全懂你的困惑!其实核心问题就在于filter的使用时机和针对Row的正确访问方式,咱们一步步来理清楚:
最优方案:先过滤,再转换为Row
这是效率最高的做法——先把不符合条件(ab.b > 0)的数据过滤掉,再进行Row转换,这样能减少后续map操作要处理的数据量。
假设你的原始RDD元素是包含a和b属性的对象(比如自定义的case class),代码应该这么写:
// 先过滤出b>0的元素,再映射为Row val rst = rdd.filter(ab => ab.b > 0).map(ab => Row(ab.a, ab.b))
如果喜欢用下划线简化语法,也可以写成:
val rst = rdd.filter(_.b > 0).map(ab => Row(ab.a, ab.b))
如果已经转换为Row,如何过滤?
要是你已经先做了map转换,那过滤的时候就得针对Row里的元素位置来判断了。因为Row是基于索引访问的(从0开始),你的ab.b对应Row的第1个位置,所以代码要这么写:
// 先转Row,再过滤(不推荐,因为多处理了不需要的数据) val rst = rdd.map(ab => Row(ab.a, ab.b)).filter(row => row.getInt(1) > 0)
如果b的类型不是Int,要换成对应的get方法,比如getDouble(1)、getLong(1),或者用泛型方法更安全:
filter(row => row.getAs[Double](1) > 0) // 假设b是Double类型
为什么之前过滤没生效?
大概率是这两个原因:
- 你把filter放在了map操作的内部(比如在map的函数里加过滤),这是错误的,map是用来转换每个元素,不是过滤;
- 要是在map之后过滤,但没正确访问Row的元素(比如用了错误的索引,或者类型不匹配),也会导致过滤逻辑不生效。
举个完整的实例,假设你的数据是case class:
case class User(id: String, score: Int) val userRDD: RDD[User] = sc.parallelize(List(User("u1", 80), User("u2", -5), User("u3", 90))) // 先过滤再转Row val validUserRowRDD = userRDD.filter(_.score > 0).map(user => Row(user.id, user.score)) // 输出结果会包含(u1,80)和(u3,90)的Row
内容的提问来源于stack exchange,提问作者Jill Clover
相关产品推荐
相关产品推荐

