【问题标题】:Task not serializable: Json strings processing using Spark Streaming任务不可序列化:使用 Spark Streaming 处理 Json 字符串
【发布时间】:2016-07-20 13:47:43
【问题描述】:

我在 Spark Streaming (Scala) 中接收来自 Kafka 的 Json 字符串。每个字符串的处理都需要一些时间,所以我想将处理分布在 X 个集群上。

目前我只是在我的笔记本电脑上进行测试。因此,为简单起见,我们假设我应该对每个 Json 字符串应用的处理只是字段的一些规范化:

  def normalize(json: String): String = {
    val parsedJson = Json.parse(json)
    val parsedRecord = (parsedJson \ "records")(0)
    val idField = parsedRecord \ "identifier"
    val titleField = parsedRecord \ "title"

    val output = Json.obj(
      "id" -> Json.parse(idField.get.toString().replace('/', '-')),
      "publicationTitle" -> titleField.get
    )
    output.toString()
  }

这是我尝试将操作normalize 分配到“集群”上(每个 Json 字符串都应该被完全处理;Json 字符串不能被拆分)。如何处理val callRDD = JSONstrings.map(normalize(_))Task not serializable的问题?

val conf = new SparkConf().setAppName("My Spark Job").setMaster("local[*]")
val ssc = new StreamingContext(conf, Seconds(5))

val topicMap = topic.split(",").map((_, numThreads)).toMap

val JSONstrings = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)

val callRDD = JSONstrings.map(normalize(_))

ssc.start()
ssc.awaitTermination()

更新

这是完整的代码:

package org.consumer.kafka

import java.util.Properties
import java.util.concurrent._
import com.typesafe.config.ConfigFactory
import kafka.consumer.{Consumer, ConsumerConfig}
import kafka.utils.Logging
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerConfig, ProducerRecord}
import org.apache.log4j.{Level, Logger}
import org.apache.spark.SparkConf
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Seconds, StreamingContext}
import play.api.libs.json.{JsObject, JsString, JsValue, Json}
import scalaj.http.{Http, HttpResponse}

class KafkaJsonConsumer(val datasource: String,
                        val apiURL: String,
                        val zkQuorum: String,
                        val group: String,
                        val topic: String) extends Logging
{
  val delay = 1000
  val config = createConsumerConfig(zkQuorum, group)
  val consumer = Consumer.create(config)
  var executor: ExecutorService = null

  def shutdown() = {
    if (consumer != null)
      consumer.shutdown();
    if (executor != null)
      executor.shutdown();
  }

  def createConsumerConfig(zkQuorum: String, group: String): ConsumerConfig = {
    val props = new Properties()
    props.put("zookeeper.connect", zkQuorum);
    props.put("group.id", group);
    props.put("auto.offset.reset", "largest");
    props.put("zookeeper.session.timeout.ms", "2000");
    props.put("zookeeper.sync.time.ms", "200");
    props.put("auto.commit.interval.ms", "1000");
    val config = new ConsumerConfig(props)
    config
  }

  def run(numThreads: Int) = {
    val conf = new SparkConf()
                              .setAppName("TEST")
                              .setMaster("local[*]")
                              //.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
    val ssc = new StreamingContext(conf, Seconds(5))
    ssc.checkpoint("checkpoint")

    val topicMap = topic.split(",").map((_, numThreads)).toMap

    val rawdata = KafkaUtils.createStream(ssc, zkQuorum, group, topicMap).map(_._2)

    val parsed = rawdata.map(Json.parse(_))

    val result = parsed.map(record => {
      val parsedRecord = (record \ "records")(0)
      val idField = parsedRecord \ "identifier"
      val titleField = parsedRecord \ "title"
      val journalTitleField = parsedRecord \ "publicationName"
      Json.obj(
        "id" -> Json.parse(idField.get.toString().replace('/', '-')),
        "publicationTitle" -> titleField.get,
        "journalTitle" -> journalTitleField.get)
    })

    result.print

    val callRDD = result.map(JsonUtils.normalize(_))

    callRDD.print()

    ssc.start()
    ssc.awaitTermination()
  }

  object JsonUtils {
    def normalize(json: JsValue): String = {
      (json \ "id").as[JsString].value
    }
  }

}

我启动这个 calss KafkaJsonConsumer 的执行如下:

package org.consumer

import org.consumer.kafka.KafkaJsonConsumer

object TestConsumer {

  def main(args: Array[String]) {

    if (args.length < 6) {
      System.exit(1)
    }

    val Array(datasource, apiURL, zkQuorum, group, topic, numThreads) = args

    val processor = new KafkaJsonConsumer(datasource, apiURL, zkQuorum, group, topic)
    processor.run(numThreads.toInt)

    //processor.shutdown()

  }

}

【问题讨论】:

    标签: scala apache-spark spark-streaming


    【解决方案1】:

    看起来normalize 方法是某个类的一部分。在map 操作中使用它的行中,Spark 不仅需要序列化方法本身,还需要序列化它所属的整个实例。最简单的解决方案是将normalize 移动到某个单例对象:

    object JsonUtils {
      def normalize(json: String): String = ???
    }
    

    然后像这样调用:

    val callRDD = JSONstrings.map(JsonUtils.normalize(_))
    

    【讨论】:

    • 它仍然说任务不可序列化。尝试 Kryo 有意义吗?
    • 不,这个问题与使用或不使用Kryo无关。您是在 spark-shell 中还是在应用程序中运行代码?
    • 我从 Intellij 运行我的应用程序,而不是 spark-shell。我使用 local[*] 只是为了测试。一旦它在本地工作,我的想法是使用纱线并将计算分布在不同的机器上。我已经在这个问题上花了很多时间,最后我不知道如何处理它。
    • 你不认为我的问题与这篇文章有关吗?:http://stackoverflow.com/questions/28554141/how-to-let-spark-serialize-an-object-using-kryo?rq=1。如果它可能有助于检测问题,我现在收到此错误:Caused by: java.io.NotSerializableException: org.consumer.kafka.KafkaJsonConsumer,其中KafkaJsonConsumer 是包含所有提到的代码的类
    • 你应该提供完整的代码,以及类。
    猜你喜欢
    • 1970-01-01
    • 2018-04-06
    • 2021-03-16
    • 1970-01-01
    • 2015-12-16
    • 2017-03-21
    • 2016-12-30
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多