【问题标题】:Exception while reading text file in cluster mode在集群模式下读取文本文件时出现异常
【发布时间】: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


    【解决方案1】:

    当您在本地系统/独立模式下运行 Spark 作业时,所有数据将在同一台机器上,因此您将能够迭代和打印数据。

    当 Spark Job 在集群模式/环境中运行时,数据将被分割成碎片并分发到集群中的所有机器(RDD - Resilient Distributed Datasets)。因此,要以这种方式打印,您必须使用 foreach() 函数。

    试试这个:

    trs_testing.foreach(new VoidFunction<String>(){ 
              public void call(String line) {
                  System.out.println(line);
              }
    });
    

    【讨论】:

    • 我也尝试过 foreach()、foreachpartition() 和 collect()。但仍然面临同样的问题。当我尝试处理数据时,它正在抛出异常。你能建议任何其他的吗替代读取数据?
    • 代码似乎没问题。您还可以提供用于运行 spark 作业的命令吗?工作节点和主节点的配置。您也尝试从 hdfs 加载的数据大小?
    【解决方案2】:

    我得到了解决方案。我通过我的机器执行代码,而我的主从服务器在远程服务器上运行。我将代码导出到远程服务器并最终能够处理数据。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-04-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多