【问题标题】:kafka.cluster.BrokerEndPoint cannot be cast to kafka.cluster.Broker issuekafka.cluster.BrokerEndPoint 无法转换为 kafka.cluster.Broker 问题
【发布时间】:2017-11-07 13:35:12
【问题描述】:

我正在使用 kafka2.11-0.11.0.1、scala 2.11 和 spark 2.2.0。我在eclipse的java构建路径中添加了以下jar:

kafka-streams-0.11.0.1,
kafka-tools-0.11.0.1,
spark-streaming_2.11-2.2.0,
spark-streaming-kafka_2.11-1.6.3,
spark-streaming-kafka-0-10_2.11-2.2.0,
kafka_2.11-0.11.0.1.

我的代码如下:

import kafka.serializer.StringDecoder
import kafka.api._
import kafka.api.ApiUtils._
import org.apache.spark.SparkConf
import org.apache.spark._
import org.apache.spark.streaming._
import org.apache.spark.streaming.dstream._
import org.apache.spark.streaming.kafka
import org.apache.spark.streaming.kafka._
import org.apache.spark.streaming.kafka.KafkaUtils
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.storage.StorageLevel
import org.apache.spark.SparkContext._


object KafkaExample {

  def main(args: Array[String]) {

    val ssc = new StreamingContext("local[*]", "KafkaExample", Seconds(1))

    val kafkaParams = Map("bootstrap.servers" -> "kafkaIP:9092")

    val topics = List("logstash_log").toSet

    val stream = KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder](ssc,kafkaParams,topics).map(_._2)

    stream.print()

    ssc.checkpoint("C:/checkpoint/")
    ssc.start()
    ssc.awaitTermination()
  }
}

这是连接 spark 和 kafka 的非常简单的代码。但是,我收到此错误:

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:146)
    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)
    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
    at scala.collection.IndexedSeqOptimized$class.foreach(IndexedSeqOptimized.scala:33)
    at scala.collection.mutable.WrappedArray.foreach(WrappedArray.scala:35)
    at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
    at scala.collection.AbstractTraversable.flatMap(Traversable.scala:104)
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2$$anonfun$3.apply(KafkaCluster.scala:87)
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2$$anonfun$3.apply(KafkaCluster.scala:86)
    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
    at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:241)
    at scala.collection.immutable.Set$Set1.foreach(Set.scala:94)
    at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
    at scala.collection.AbstractTraversable.flatMap(Traversable.scala:104)
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2.apply(KafkaCluster.scala:86)
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$2.apply(KafkaCluster.scala:85)
    at scala.util.Either$RightProjection.flatMap(Either.scala:522)
    at org.apache.spark.streaming.kafka.KafkaCluster.findLeaders(KafkaCluster.scala:85)
    at org.apache.spark.streaming.kafka.KafkaCluster.getLeaderOffsets(KafkaCluster.scala:179)
    at org.apache.spark.streaming.kafka.KafkaCluster.getLeaderOffsets(KafkaCluster.scala:161)
    at org.apache.spark.streaming.kafka.KafkaCluster.getLatestLeaderOffsets(KafkaCluster.scala:150)
    at org.apache.spark.streaming.kafka.KafkaUtils$$anonfun$5.apply(KafkaUtils.scala:215)
    at org.apache.spark.streaming.kafka.KafkaUtils$$anonfun$5.apply(KafkaUtils.scala:211)
    at scala.util.Either$RightProjection.flatMap(Either.scala:522)
    at org.apache.spark.streaming.kafka.KafkaUtils$.getFromOffsets(KafkaUtils.scala:211)
    at org.apache.spark.streaming.kafka.KafkaUtils$.createDirectStream(KafkaUtils.scala:484)
    at com.defne.KafkaExample$.main(KafkaExample.scala:28)
    at com.defne.KafkaExample.main(KafkaExample.scala)

我哪里做错了?

注意:我尝试了“metadata.broker.list”而不是“bootstrap.server”,但没有任何变化。

【问题讨论】:

    标签: scala apache-spark apache-kafka


    【解决方案1】:

    您的问题是加载了太多 Kafka 依赖项,而在运行时获取的依赖项与 Spark 期望的版本不兼容。

    您的实际问题是PartitionMetadata 类。在 0.8.2 中它看起来像这样(这是你从 spark-streaming-kafka_2.11-1.6.3 得到的):

    case class PartitionMetadata(partitionId: Int, 
                                 val leader: Option[Broker], 
                                 replicas: Seq[Broker], 
                                 isr: Seq[Broker] = Seq.empty,
                                 errorCode: Short = ErrorMapping.NoError) extends Logging
    

    在 > 0.10.0.0 中是这样的:

    case class PartitionMetadata(partitionId: Int,
                                 leader: Option[BrokerEndPoint],
                                 replicas: Seq[BrokerEndPoint],
                                 isr: Seq[BrokerEndPoint] = Seq.empty,
                                 errorCode: Short = Errors.NONE.code) extends Logging
    

    看看leader 如何从Option[Broker] 变为Option[BrokerEndPoint]?这就是 Spark 大喊大叫的原因。

    你必须清理你的依赖,你需要的是(如果你使用的是 Spark 2.2)是:

    spark-streaming_2.11-2.2.0,
    spark-streaming-kafka-0-10_2.11-2.2.0
    

    【讨论】:

    • 感谢 Yuval 的快速响应。当我删除 spark-streaming-kafka_2.11-1.6.3 (import .....KafkaUtils) 时出错。当我删除 kafka_2.11-0.11.0.1 前 3 个导入时会出错。所以,我不能在 main.js 中使用方法。我删除了前 2 个罐子,没关系。但是有了这 5 个罐子,我得到了这个:线程“主”java.lang.NoClassDefFoundError 中的异常:org/apache/kafka/common/network/Send。我完全糊涂了。我不能删除它们,如果我包含,仍然不起作用。 :(
    • 如果只使用spark-streaming-kafka-0-10_2.11-2.2.0spark-streaming_2.11-2.2.0, 会出错?那不太可能。
    • 我的代码错了吗?我想我需要kafkautils。它在 kafka_2.11-0.11.0.1 中。所以我需要stringdecoder,它在kafka流中。如果我不需要这 2 个代码部分,我可以删除它们。在上述情况下,我不能。实际上,正如您所猜测的,我是初学者,需要一个快速指南来连接 kafka 和 spark。可以举个例子吗?
    • @ÖmerFarukAktaş KafkaUtils 位于 spark-streaming-kafka-0-10_2.11-2.2.0 内。同样,正如我所说,删除所有其他 jar 依赖项。
    • @YuvalItzchakov。谢谢。我有同样的错误。你的解决方案让我纠正它。你是对的,我们只需要指出这些依赖 jar(spark-streaming-kafka-0-8_2.11-2.3.0,spark-streaming-kafka-0-8-assembly_2.11-2.3.0,spark-streaming -kafka-0-10_2.11-2.3.0)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-02-19
    • 1970-01-01
    • 2014-05-12
    • 2011-10-22
    • 2020-03-02
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多