【发布时间】:2016-08-22 03:02:52
【问题描述】:
我有一个非常简单的类,我正在尝试使用 spark 来减少它。 由于某种原因,它不断抛出异常无法序列化类。 这是我的课:
@SerialVersionUID(1000L)
class TimeRange(val start: Long, val end: Long) extends Serializable {
def this(){
this(0,0)
}
def mergeOverlapping(rangesSet : Set[TimeRange]) = {
def minMax(t1: TimeRange, t2: TimeRange) : TimeRange = {
new TimeRange(if(t1.start < t2.start) t1.start else t2.start, if(t1.end > t2.end) t1.end else t2.end)
}
(rangesSet ++ Set(this)).reduce(minMax)
}
def containsSlice(timeRange: TimeRange): Boolean ={
(start < timeRange.start && end > timeRange.start) ||
(start > timeRange.start && start < timeRange.end)
}
override def toString = {
"("+ start + ", " + end + ")"
}
}
我也尝试过 @SerialVersionUID(2L) 和 Kryo:
conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
conf.registerKryoClasses(Array(classOf[TimeRange], classOf[Set[TimeRange]]))
我正在使用带有 Spark 1.6.1 的 Scala 2.11。
编辑: 使用相同的映射器和缩减器,但不是使用类 TimeRange 使用 (Long, Long) 作品。
def mergeOverlapping(currRange : (Long, Long), rangesSet : Set[(Long, Long)]) = {
def minMax(t1: (Long, Long), t2: (Long, Long)) : (Long, Long) = {
(if(t1._1 < t2._1) t1._1 else t2._1, if(t1._2 > t2._2) t1._2 else t2._2)
}
(rangesSet ++ Set(currRange)).reduce(minMax)
}
def containsSlice(t1: (Long, Long), t2 : (Long, Long)): Boolean ={
(t1._1 < t2._1 && t1._2 > t2._1) ||
(t1._1 > t2._1 && t1._1 < t2._2)
}
【问题讨论】:
-
你能添加错误和堆栈跟踪吗? Sparks CloserCleaner 通常会告诉您哪个字段引起了问题。
标签: scala serialization apache-spark kryo