【问题标题】:Apache Spark iterating through RDD gives error using mappartitionstopair使用 mappartitionstopair 遍历 RDD 的 Apache Spark 出现错误
【发布时间】:2015-07-07 20:47:23
【问题描述】:

我正在为 JavaPairRDD 使用 mapPartitionstoPair 函数,如下所示:

JavaPairRDD<MyKeyClass, MyValueClass> myRDD;

JavaPairRDD<Integer, Double> myResult = myRDD.mapPartitionsToPair(new PairFlatMapFunction<Iterator<Tuple2<MyKeyClass,MyValueClass>>, Integer, Double>(){

  public Iterable<Tuple2<MyInteger, MyDouble>> call(Iterator<Tuple2<MyKeyClass, MyValueClass>> arg0) throws Exception {

  Tuple2<MyKeyClass, MyValueClass> temp = arg0.next(); //The error is coming here...
  TreeMap<Integer, Double> dic = new TreeMap<Integer, Double>();

  do{

   ........
   // Some Code to compute to newIntegerValue and newDoubleValue from temp
   ........

   dic.put(newIntegerValue, newDoubleValue)
   temp = arg0.next();

   }while(arg0.hasNext());

  }

}

我可以在 Apache Spark 伪分布式模式下运行它。我无法在我的集群上运行上述代码。我收到以下错误:

java.util.NoSuchElementException: next on empty iterator
    at scala.collection.Iterator$$anon$2.next(Iterator.scala:39)
    at scala.collection.Iterator$$anon$2.next(Iterator.scala:37)
    at scala.collection.IndexedSeqLike$Elements.next(IndexedSeqLike.scala:64)
    at org.apache.spark.InterruptibleIterator.next(InterruptibleIterator.scala:43)
    at scala.collection.convert.Wrappers$IteratorWrapper.next(Wrappers.scala:30)
    at IncrementalGraph$6.call(MySparkJob.java:584)
    at IncrementalGraph$6.call(MySparkJob.java:573)
    at org.apache.spark.api.java.JavaRDDLike$$anonfun$fn$9$1.apply(JavaRDDLike.scala:186)
    at org.apache.spark.api.java.JavaRDDLike$$anonfun$fn$9$1.apply(JavaRDDLike.scala:186)
    at org.apache.spark.rdd.RDD$$anonfun$13.apply(RDD.scala:601)
    at org.apache.spark.rdd.RDD$$anonfun$13.apply(RDD.scala:601)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:35)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:263)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:230)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:61)
    at org.apache.spark.scheduler.Task.run(Task.scala:56)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:196)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1110)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:603)
    at java.lang.Thread.run(Thread.java:722)

我在 Hadoop 2.2.0 上使用 Spark 1.2.0。

谁能帮我解决这个问题??

更新: hasNext() 在迭代器上调用 next() 之前给出 true

【问题讨论】:

    标签: iterator apache-spark hadoop2 rdd nosuchelementexception


    【解决方案1】:

    我找到了答案。

    我将 myRDD 存储级别设为 MEMORY_ONLY。在 mapPartitonsToPair 转换开始之前,我的代码中有以下行:

    myRDD.persist(StorageLevel.MEMORY_ONLY());
    

    我删除了它并修复了程序。

    我不知道为什么它修复了它。如果有人能解释一下,不胜感激。

    【讨论】:

      【解决方案2】:

      您的代码假设将传入的所有迭代器都将包含一些元素,但事实并非如此。某些分区可能是空的(尤其是对于小型测试数据集)。这是一种非常常见的模式,只检查迭代器是否为空,如果在 mapPartitions 代码的开头是这种情况,则返回一个空的迭代器。希望有帮助:)

      【讨论】:

      • 你是对的。我是这么认为的。我将相应地更改我的代码并检查它是否修复它。实际上,我为 RDD 及其范围分区提供了固定数量的分区。我的测试数据分布均匀,所以所有 RDD 分区都应该有一些东西。
      • 我更改了代码以检查迭代器是否为空。之后它也给出了错误。
      猜你喜欢
      • 2014-10-26
      • 1970-01-01
      • 1970-01-01
      • 2014-05-13
      • 1970-01-01
      • 2015-07-17
      • 1970-01-01
      • 2015-07-05
      • 1970-01-01
      相关资源
      最近更新 更多