【问题标题】:'Exception in thread "dispatcher-event-loop-0" java.lang.OutOfMemoryError: Java heap space ' error in Spark Scala codeSpark Scala代码中的“线程“dispatcher-event-loop-0”java.lang.OutOfMemoryError:Java堆空间中的异常”错误
【发布时间】:2020-02-29 14:33:25
【问题描述】:
val data = spark.read
    .text(filePath)
    .toDF("val")
    .withColumn("id", monotonically_increasing_id())



    val count = data.count()



    val header = data.where("id==1").collect().map(s => s.getString(0)).apply(0)



    val columns = header
    .replace("H|*|", "")
    .replace("|##|", "")
    .split("\\|\\*\\|")


    val structSchema = StructType(columns.map(s=>StructField(s, StringType, true)))



    var correctData = data.where('id > 1 && 'id < count-1).select("val")
    var dataString = correctData.collect().map(s => s.getString(0)).mkString("").replace("\\\n","").replace("\\\r","")
    var dataArr = dataString.split("\\|\\#\\#\\|").map(s =>{ 
                                                          var arr = s.split("\\|\\*\\|")
                                                          while(arr.length < columns.length) arr = arr :+ ""
                                                          RowFactory.create(arr:_*)
                                                         })
    val finalDF = spark.createDataFrame(sc.parallelize(dataArr),structSchema)

    display(finalDF)

这部分代码报错:

线程“dispatcher-event-loop-0”java.lang.OutOfMemoryError 中的异常:Java 堆空间

经过几个小时的调试主要是部分:

var dataArr = dataString.split("\\|\\#\\#\\|").map(s =>{ 
                                                          var arr = s.split("\\|\\*\\|")
                                                          while(arr.length < columns.length) arr = arr :+ ""
                                                          RowFactory.create(arr:_*)
                                                         })
    val finalDF = spark.createDataFrame(sc.parallelize(dataArr),structSchema)

导致错误。

我把部分改成

var dataArr = dataString.split("\\|\\#\\#\\|").map(s =>{
                                                          var arr = s.split("\\|\\*\\|")
                                                          while(arr.length < columns.length) arr = arr :+ ""
                                                          RowFactory.create(arr:_*)
                                                         }).toList
  val finalDF = sqlContext.createDataFrame(sc.makeRDD(dataArr),structSchema)

但错误仍然相同。我应该改变什么来避免这种情况?

当我运行此代码是 databricks spark 集群时,特定作业给出了此 Spark 驱动程序错误:

作业因阶段故障而中止:序列化任务 45:0 为 792585456 字节,超过了允许的最大值:spark.rpc.message.maxSize(268435456 字节)。

我添加了这部分代码:

spark.conf.set("spark.rpc.message.maxSize",Int.MaxValue)

但没用。

【问题讨论】:

标签: java scala apache-spark apache-spark-sql out-of-memory


【解决方案1】:

我的猜测是

var dataString = correctData.collect().map(s => s.getString(0)).mkString("").replace("\\\n","").replace("\\\r","")

是问题所在,因为您将(几乎)所有数据收集到驱动程序,即 1 个单一 JVM。

也许这条线会运行,但对dataString 的后续操作将超出您的内存限制。你不应该收集你的数据!相反,使用分布式“数据结构”,例如 Dataframe 或 RDD。

我认为你可以省略上面一行中的collect

【讨论】:

  • 如果我删除 collect 代码将失败。因为我是这个领域的新手,你能帮我解决这个问题吗?我应该在这里做什么
猜你喜欢
  • 2016-07-29
  • 2017-06-12
  • 2018-03-27
  • 2018-08-13
  • 2016-12-16
  • 2022-10-01
  • 2021-11-20
  • 2014-02-04
  • 1970-01-01
相关资源
最近更新 更多