【问题标题】:code execution inside Spark foreachSpark foreach 中的代码执行
【发布时间】:2014-12-22 00:14:31
【问题描述】:

我有两个 RDD:pointspointsWithinEpspoints 中的每个点代表x, y 坐标。 pointsWithinEps 代表两点和它们之间的距离:((x, y), distance)。我想循环所有点,并且对于每个点,只过滤pointsWithinEps 中的元素作为x(第一个)坐标。所以我做了以下事情:

    points.foreach(p =>
      val distances = pointsWithinEps.filter{
        case((x, y), distance) => x == p
      }
      if (distances.count() > 3) {
//        do some other actions
      }
    )

但是这种语法是无效的。据我了解,不允许在 Spark foreach 中创建变量。我应该这样做吗?

for (i <- 0 to points.count().toInt) {
  val p = points.take(i + 1).drop(i) // take the point
  val distances = pointsWithinEps.filter{
    case((x, y), distance) => x == p
  }
  if (distances.count() > 3) {
    //        do some other actions
  }
}

或者有更好的方法来做到这一点?完整代码托管在这里:https://github.com/timasjov/spark-learning/blob/master/src/DBSCAN.scala

编辑:

points.foreach({ p =>
  val pointNeighbours = pointsWithinEps.filter {
    case ((x, y), distance) => x == p
  }
  println(pointNeighbours)
})

现在我有以下代码,但它会引发 NullPointerException (pointsWithinEps)。为什么pointsWithinEps为空(在foreach之前有元素),如何解决?

【问题讨论】:

  • 我是否理解正确,对于points 上的每个点 (x,y),您想要来自 pointsWithinEps 的所有 ((x,y),distance) 元组源自同一 (x ) ?
  • 是的。基本上对于每个点,我都想找出哪些其他点是它的邻居(在 epsilon 内的点)。在我的情况下,它是点本身和 ((x, y), distance) 结构中的 x。代码在 github 中,因此例如您可以执行它并在调试器中找到确切的值。

标签: scala apache-spark rdd


【解决方案1】:

为了收集从给定坐标开始的所有距离点,一种简单的分布式方法是通过该坐标x 对点进行键控,并按该键对它们进行分组,如下所示:

val pointsWithinEpsByX = pointsWithinEps.map{case ((x,y),distance) => (x,((x,y),distance))}
val xCoordinatesWithDistance = pointsWithinEpsByX.groupByKey

然后将点的RDD与上一次转换的结果左连接:

val pointsWithCoordinatesWithDistance = points.leftOuterJoin(xCoordinatesWithDistance)

【讨论】:

  • 编译器显示没有 groupByKey 方法,只有 groupBy。 leftOuterJoin 方法也是如此。我正在使用为 Hadoop 1.X 预构建的 spark 1.1.0
  • 另外,我应该把它放到foreach循环中吗?
  • 这些是通过隐式转换在 (key, value) 对的 RDD 上可用的函数。在程序顶部导入org.apache.spark.SparkContext._ 以使用这些功能。 - 另外,不需要循环。这个功能管道通过对整个数据集应用转换和分组来完成工作。
  • 哦,现在我有 groupByKey 方法,但我仍然没有 leftOuterJoin 方法。我不明白 leftOuterJoin 方法在哪里。哦,我明白了,从文件读取后点不是 RDD 类型,我无法将 leftOuterJoin 应用于它们。如何将它们转换为 RDD?
  • 我假设点的形式是(x,y) ??隐式转换应用于Tuple2 类型的RDD 或(x,y) 之类的“对”
【解决方案2】:

声明变量意味着你有一个块,而不仅仅是一个表达式,所以你需要使用大括号{},例如

point.foreach({p => ... })

【讨论】:

  • 谢谢!但是您也可以帮助我描述的功能吗?
猜你喜欢
  • 2017-01-26
  • 1970-01-01
  • 2020-05-22
  • 1970-01-01
  • 2019-08-11
  • 1970-01-01
  • 2019-11-29
  • 1970-01-01
相关资源
最近更新 更多