【问题标题】:Java Execption while parsing JSON RDD from Kafka Stream从 Kafka Stream 解析 JSON RDD 时出现 Java 异常
【发布时间】:2016-09-27 19:18:28
【问题描述】:

我正在尝试使用 Spark 流库从 kafka 读取 json 字符串。该代码能够连接到 kafka 代理,但在解码消息时失败。代码灵感来自

https://github.com/killrweather/killrweather/blob/master/killrweather-examples/src/main/scala/com/datastax/killrweather/KafkaStreamingJson.scala

val kStream = KafkaUtils.createDirectStream[String, String, StringDecoder, 
         StringDecoder](ssc, kParams, kTopic).map(_._2)
  println("Starting to read from kafka topic:" + topicStr)
kStream.foreachRDD { rdd =>

   if (rdd.toLocalIterator.nonEmpty) {

          val sqlContext = new org.apache.spark.sql.SQLContext(sc)
            sqlContext.read.json(rdd).registerTempTable("mytable")
            if (firstTime) {
                sqlContext.sql("SELECT * FROM mytable").printSchema()
            }
            val df = sqlContext.sql(selectStr)
            df.collect.foreach(println)
            df.rdd.saveAsTextFile(fileName)
            mergeFiles(fileName, firstTime)
            firstTime = false
           println(rdd.name)
        }

java.lang.NoSuchMethodError: kafka.message.MessageAndMetadata.(Ljava/lang/String;ILkafka/message/Message;JLkafka/serializer/Decoder;Lkafka/serializer/Decoder;)V 在 org.apache.spark.streaming.kafka.KafkaRDD$KafkaRDDIterator.getNext(KafkaRDD.scala:222) 在 org.apache.spark.util.NextIterator.hasNext(NextIterator.scala:73) 在 scala.collection.Iterator$$anon$11.hasNext(Iterator.scala:327) 在 scala.collection.Iterator$class.foreach(Ite​​rator.scala:727) 在 scala.collection.AbstractIterator.foreach(Ite​​rator.scala:1157) 在 scala.collection.generic.Growable$class.$plus$plus$eq(Growable.scala:48) 在 scala.collection.mutable.ArrayBuffer.$plus$plus$eq(ArrayBuffer.scala:103) 在 scala.collection.mutable.ArrayBuffer.$plus$plus$eq(ArrayBuffer.scala:47) 在 scala.collection.TraversableOnce$class.to(TraversableOnce.scala:27​​3) 在 scala.collection.AbstractIterator.to(Iterator.scala:1157) 在 scala.collection.TraversableOnce$class.toBuffer(TraversableOnce.scala:265)

【问题讨论】:

  • 您是如何完成这项工作的?似乎 kafka 在运行时不可用
  • Kafka 可用并建立连接。我通过更改为随机卡夫卡经纪人对此进行了负面测试。异常来自 if (rdd.toLocalIterator.nonEmpty) { 行

标签: json scala apache-spark apache-kafka spark-streaming


【解决方案1】:

问题出在所使用的 Kafka jar 版本上,使用 0.9.0.0 解决了这些问题。 0.8.2.0 引入了 kafka.message.MessageAndMetadata 类。

【讨论】:

    猜你喜欢
    • 2018-08-11
    • 1970-01-01
    • 2019-06-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-05-30
    相关资源
    最近更新 更多