【问题标题】:Unable to serialize SparkContext in foreachRDD无法在 foreachRDD 中序列化 SparkContext
【发布时间】:2016-12-12 22:07:24
【问题描述】:

我正在尝试将流数据从 Kafka 保存到 cassandra。我能够读取和解析数据,但是当我调用下面的行来保存数据时,我得到了 Task not Serializable 异常。我的课程正在扩展可序列化但不知道为什么我会看到这个错误,在谷歌搜索 3 小时后没有得到太多帮助,有人可以给出任何指示吗?

val collection = sc.parallelize(Seq((obj.id, obj.data)))
collection.saveToCassandra("testKS", "testTable ", SomeColumns("id", "data"))` 


import org.apache.spark.SparkConf
import org.apache.spark.SparkContext
import org.apache.spark.sql.SaveMode
import org.apache.spark.streaming.Seconds
import org.apache.spark.streaming.StreamingContext
import org.apache.spark.streaming.kafka.KafkaUtils
import com.datastax.spark.connector._

import kafka.serializer.StringDecoder
import org.apache.spark.rdd.RDD
import com.datastax.spark.connector.SomeColumns
import java.util.Formatter.DateTime

object StreamProcessor extends Serializable {
  def main(args: Array[String]): Unit = {
    val sparkConf = new SparkConf().setMaster("local[2]").setAppName("StreamProcessor")
      .set("spark.cassandra.connection.host", "127.0.0.1")

    val sc = new SparkContext(sparkConf)

    val ssc = new StreamingContext(sc, Seconds(2))

    val sqlContext = new SQLContext(sc)

    val kafkaParams = Map("metadata.broker.list" -> "localhost:9092")

    val topics = args.toSet

    val stream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](
      ssc, kafkaParams, topics)

    stream.foreachRDD { rdd =>

      if (!rdd.isEmpty()) {
        try {

          rdd.foreachPartition { iter =>
            iter.foreach {
              case (key, msg) =>

                val obj = msgParseMaster(msg)

                val collection = sc.parallelize(Seq((obj.id, obj.data)))
                collection.saveToCassandra("testKS", "testTable ", SomeColumns("id", "data"))

            }
          }

        }

      }
    }

    ssc.start()
    ssc.awaitTermination()

  }

  import org.json4s._
  import org.json4s.native.JsonMethods._
  case class wordCount(id: Long, data: String) extends serializable
  implicit val formats = DefaultFormats
  def msgParseMaster(msg: String): wordCount = {
    val m = parse(msg).extract[wordCount]
    return m

  }

}

我来了

org.apache.spark.SparkException:任务不可序列化

以下是完整的日志

16/08/06 10:24:52 错误 JobScheduler:运行作业流作业时出错 1470504292000 ms.0 org.apache.spark.SparkException:任务不可序列化 在 org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:304) 在 org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:294) 在 org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:122) 在 org.apache.spark.SparkContext.clean(SparkContext.scala:2055) 在 org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1.apply(RDD.scala:919) 在 org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1.apply(RDD.scala:918) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:150) 在 org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:111) 在 org.apache.spark.rdd.RDD.withScope(RDD.scala:316) 在 org.apache.spark.rdd.RDD.foreachPartition(RDD.scala:918) 在

【问题讨论】:

    标签: scala apache-spark spark-streaming spark-cassandra-connector


    【解决方案1】:

    您不能在传递给foreachPartition 的函数中调用sc.parallelize - 该函数必须被序列化并发送到每个执行程序,并且SparkContext(故意)不可序列化(它应该只驻留在驱动程序中应用程序,而不是执行者)。

    【讨论】:

    • 我正在研究 Databricks 笔记本以获得视觉帮助。代码在我的 Scala IDE 中运行良好,但在 databricks 中抛出 NotSerializableException。 val predictionAndLabels = testData.map { case LabeledPoint(label, features) => val prediction = model.predict(features) (prediction, label) } //predictionAndLabels.saveAsTextFile("/Sparkspace/result")
    • 我正在研究 Databricks 笔记本以获得视觉帮助。代码在我的 Scala IDE 中运行良好,但在 databricks 中抛出 NotSerializableException。 val predictionAndLabels = testData.map { case LabeledPoint(label, features) => val prediction = model.predict(features) (prediction, label) } val metrics = new BinaryClassificationMetrics(predictionAndLabels) java.io.NotSerializableException:org.apache.spark .mllib.evaluation.BinaryClassificationMetrics
    【解决方案2】:

    SparkContext 不可序列化,你不能在foreachRDD 中使用它,并且从你的图表的使用中你不需要它。相反,您可以简单地映射每个 RDD,解析出相关数据并将新的 RDD 保存到 cassandra:

    stream
      .map { 
        case (_, msg) => 
          val result = msgParseMaster(msg)
          (result.id, result.data)
       }
      .foreachRDD(rdd => if (!rdd.isEmpty)
                           rdd.saveToCassandra("testKS",
                                               "testTable",
                                                SomeColumns("id", "data")))
    

    【讨论】:

    • @Suresh 以上对我来说看起来不错,但实际上您应该测试一下它是否符合您的性能目标。
    • @尤瓦尔。当然谢谢。
    • 嗨 Yuval,有没有办法过滤数据?如果 result.id 值是偶数,我想保存到 testTable_Even,如果 result.id 值是奇数,那么我想保存到 testTable_Odd
    • @Suresh 我建议打开一个新问题,因为它与这个问题不同。
    • @Yuval 发布了新问题stackoverflow.com/questions/38811434/…
    猜你喜欢
    • 1970-01-01
    • 2019-11-11
    • 1970-01-01
    • 2017-07-12
    • 1970-01-01
    • 1970-01-01
    • 2019-06-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多