【发布时间】: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