【问题标题】:Operate on neighbor elements in RDD in Spark在 Spark 中对 RDD 中的相邻元素进行操作
【发布时间】:2018-05-23 03:28:36
【问题描述】:

因为我有收藏:

List(1, 3,-1, 0, 2, -4, 6)

很容易将其排序为:

List(-4, -1, 0, 1, 2, 3, 6)

然后我可以通过计算 6 - 3、3 - 2、2 - 1、1 - 0 等来构造一个新集合,如下所示:

for(i <- 0 to list.length -2) yield {
    list(i + 1) - list(i)
}

并得到一个向量:

Vector(3, 1, 1, 1, 1, 3)

也就是说,我想让下一个元素减去当前元素。

但是如何在 RDD on Spark 中实现呢?

我知道收藏:

List(-4, -1, 0, 1, 2, 3, 6)

集合会有一些分区,每个分区是有序的,我可以对每个分区做类似的操作,并在每个分区上一起收集结果吗?

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    最有效的解决方案是使用sliding方法:

    import org.apache.spark.mllib.rdd.RDDFunctions._
    
    val rdd = sc.parallelize(Seq(1, 3,-1, 0, 2, -4, 6))
      .sortBy(identity)
      .sliding(2)
      .map{case Array(x, y) => y - x}
    

    【讨论】:

      【解决方案2】:

      假设你有类似的东西

      val seq = sc.parallelize(List(1, 3, -1, 0, 2, -4, 6)).sortBy(identity)
      

      让我们像 Ton Torres 建议的那样,以 index 作为 key 创建第一个集合

      val original = seq.zipWithIndex.map(_.swap)
      

      现在我们可以构建移动一个元素的集合。

      val shifted = original.map { case (idx, v) => (idx - 1, v) }.filter(_._1 >= 0)
      

      接下来我们可以按索引降序计算所需的差异

      val diffs = original.join(shifted)
            .sortBy(_._1, ascending = false)
            .map { case (idx, (v1, v2)) => v2 - v1 }
      

      所以

       println(diffs.collect.toSeq)
      

      显示

      WrappedArray(3, 1, 1, 1, 1, 3)
      

      请注意,如果反转不是关键,您可以跳过sortBy 步骤。

      另请注意,对于本地收集,这可以更简单地计算,例如:

      val elems = List(1, 3, -1, 0, 2, -4, 6).sorted  
      
      (elems.tail, elems).zipped.map(_ - _).reverse
      

      但在RDD 的情况下,zip 方法要求每个集合的每个分区都应包含相等的元素计数。因此,如果您要实现tail 之类的

      val tail = seq.zipWithIndex().filter(_._2 > 0).map(_._1)  
      

      tail.zip(seq) 不起作用,因为这两个集合都需要每个分区的元素数量相等,并且每个分区都有一个元素应该移动到前一个分区。

      【讨论】:

      • 它会起作用,但效率极低。您可以使用rangePartitioner 降低成本,但我怀疑这是否值得大惊小怪。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2021-07-13
      • 2021-06-04
      • 2017-10-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多