【问题标题】:Spark RDD any() and all() methods?Spark RDD any() 和 all() 方法?
【发布时间】:2014-12-20 10:32:30
【问题描述】:

我有一个RDD[T] 和一个谓词T => Boolean。 如何计算所有项目是否适合/不适合谓词?

当然可以这样:

rdd
 .map(predicate)
 .reduce(_ && _)

但这将需要完整的集合来迭代,这是一个矫枉过正。

我尝试了另一种适用于 local[1] 的方法,但似乎也遍历了真实集群上的所有内容:

rdd
 .map(predicate)
 .first()

[如果找不到任何需要,则失败并出现异常]

实现这一目标的规范方法是什么?

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    你可以使用aggregate:

    def forAll[T](rdd:RDD[T])(p:T => Boolean): Boolean = {
      rdd.aggregate(true)((b, t) => b && p(t), _ && _)
    }
    

    附带说明,在 spark 中没有真正的提前终止方法,您将作业发送到集群并执行。聚合只是做你想做的事的好方法。

    【讨论】:

    • 您的解决方案和我的选择 (1) 有什么区别?为什么聚合会在第一个“假”时停止?
    • @IlyaSmagin 我明白你的意思了。没有真正的方法可以提前终止 Spark,您将作业发送到集群并执行。聚合只是做你想做的事的好方法。
    • 您想单独发布一下,以便我将其标记为答案吗?
    猜你喜欢
    • 2015-10-07
    • 2015-12-07
    • 1970-01-01
    • 1970-01-01
    • 2019-06-16
    • 1970-01-01
    • 1970-01-01
    • 2015-02-10
    相关资源
    最近更新 更多