【问题标题】:mapPartitions returns empty arraymapPartitions 返回空数组
【发布时间】:2015-08-18 04:22:11
【问题描述】:

我有以下具有 4 个分区的 RDD:-

val rdd=sc.parallelize(1 to 20,4)

现在我尝试在此调用 mapPartitions:-

scala> rdd.mapPartitions(x=> { println(x.size); x }).collect
5
5
5
5
res98: Array[Int] = Array()

为什么返回空数组? anonymoys 函数只是简单地返回它接收到的相同迭代器,那么它如何返回空数组呢?有趣的是,如果我删除 println 语句,它确实返回非空数组:-

scala> rdd.mapPartitions(x=> { x }).collect
res101: Array[Int] = Array(1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11, 12, 13, 14, 15, 16, 17, 18, 19, 20)

这个我不明白。 println(它只是打印迭代器的大小)的存在如何影响函数的最终结果?

【问题讨论】:

    标签: apache-spark rdd


    【解决方案1】:

    那是因为x 是一个TraversableOnce,这意味着你通过调用size 来遍历它,然后返回它......空。

    您可以通过多种方式解决此问题,但这里有一种:

    rdd.mapPartitions(x=> {
      val list = x.toList;
      println(list.size);
      list.toIterator
    }).collect
    

    【讨论】:

      【解决方案2】:

      要了解发生了什么,我们必须查看您传递给 mapPartitions 的函数的签名:

      (Iterator[T]) ⇒ Iterator[U]
      

      那么Iterator 是什么?如果你看一下Iterator documentation,你会发现它是一个扩展TraversableOnce的特征:

      trait Iterator[+A] extends TraversableOnce[A]
      

      上面应该给你一个提示,你的情况会发生什么。迭代器提供了两种方法hasNextnext。要获得迭代器的size,您必须简单地对其进行迭代。之后hasNext 返回false,结果是一个空的Iterator

      【讨论】:

        猜你喜欢
        • 2016-02-06
        • 2017-07-12
        • 1970-01-01
        • 1970-01-01
        • 2021-08-11
        • 2017-02-27
        • 2020-04-01
        • 2019-04-17
        • 2017-04-24
        相关资源
        最近更新 更多