【发布时间】:2019-01-27 16:08:53
【问题描述】:
我使用 spark 读取了一个文本文件并将其保存在 JavaRDD 中,并尝试打印保存在 RDD 中的数据。我在一个主服务器和两个从服务器的集群中运行我的代码。但是我遇到了异常,例如, 容器超过阈值,同时遍历 RDD。代码在独立模式下完美运行。
我的代码:
SparkContext sc = new SparkContext("spark://master.com:7077","Spark-Phoenix");
JavaSparkContext jsc = new JavaSparkContext(sc);
JavaRDD<String> trs_testing = jsc.textFile("hdfs://master.com:9000/Table/sample");
//using iterator
Iterator<String> iStr= trs_testing.toLocalIterator();
while(iStr.hasNext()){ //here I am getting exception
System.out.println("itr next : " + iStr.next());
}
//using foreach()
trs_testing.foreach(new VoidFunction<String>() {//here I am getting exception
private static final long serialVersionUID = 1L;
@Override public void call(String line) throws Exception {
System.out.println(line);
}
});
//using collect()
for(String line:trs_testing.collect()){//here I am getting exception
System.out.println(line);
}
//using foreachPartition()
trs_testing.foreachPartition(new VoidFunction<Iterator<String>>() {//here I am getting exception
private static final long serialVersionUID = 1L;
@Override public void call(Iterator<String> arg0) throws Exception {
while (arg0.hasNext()) {
String line = arg0.next();
System.out.println(line);
}
}
});
例外:
错误 TaskSchedulerImpl 在 master.com 上丢失了执行程序 0:远程 RPC 客户解除关联。可能是由于容器超过阈值, 或网络问题。检查驱动程序日志以获取 WARN 消息。错误 slave1.com 上的 TaskSchedulerImpl 丢失执行程序 1:远程 RPC 客户端 解离。可能是由于容器超过阈值,或 网络问题。检查驱动程序日志以获取 WARN 消息。错误 master.com 上的 TaskSchedulerImpl Lost executor 2:远程 RPC 客户端 解离。可能是由于容器超过阈值,或 网络问题。检查驱动程序日志以获取 WARN 消息。错误 slave2.com 上的 TaskSchedulerImpl Lost executor 3:远程 RPC 客户端 解离。可能是由于容器超过阈值,或 网络问题。检查驱动程序日志以获取 WARN 消息。错误 TaskSetManager Task 0 在 stage 0.0 失败 4 次;中止工作 线程“主”org.apache.spark.SparkException 中的异常:作业 由于阶段失败而中止:阶段 0.0 中的任务 0 失败了 4 次,大多数 最近失败:在 0.0 阶段丢失任务 0.3(TID 3,slave1.com): ExecutorLostFailure (executor 3 exited 由其中一个运行引起 tasks) 原因:远程 RPC 客户端已解除关联。可能由于 容器超过阈值或网络问题。检查驱动程序日志 警告消息。驱动程序堆栈跟踪:在 org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1454) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1442) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1441) 在 scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59) 在 scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48) 在 org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1441) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:811) 在 org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:811) 在 scala.Option.foreach(Option.scala:257) 在 org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:811) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:1667) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1622) 在 org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:1611) 在 org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48) 在 org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:632) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1890) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1903) 在 org.apache.spark.SparkContext.runJob(SparkContext.scala:1916) 在 org.apache.spark.rdd.RDD$$anonfun$take$1.apply(RDD.scala:1324) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112) 在 org.apache.spark.rdd.RDD.withScope(RDD.scala:358) 在 org.apache.spark.rdd.RDD.take(RDD.scala:1298) 在 org.apache.spark.rdd.RDD$$anonfun$first$1.apply(RDD.scala:1338) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112) 在 org.apache.spark.rdd.RDD.withScope(RDD.scala:358) 在 org.apache.spark.rdd.RDD.first(RDD.scala:1337) 在 com.test.java.InsertE.main(InsertE.java:147)
【问题讨论】:
标签: apache-spark cluster-computing rdd