【发布时间】: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)
但没用。
【问题讨论】:
-
请阅读Under what circumstances may I add “urgent” or other similar phrases to my question, in order to obtain faster answers? - 总结是这不是解决志愿者的理想方式,并且可能会适得其反。请不要将此添加到您的问题中。
-
对不起。我读了,不会再用了。
标签: java scala apache-spark apache-spark-sql out-of-memory