【问题标题】:spark java.io.NotSerializableException: org.apache.spark.SparkContext火花java.io.NotSerializableException:org.apache.spark.SparkContext
【发布时间】:2016-06-07 17:03:28
【问题描述】:

我正在尝试通过 spark 流从 kafka 接收到的消息来检查存在记录,现在当我运行 RunReadLogByKafka 对象时,抛出了 SparkContext 的 NotSerializableException,我用谷歌搜索它,但我仍然没有知道如何解决它,有人可以建议我如何重写它吗?提前致谢。

package com.test.spark.hbase

import java.sql.{DriverManager, PreparedStatement, Connection}
import java.text.SimpleDateFormat
import com.powercn.spark.LogRow
import com.powercn.spark.SparkReadHBaseTest.{SensorStatsRow, SensorRow}
import kafka.serializer.{DefaultDecoder, StringDecoder}
import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.hbase.{TableName, HBaseConfiguration}
import org.apache.hadoop.hbase.client.{ConnectionFactory, Result, Put}
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.hadoop.hbase.mapreduce.{TableInputFormat, TableOutputFormat}
import org.apache.hadoop.hbase.util.Bytes
import org.apache.hadoop.io.Text
import org.apache.hadoop.mapred.JobConf
import org.apache.spark.sql.{SQLContext, Row}
import org.apache.spark.{SparkContext, SparkConf}
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Seconds, StreamingContext}

case class LogRow(rowkey: String, content: String)

object LogRow {
  def parseLogRow(result: Result): LogRow = {
    val rowkey = Bytes.toString(result.getRow())
    val p0 = rowkey
    val p1 = Bytes.toString(result.getValue(Bytes.toBytes("data"), Bytes.toBytes("content")))
    LogRow(p0, p1)
  }
}

 class ReadLogByKafka(sct:SparkContext) extends Serializable {

   implicit def func(records: String) {
    @transient val conf = HBaseConfiguration.create()
    conf.set(TableInputFormat.INPUT_TABLE, "log")
    @transient val sc = sct
    @transient val sqlContext = SQLContext.getOrCreate(sc)
    import sqlContext.implicits._
    try {
      //query info table
      val hBaseRDD = sc.newAPIHadoopRDD(conf, classOf[TableInputFormat],
        classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable],
        classOf[org.apache.hadoop.hbase.client.Result])
      println(hBaseRDD.count())
      // transform (ImmutableBytesWritable, Result) tuples into an RDD of Results
      val resultRDD = hBaseRDD.map(tuple => tuple._2)
      println(resultRDD.count())
      val logRDD = resultRDD.map(LogRow.parseLogRow)
      val logDF = logRDD.toDF()
      logDF.printSchema()
      logDF.show()
      // register the DataFrame as a temp table
      logDF.registerTempTable("LogRow")
      val logAdviseDF = sqlContext.sql("SELECT rowkey, content as content FROM LogRow ")
      logAdviseDF.printSchema()
      logAdviseDF.take(5).foreach(println)
    } catch {
      case e: Exception => e.printStackTrace()
    } finally {

    }
  }

}

package com.test.spark.hbase

import kafka.serializer.{DefaultDecoder, StringDecoder}
import org.apache.hadoop.conf.Configuration
import org.apache.hadoop.hbase.HBaseConfiguration
import org.apache.hadoop.hbase.mapreduce.{TableInputFormat, TableOutputFormat}
import org.apache.hadoop.hbase.util.Bytes
import org.apache.hadoop.io.Text
import org.apache.spark.SparkConf
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Seconds, StreamingContext}


object RunReadLogByKafka extends Serializable {

  def main(args: Array[String]): Unit = {
    val broker = "192.168.13.111:9092"
    val topic = "log"

    @transient val sparkConf = new SparkConf().setAppName("RunReadLogByKafka")
    @transient val streamingContext = new StreamingContext(sparkConf, Seconds(2))
    @transient val sc = streamingContext.sparkContext

    val kafkaConf = Map("metadata.broker.list" -> broker,
      "group.id" -> "group1",
      "zookeeper.connection.timeout.ms" -> "3000",
      "kafka.auto.offset.reset" -> "smallest")

    // Define which topics to read from
    val topics = Set(topic)

    val messages = KafkaUtils.createDirectStream[Array[Byte], String, DefaultDecoder, StringDecoder](
      streamingContext, kafkaConf, topics).map(_._2)

    messages.foreachRDD(rdd => {
      val readLogByKafka =new ReadLogByKafka(sc)
      //parse every message, it will throw NotSerializableException
      rdd.foreach(readLogByKafka.func)
    })

    streamingContext.start()
    streamingContext.awaitTermination()
  }


}

