【问题标题】:Spark zipPartitions on the same RDDSpark zipPartitions 在同一个 RDD 上
【发布时间】:2015-03-13 12:16:51
【问题描述】:

我是 Spark 的新手,我在执行 cartesian 之类的操作时遇到了一些问题,但仅限于同一分区内。也许一个例子可以清楚地说明我想要做什么:假设我们有一个用sc.parallelize(1,2,3,4,5,6) 制作的RDD,这个RDD 被分成三个分区,分别包含:(1,2); (3,4) ; (5,6)。比我想获得以下结果:((1,1),(1,2),(2,1),(2,2)); ((3,3),(3,4),(4,3),(4,4)) ; ((5,5),(5,6),(6,5),(6,6)).

到目前为止我所做的是:

 partitionedData.zipPartitions(partitionedData)((aiter, biter) => {
  var res = new ListBuffer[(Double,Double)]()
  while(aiter.hasNext){
    val a = aiter.next()
    while(biter.hasNext){
      val b = biter.next()
      res+=(a,b)
    }
  }
  res.iterator
})

但它不起作用,因为 aiter 和 biter 是同一个迭代器...所以我只得到结果的第一行。

有人可以帮我吗?

谢谢。

【问题讨论】:

    标签: apache-spark rdd


    【解决方案1】:

    使用RDD.mapPartitions:

    val rdd = sc.parallelize(1 to 6, 3)
    val res = rdd.mapPartitions { iter =>
      val seq = iter.toSeq
      val res = for (a <- seq; b <- seq) yield (a, b)
      res.iterator
    }
    res.collect
    

    打印:

    res0: Array[(Int, Int)] = Array((1,1), (1,2), (2,1), (2,2), (3,3), (3,4), (4,3), (4,4), (5,5), (5,6), (6,5), (6,6))
    

    【讨论】:

    • 感谢您的回复,我是 Scala 的新手。什么是序列?我的意思是它在内存使用方面的表现如何?谢谢。
    • Seq 是一个“特征”(一个接口),所以你永远无法确定:)。看起来Iterator.toSeq创建了一个Stream,就像一个惰性链表。 Scala 喜欢链表。如果您担心内存使用情况,可以尝试 toArray 而不是 toSeq。
    猜你喜欢
    • 2014-07-01
    • 2015-12-05
    • 1970-01-01
    • 2016-12-23
    • 2015-10-26
    • 2017-08-19
    • 2015-07-05
    • 1970-01-01
    • 2015-07-08
    相关资源
    最近更新 更多