【问题标题】:Scala/Spark serializable error - join don't workScala/Spark 可序列化错误 - 加入不起作用
【发布时间】:2017-11-19 07:52:52
【问题描述】:

我正在尝试在 2 个 RDD 之间使用 join 方法并将其保存到 cassandra,但我的代码不起作用。一开始,我得到了一个巨大的 Main 方法,一切都运行良好,但是当我使用函数和类时,这不起作用。我是 scala 和 spark 的新手

代码是:

  class Migration extends Serializable { 

  case class userId(offerFamily: String, bp: String, pdl: String) extends Serializable
  case class siteExternalId(site_external_id: Option[String]) extends Serializable
  case class profileData(begin_ts: Option[Long], Source: Option[String]) extends Serializable

  def SparkMigrationProfile(sc: SparkContext) = {

    val test  = sc.cassandraTable[siteExternalId](KEYSPACE,TABLE)
     .keyBy[userId]
     .filter(x => x._2.site_external_id != None)

    val profileRDD = sc.cassandraTable[profileData](KEYSPACE,TABLE)
    .keyBy[userId]

    //dont work
   test.join(profileRDD)
    .foreach(println)

    // don't work
     test.join(profileRDD)
     .saveToCassandra(keyspace, table)

   }

在开始时,我得到了著名的:线程“主”org.apache.spark.SparkException 中的异常:Task not serializable at 。 . . 所以我扩展了我的主类和案例类,但仍然不起作用。

【问题讨论】:

    标签: scala apache-spark serialization rdd


    【解决方案1】:

    我认为您应该将案例类从 Migration 类移动到专用文件和/或对象。这应该可以解决您的问题。此外,Scala 案例类默认是可序列化的。

    【讨论】:

    • 它的工作!我现在太傻了。 . .你能解释一下为什么吗?
    • 嗨@user3394825,这很难说,因为我没有将Spark与Cassandra一起使用。根据我的经验,我在使用其他类中定义的案例类时遇到了类似的问题。在您的情况下,为 cassandraTable 函数 (github.com/datastax/spark-cassandra-connector/blob/master/…) 创建隐式参数可能存在一些问题,例如rrf: RowReaderFactory[T], ev: ValidRDDType[T],但我只是猜测。我知道在使用 Spark SQL Encoder 时也会出现类似的异常。
    • 案例类在技术上是内部类,可以访问封闭的迁移实例。当它们被序列化时,伴随的 Migration 对象也被序列化。即使它被标记为 Serializable,它内部的某处也可能有一些实例变量不是。通常罪魁祸首是 SparkContext 对象。
    猜你喜欢
    • 2017-09-21
    • 2020-02-04
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-16
    • 1970-01-01
    • 1970-01-01
    • 2012-06-18
    相关资源
    最近更新 更多