Exception in thread "main" org.apache.spark.SparkException: Task not serializable
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:304)
    at org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:294)
    at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:122)
    at org.apache.spark.SparkContext.clean(SparkContext.scala:2021)
    at org.apache.spark.rdd.RDD$$anonfun$foreach$1.apply(RDD.scala:889)
    at org.apache.spark.rdd.RDD$$anonfun$foreach$1.apply(RDD.scala:888)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:147)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:108)
    at org.apache.spark.rdd.RDD.withScope(RDD.scala:306)
    at org.apache.spark.rdd.RDD.foreach(RDD.scala:888)
    at com.test.spark.hbase.RunReadLogByKafka$$anonfun$main$1.apply(RunReadLogByKafka.scala:38)
    at com.test.spark.hbase.RunReadLogByKafka$$anonfun$main$1.apply(RunReadLogByKafka.scala:35)
    at org.apache.spark.streaming.dstream.DStream$$anonfun$foreachRDD$1$$anonfun$apply$mcV$sp$3.apply(DStream.scala:631)
    at org.apache.spark.streaming.dstream.DStream$$anonfun$foreachRDD$1$$anonfun$apply$mcV$sp$3.apply(DStream.scala:631)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1$$anonfun$apply$mcV$sp$1.apply$mcV$sp(ForEachDStream.scala:42)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1$$anonfun$apply$mcV$sp$1.apply(ForEachDStream.scala:40)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1$$anonfun$apply$mcV$sp$1.apply(ForEachDStream.scala:40)
    at org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:399)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply$mcV$sp(ForEachDStream.scala:40)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply(ForEachDStream.scala:40)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply(ForEachDStream.scala:40)
    at scala.util.Try$.apply(Try.scala:161)
    at org.apache.spark.streaming.scheduler.Job.run(Job.scala:34)
    at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler$$anonfun$run$1.apply$mcV$sp(JobScheduler.scala:207)
    at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler$$anonfun$run$1.apply(JobScheduler.scala:207)
    at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler$$anonfun$run$1.apply(JobScheduler.scala:207)
    at scala.util.DynamicVariable.withValue(DynamicVariable.scala:57)
    at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler.run(JobScheduler.scala:206)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
    at java.lang.Thread.run(Thread.java:745)
Caused by: java.io.NotSerializableException: org.apache.spark.SparkContext
Serialization stack:
    - object not serializable (class: org.apache.spark.SparkContext, value: org.apache.spark.SparkContext@5207878f)
    - field (class: com.test.spark.hbase.ReadLogByKafka, name: sct, type: class org.apache.spark.SparkContext)
    - object (class com.test.spark.hbase.ReadLogByKafka, com.test.spark.hbase.ReadLogByKafka@60212100)
    - field (class: com.test.spark.hbase.RunReadLogByKafka$$anonfun$main$1$$anonfun$apply$1, name: readLogByKafka$1, type: class com.test.spark.hbase.ReadLogByKafka)
    - object (class com.test.spark.hbase.RunReadLogByKafka$$anonfun$main$1$$anonfun$apply$1, <function1>)
    at org.apache.spark.serializer.SerializationDebugger$.improveException(SerializationDebugger.scala:40)
    at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:47)
    at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:84)
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:301)
    ... 30 more

【问题讨论】:

  • 我已阅读此链接,但我不知道如何重写我的代码,我尝试为 sparkContext 添加@transient,但它仍然抛出 NotSerializableException。
  • 尝试启用spark.driver.allowMultipleContexts,然后在foreachReadLogByKafka 中创建SparkContext注意 这不是推荐的方法,但在某些极端情况下您可能需要它。
  • 谢谢你的帮助,是的,如果elableallowMultipleContexts,SparkContext可以在ReadLogByKafka的foreach中调用,不知道能不能只用一个context实现。

标签: scala apache-spark spark-streaming


【解决方案1】:

您在foreachRDD 中作为参数获得的对象是标准的Spark RDD,因此您必须遵守与往常完全相同的规则,包括没有嵌套操作或转换,也无法访问SparkContext。目前尚不清楚您尝试实现的目标(看起来ReadLogByKafka.func 没有做任何有用的事情),但我猜您正在寻找某种join

【讨论】:

  • 感谢您的解释,事实上,我想在这里通过 spark hive 接收到的消息来执行查询 hbase,所以我想在 'ReadLogByKafka.func' 中的 'foreach' 中获取 'sparkContext'
  • 您可以直接尝试,但不能使用Spark工具。
猜你喜欢
  • 2011-01-22
  • 1970-01-01
  • 2017-04-27
  • 1970-01-01
  • 2016-10-20
  • 2023-03-26
  • 2016-12-18
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多