【问题标题】:Scala IllegalArgumentException: can't serialize classScala IllegalArgumentException:无法序列化类
【发布时间】: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


【解决方案1】:

我认为这个类是可以序列化的,但是当你在 Spark 中使用它时,你必须确保闭包集中的所有变量都是可序列化的。我在学习Spark的时候遇到过这种问题,而且大部分是因为在闭包中涉及到了另一个不可序列化的变量(map或者reduce方法中的代码)。如果你能展示你如何使用这个类的代码,那可能会有所帮助。

【讨论】:

  • 如果是这种情况,我将无法在没有类包装器的情况下使用类对象。这意味着我可以使用 (Long, Long) 而不是 TimeRange 并且一切正常。我将添加匹配 (Long, Long) 变量的 containsSlice 和 mergeOverlapping
  • 我的意思是,你如何在火花相关代码中使用它。我测试了一个简单的例子,没有报错。 ` val sc = new SparkContext(new SparkConf().setMaster("local[2]").setAppName("test")) val tset = Set[TimeRange]() val ret = sc.parallelize(Range(0,100)) .map((,new TimeRange(1,2))).reduceByKey( (v1,v2) => v1.mergeOverlapping(tset).mergeOverlapping(tset)).collect() ret.foreach(println()) ` 另外对于reduce函数,一般是(A,A)=>A的格式,不明白spark中mergeOverlapping怎么用。
猜你喜欢
  • 1970-01-01
  • 2018-03-25
  • 2011-09-21
  • 2016-12-14
  • 1970-01-01
  • 2019-08-15
  • 2014-03-05
  • 1970-01-01
  • 2017-03-29
相关资源
最近更新 更多