【问题标题】:Couldn't find leaders for Set([TOPICNNAME,0])) When we are uisng in Apache Spark当我们在 Apache Spark 中使用时,找不到 Set([TOPIC NAME,0])) 的领导者
【发布时间】:2016-02-22 12:48:01
【问题描述】:

我们正在使用 Apache Spark 1.5.1 和 kafka_2.10-0.8.2.1 以及 Kafka DirectStream API 来使用 Spark 从 Kafka 获取数据。

我们使用以下设置在 Kafka 中创建主题

ReplicationFactor :1 和 Replica :1

当所有 Kafka 实例都在运行时,Spark 作业工作正常。但是,当集群中的一个 Kafka 实例关闭时,我们会得到下面重现的异常。一段时间后,我们重新启动了禁用的 Kafka 实例并尝试完成 Spark 作业,但 Spark 已经因为异常而终止。因此,我们无法读取 Kafka 主题中的剩余消息。

ERROR DirectKafkaInputDStream:125 - ArrayBuffer(org.apache.spark.SparkException: Couldn't find leaders for Set([normalized-tenant4,0]))
ERROR JobScheduler:96 - Error generating jobs for time 1447929990000 ms
org.apache.spark.SparkException: ArrayBuffer(org.apache.spark.SparkException: Couldn't find leaders for Set([normalized-tenant4,0]))
        at org.apache.spark.streaming.kafka.DirectKafkaInputDStream.latestLeaderOffsets(DirectKafkaInputDStream.scala:123)
        at org.apache.spark.streaming.kafka.DirectKafkaInputDStream.compute(DirectKafkaInputDStream.scala:145)
        at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:350)
        at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:350)
        at scala.util.DynamicVariable.withValue(DynamicVariable.scala:57)
        at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:349)
        at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:349)
        at org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:399)
        at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:344)
        at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:342)
        at scala.Option.orElse(Option.scala:257)
        at org.apache.spark.streaming.dstream.DStream.getOrCompute(DStream.scala:339)
        at org.apache.spark.streaming.dstream.ForEachDStream.generateJob(ForEachDStream.scala:38)
        at org.apache.spark.streaming.DStreamGraph$$anonfun$1.apply(DStreamGraph.scala:120)
        at org.apache.spark.streaming.DStreamGraph$$anonfun$1.apply(DStreamGraph.scala:120)
        at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:251)
        at scala.collection.TraversableLike$$anonfun$flatMap$1.apply(TraversableLike.scala:251)
        at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
        at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:47)
        at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:251)
        at scala.collection.AbstractTraversable.flatMap(Traversable.scala:105)
        at org.apache.spark.streaming.DStreamGraph.generateJobs(DStreamGraph.scala:120)
        at org.apache.spark.streaming.scheduler.JobGenerator$$anonfun$2.apply(JobGenerator.scala:247)
        at org.apache.spark.streaming.scheduler.JobGenerator$$anonfun$2.apply(JobGenerator.scala:245)
        at scala.util.Try$.apply(Try.scala:161)
        at org.apache.spark.streaming.scheduler.JobGenerator.generateJobs(JobGenerator.scala:245)
        at org.apache.spark.streaming.scheduler.JobGenerator.org$apache$spark$streaming$scheduler$JobGenerator$$processEvent(JobGenerator.scala:181)
        at org.apache.spark.streaming.scheduler.JobGenerator$$anon$1.onReceive(JobGenerator.scala:87)
        at org.apache.spark.streaming.scheduler.JobGenerator$$anon$1.onReceive(JobGenerator.scala:86)
        at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48)

提前致谢。请帮助解决此问题。

【问题讨论】:

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


    【解决方案1】:

    这是预期的行为。您已通过将 ReplicationFactor 设置为 1 来请求将每个主题存储在一台机器上。当恰好存储topic normalized-tenant4的那台机器宕机时,consumer找不到topic的leader。

    http://kafka.apache.org/documentation.html#intro_guarantees

    【讨论】:

      【解决方案2】:

      无法为指定主题找到领导者的此类错误的原因之一是一个人的 Kafka 服务器配置有问题。

      打开您的 Kafka 服务器配置:

      vim ./kafka/kafka-<your-version>/config/server.properties
      

      在“套接字服务器设置”部分中,如果您的主机缺少 IP,请提供它:

      listeners=PLAINTEXT://{host-ip}:{host-port}
      

      我正在使用 MapR 沙箱提供的 Kafka 设置,并试图通过 spark 代码访问 kafka。我在访问我的 kafka 时遇到了同样的错误,因为我的配置缺少 IP。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2016-03-21
        • 2017-01-09
        • 2015-10-25
        • 2017-07-12
        • 1970-01-01
        • 2021-09-15
        • 1970-01-01
        相关资源
        最近更新 更多