【发布时间】: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