【问题标题】:Spark - Multiple filters on RDD in one passSpark - 一次通过 RDD 上的多个过滤器
【发布时间】:2015-07-07 10:19:35
【问题描述】:

我的 RDD 为 Map[String, String];有没有一种方法可以多次调用filter 而无需多次通过 RDD?

例如,我想做这样的事情:

val stateNY = mapRDD.filter(person => person("state").equals("NY"))
val stateOR = mapRDD.filter(person => person("state").equals("OR"))
val stateMA = mapRDD.filter(person => person("state").equals("MA"))
val stateWA = mapRDD.filter(person => person("state").equals("WA"))

还有这个:

val wage10to20 = mapRDD.filter(person => person("wage").toDouble > 10 && person("wage").toDouble <= 20)
val wage20to30 = mapRDD.filter(person => person("wage").toDouble > 20 && person("wage").toDouble <= 30)
val wage30to40 = mapRDD.filter(person => person("wage").toDouble > 30 && person("wage").toDouble <= 40)
val wage40to50 = mapRDD.filter(person => person("wage").toDouble > 40 && person("wage").toDouble <= 50)

其中mapRDD 的类型为RDD[Map[String, String]],一次通过。

【问题讨论】:

  • 使用分布式集合需要改变思维模型。可能您不需要进行此类过滤选择。考虑替代方案,将事物分组。

标签: scala apache-spark


【解决方案1】:

我假设你的意思是你想为每个值返回单独的 RDD(即不简单地做person =&gt; Set("NY", "OR", "MA", "WA").contains(person("state"))

通常使用Pair RDDs 可以实现您想要实现的目标

在您的第一个示例中,您可以使用:

val keyByState = mapRDD.keyBy(_("state"))

然后做groupByKey、reduceByKey等操作

或者在您的第二个示例中,以工资四舍五入到最接近的 10 为关键字。

【讨论】:

  • 谢谢!一个问题:当我执行keyBy 后跟groupByKey 时,我最终得到PairRDDStringCompactBuffer。然后如何将其转换为 Map[String, String] 的多个 RDD?
【解决方案2】:

如果您最终需要将它们放在单独的 RDD 中,则有时需要单独的过滤器和多次扫描。您应该缓存您正在遍历的 RDD(第一个示例中的 mapRDD),以防止它被多次读取。

在编写过滤器时执行过滤器与在另一个答案中建议的分组相比有一个优势,因为过滤器可以发生在地图一侧,而分组后过滤需要对所有数据进行混洗(包括与您不使用的状态相关的数据'不需要...)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-08-19
    • 1970-01-01
    • 2015-06-15
    • 2015-06-27
    • 2020-12-14
    • 2020-03-11
    • 2015-10-26
    • 2017-02-02
    相关资源
    最近更新 更多