【问题标题】:Unexpected behavior inside the foreachPartition method of a RDDRDD 的 foreachPartition 方法中的意外行为
【发布时间】:2016-08-21 10:14:10
【问题描述】:

我通过 spark-shell 评估了以下几行 scala 代码:

val a = sc.parallelize(Array(1,2,3,4,5,6,7,8,9,10))
val b = a.coalesce(1)
b.foreachPartition { p => 
  p.map(_ + 1).foreach(println)
  p.map(_ * 2).foreach(println)
}

输出如下:

2
3
4
5
6
7
8
9
10
11

为什么分区p在第一个map之后就变空了?

【问题讨论】:

    标签: scala apache-spark rdd


    【解决方案1】:

    它对我来说并不奇怪,因为 pIterator,当你用 map 遍历它时,它没有更多的值,并且考虑到 lengthsize 的快捷方式,其实现方式如下:

    def size: Int = {
      var result = 0
      for (x <- self) result += 1
      result
    }
    

    你得到 0。

    【讨论】:

      【解决方案2】:

      答案在 scala 文档http://www.scala-lang.org/api/2.11.8/#scala.collection.Iterator 中。它明确指出迭代器(p 是一个迭代器)在调用 map 方法后必须丢弃。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2016-11-08
        • 2023-02-26
        • 1970-01-01
        • 2017-05-24
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多