【问题标题】:Spark Streaming Kafka consumer doesn't like DStreamSpark Streaming Kafka 消费者不喜欢 DStream
【发布时间】:2018-12-01 16:56:39
【问题描述】:

我正在使用 Spark Shell(Scala 2.10 和 Spark Streaming org.apache.spark:spark-streaming-kafka-0-10_2.10:2.0.1)来测试 Spark/Kafka 消费者:

import org.apache.kafka.clients.consumer.ConsumerRecord
import org.apache.kafka.common.serialization.StringDeserializer
import org.apache.spark.streaming.kafka010._
import org.apache.spark.streaming.kafka010.LocationStrategies.PreferConsistent
import org.apache.spark.streaming.kafka010.ConsumerStrategies.Subscribe
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.dstream.DStream

val kafkaParams = Map[String, Object](
  "bootstrap.servers" -> "mykafka01.example.com:9092",
  "key.deserializer" -> classOf[StringDeserializer],
  "value.deserializer" -> classOf[StringDeserializer],
  "group.id" -> "mykafka",
  "auto.offset.reset" -> "latest",
  "enable.auto.commit" -> (false: java.lang.Boolean)
)

val topics = Array("mytopic")

def createKafkaStream(ssc: StreamingContext, topics: Array[String], kafkaParams: Map[String,Object]) : DStream[(String, String)] = {
    KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics, kafkaParams))
}

def messageConsumer(): StreamingContext = {
    val ssc = new StreamingContext(SparkContext.getOrCreate(), Seconds(10))

    createKafkaStream(ssc, topics, kafkaParams).foreachRDD(rdd => {
        rdd.collect().foreach { msg =>
            try {
                println("Received message: " + msg._2)
            } catch {
                case e @ (_: Exception | _: Error | _: Throwable) => {
                println("Exception: " + e.getMessage)
                e.printStackTrace()
            }
          }
        }
    })

    ssc
}

val ssc = StreamingContext.getActiveOrCreate(messageConsumer)
ssc.start()
ssc.awaitTermination()

当我运行它时,我得到以下异常:

<console>:60: error: type mismatch;
 found   : org.apache.spark.streaming.dstream.InputDStream[org.apache.kafka.clients.consumer.ConsumerRecord[String,String]]
 required: org.apache.spark.streaming.dstream.DStream[(String, String)]
                  KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics, kafkaParams))
                                                               ^

我一遍又一遍地检查了 Scala/API 文档,这段代码看起来应该可以正确执行。知道我哪里出错了吗?

【问题讨论】:

    标签: apache-spark apache-kafka spark-streaming kafka-consumer-api spark-streaming-kafka


    【解决方案1】:

    Subscribetopics 参数作为Array[String],您将按照def createKafkaStream(ssc: StreamingContext, topics: String, 传递单个字符串。将参数类型更改为 Array[String](并适当地调用它)将解决问题。

    【讨论】:

    • @smeeb 我没有那个版本的 kafka lib 可以尝试,但是你可以查看重载的 createDirectStream 方法,看看它们的返回类型是什么以及它们采用什么参数。您调用的方法似乎返回 DStream[ConsumerRecord[K, V]] 而不是您期望的 DStream[K, V]。或者,如果这是唯一的选择,请将您的代码更改为接受DStream[ConsumerRecord[K, V]],然后映射到(K, V
    • 再次感谢 @khachik (+1) - 当您说“您调用的方法似乎返回 DStream[ConsumerRecord[K,V]]”...您在哪里看到那。我正在查看我 认为correct javadocs 并且我没有看到任何返回 DStream[ConsumerRecord[K,V]]s 的重载 createDirectStream 方法,想法?再次感谢!!!
    • @smeeb 看着the integration guide,你调用的方法返回ConsumerRecords 的流,你应该映射它来获取键/值对:stream.map(record =&gt; (record.key, record.value))。您发布的 javadocs 似乎是针对不同版本的。
    猜你喜欢
    • 2017-02-23
    • 2018-12-18
    • 2014-12-30
    • 2017-05-04
    • 1970-01-01
    • 2020-06-13
    • 2018-07-02
    • 2017-08-08
    • 1970-01-01
    相关资源
    最近更新 更多