【问题标题】:Spark: How to map an RDD when access to another RDD is requiredSpark:当需要访问另一个 RDD 时如何映射一个 RDD
【发布时间】:2015-05-27 14:29:44
【问题描述】:

给定两个大型键值对 RDD(d1d2),它们都由唯一的 ID 键和 vector 值组成(例如 RDD[Int,DenseVector]) ,我需要映射d1,以便使用向量之间的欧几里德距离度量为它的每个元素获取d2中最近元素的ID

我还没有找到使用标准 RDD 转换的方法。我知道 Spark 中不允许嵌套 RDD,但是,如果可能的话,一个简单的解决方案是:

d1.map((k,v) => (k, d2.map{case (k2, v2) => val diff = (v - v2); (k2, sqrt(diff dot diff))} 
                      .takeOrdered(1)(Ordering.by[(Double,Double), Double](_._2))      
                      ._1))

此外,如果 d1 很小,我可以使用 Map(例如 d1.collectAsMap())并遍历其每个元素,但由于数据集的大小,这不是一个选项。

Spark 中的这种转换是否有任何替代方法?

编辑 1:

使用@holden 和@david-griffin 的建议,我使用cartesian()reduceByKey() 解决了这个问题。这是脚本(假设scSparkContext 并使用Breeze 库)。

val d1 = sc.parallelize(List((1,DenseVector(0.0,0.0)), (2,DenseVector(1.0,0.0)), (3,DenseVector(0.0,1.0))))
val d2 = sc.parallelize(List((1,DenseVector(0.0,0.75)), (2,DenseVector(0.0,0.25)), (3,DenseVector(1.0,1.0)), (4,DenseVector(0.75,0.0))))

val d1Xd2 = d1.cartesian(d2)
val pairDistances = d1Xd2.map{case ((k1, v1), (k2, v2)) => (k1, (k2, sqrt(sum(pow(v1-v2,2)))))}
val closestPoints = pairDistances.reduceByKey{case (x, y) => if (x._2 < y._2) x else y }

closestPoints.foreach(s => println(s._1 + " -> " + s._2._1))

得到的输出是:

1 -> 2
2 -> 4
3 -> 1

【问题讨论】:

  • 我会将它们转换为 DataFrame,并尝试在两者之间使用 join。否则,我认为您必须先做d1.cartesian(d2),然后使用reduce 找到每个d1._1 的最短距离
  • 几乎与此相同:stackoverflow.com/a/29953122/21755
  • 如果两个数据集都很大,则无法将每个元素与其他每个元素进行比较。您需要使用空间分区方案,然后join 分区并在分区内找到最佳匹配。
  • @Paul 您重定向到的答案也有效。它概括了当需要的邻居数量超过一个时(例如,对于 KNN)。

标签: scala nested apache-spark transformation rdd


【解决方案1】:

RDD 上的转换只能在驱动端应用,因此地图嵌套不起作用。正如@davidgriffin 指出的那样,您可以使用cartesian。对于您的用例,您可能希望使用 reduceByKey 跟进,并且在您的 reduce by key 中您可以跟踪最小距离。

【讨论】:

  • 我已根据您的建议使用解决方案更新了问题。谢谢。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-09-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-12-23
  • 2017-08-10
  • 1970-01-01
相关资源
最近更新 更多