【问题标题】:Spark Streaming Kafka streamSpark Streaming Kafka 流
【发布时间】:2016-03-12 18:12:51
【问题描述】:

我在尝试使用 spark 流从 kafka 中读取数据时遇到了一些问题。

我的代码是:

val sparkConf = new SparkConf().setMaster("local[2]").setAppName("KafkaIngestor")
val ssc = new StreamingContext(sparkConf, Seconds(2))

val kafkaParams = Map[String, String](
  "zookeeper.connect" -> "localhost:2181",
  "group.id" -> "consumergroup",
  "metadata.broker.list" -> "localhost:9092",
  "zookeeper.connection.timeout.ms" -> "10000"
  //"kafka.auto.offset.reset" -> "smallest"
)

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

我之前在 2181 端口启动了 zookeeper,在 9092 端口启动了 Kafka 服务器 0.9.0.0。 但我在 Spark 驱动程序中收到以下错误:

Exception in thread "main" java.lang.ClassCastException: kafka.cluster.BrokerEndPoint cannot be cast to kafka.cluster.Broker
at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2$$anonfun$3$$anonfun$apply$6$$anonfun$apply$7.apply(KafkaCluster.scala:90)
at scala.Option.map(Option.scala:145)
at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2$$anonfun$3$$anonfun$apply$6.apply(KafkaCluster.scala:90)
at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2$$anonfun$3$$anonfun$apply$6.apply(KafkaCluster.scala:87)

动物园管理员日志:

[2015-12-08 00:32:08,226] INFO Got user-level KeeperException when processing sessionid:0x1517ec89dfd0000 type:create cxid:0x34 zxid:0x1d3 txntype:-1 reqpath:n/a Error Path:/brokers/ids Error:KeeperErrorCode = NodeExists for /brokers/ids (org.apache.zookeeper.server.PrepRequestProcessor)

有什么提示吗?

非常感谢

【问题讨论】:

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


    【解决方案1】:

    Kafka10 流式传输 / Spark 2.1.0 / DCOS / Mesosphere

    Ugg 我花了一整天的时间在这上面,并且一定已经阅读了十几遍这篇文章。我试过 spark 2.0.0、2.0.1、Kafka 8、Kafka 10。远离 Kafka 8 和 spark 2.0.x,依赖就是一切。从下面开始。它有效。

    SBT:

    "org.apache.hadoop" % "hadoop-aws" % "2.7.3" excludeAll ExclusionRule(organization = "org.apache.hadoop", name = "hadoop-common"),
    "org.apache.spark" %% "spark-core" % "2.1.0",
    "org.apache.spark" %% "spark-sql" % "2.1.0" ,
    "org.apache.spark" % "spark-streaming-kafka-0-10_2.11" % "2.1.0",
    "org.apache.spark" % "spark-streaming_2.11" % "2.1.0"
    

    工作 Kafka/Spark 流代码:

    val spark = SparkSession
      .builder()
      .appName("ingest")
      .master("local[4]")
      .getOrCreate()
    
    import spark.implicits._
    val ssc = new StreamingContext(spark.sparkContext, Seconds(2))
    
    val topics = Set("water2").toSet
    
    val kafkaParams = Map[String, String](
      "metadata.broker.list"        -> "broker:port,broker:port",
      "bootstrap.servers"           -> "broker:port,broker:port",
      "group.id"                    -> "somegroup",
      "auto.commit.interval.ms"     -> "1000",
      "key.deserializer"            -> "org.apache.kafka.common.serialization.StringDeserializer",
      "value.deserializer"          -> "org.apache.kafka.common.serialization.StringDeserializer",
      "auto.offset.reset"           -> "earliest",
      "enable.auto.commit"          -> "true"
    )
    
    val messages = KafkaUtils.createDirectStream[String, String](ssc, PreferConsistent, Subscribe[String, String](topics, kafkaParams))
    
    messages.foreachRDD(rdd => {
      if (rdd.count() >= 1) {
        rdd.map(record => (record.key, record.value))
          .toDS()
          .withColumnRenamed("_2", "value")
          .drop("_1")
          .show(5, false)
        println(rdd.getClass)
      }
    })
    ssc.start()
    ssc.awaitTermination()
    

    如果你看到这个请点赞,这样我可以获得一些声望点。 :)

    【讨论】:

      【解决方案2】:

      问题与错误的 spark-streaming-kafka 版本有关。

      documentation中所述

      Kafka:Spark Streaming 1.5.2 与 Kafka 0.8.2.1 兼容

      所以,包括

      <dependency>
          <groupId>org.apache.kafka</groupId>
          <artifactId>kafka_2.10</artifactId>
          <version>0.8.2.2</version>
      </dependency>
      

      在我的 pom.xml(而不是版本 0.9.0.0)中解决了这个问题。

      希望对你有帮助

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-08-08
        • 2018-05-17
        • 1970-01-01
        • 1970-01-01
        • 2017-12-28
        • 2020-10-29
        • 2019-01-15
        • 2015-10-15
        相关资源
        最近更新 更多