【问题标题】:Spark 2 application failed with Couldn't find leader offsets for ErrorSpark 2 应用程序因找不到错误的领导者偏移而失败
【发布时间】:2018-02-13 19:08:12
【问题描述】:

我的 Spark 应用程序从 Kafka 读取数据并摄取到 Kudu。它已成功运行近 25 小时,并将数据摄取到 Kudu。之后,我看到从 kafka 日志中为 kafka 分区选举了新的领导者。我的应用程序进入 FINISHED 状态并出现以下错误,

org.apache.spark.SparkException: ArrayBuffer(kafka.common.NotLeaderForPartitionException, org.apache.spark.SparkException: Couldn't find leader offsets for Set([test,0]))
at org.apache.spark.streaming.kafka.DirectKafkaInputDStream.latestLeaderOffsets(DirectKafkaInputDStream.scala:133)
at org.apache.spark.streaming.kafka.DirectKafkaInputDStream.compute(DirectKafkaInputDStream.scala:158)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:342)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1$$anonfun$apply$7.apply(DStream.scala:342)
at scala.util.DynamicVariable.withValue(DynamicVariable.scala:58)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:341)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1$$anonfun$1.apply(DStream.scala:341)
at org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:416)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:336)
at org.apache.spark.streaming.dstream.DStream$$anonfun$getOrCompute$1.apply(DStream.scala:334)
at scala.Option.orElse(Option.scala:289)
at org.apache.spark.streaming.dstream.DStream.getOrCompute(DStream.scala:331)
at org.apache.spark.streaming.dstream.ForEachDStream.generateJob(ForEachDStream.scala:48)
at org.apache.spark.streaming.DStreamGraph$$anonfun$1.apply(DStreamGraph.scala:122)
at org.apache.spark.streaming.DStreamGraph$$anonfun$1.apply(DStreamGraph.scala:121)
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.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
at scala.collection.TraversableLike$class.flatMap(TraversableLike.scala:241)
at scala.collection.AbstractTraversable.flatMap(Traversable.scala:104)
at org.apache.spark.streaming.DStreamGraph.generateJobs(DStreamGraph.scala:121)
at org.apache.spark.streaming.scheduler.JobGenerator$$anonfun$3.apply(JobGenerator.scala:249)
at org.apache.spark.streaming.scheduler.JobGenerator$$anonfun$3.apply(JobGenerator.scala:247)
at scala.util.Try$.apply(Try.scala:192)
at org.apache.spark.streaming.scheduler.JobGenerator.generateJobs(JobGenerator.scala:247)
at org.apache.spark.streaming.scheduler.JobGenerator.org$apache$spark$streaming$scheduler$JobGenerator$$processEvent(JobGenerator.scala:183)
at org.apache.spark.streaming.scheduler.JobGenerator$$anon$1.onReceive(JobGenerator.scala:89)
at org.apache.spark.streaming.scheduler.JobGenerator$$anon$1.onReceive(JobGenerator.scala:88)
at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:48)

Does it mean that, whenever a new leader is elected Spark application will fail?

我在 Stackoverflow 上看到很多帖子,每个人都说他们无法启动应用程序并出现此错误。但是,就我而言,它运行了 25 小时,然后就完成了。

对可能出了什么问题有任何想法吗?我搜索了 Kafka 问题,但没有与此相关的运气。

【问题讨论】:

  • 你的 Kafka 主题中分区的复制因子是多少?
  • 复制因子为 3

标签: apache-spark apache-kafka


【解决方案1】:

当 spark 流期望的 topicPartition 的计数与简单的 Kafka 客户端提供的主题分区不匹配时引发异常 [spark 用于获取 topicPartition-Leader 偏移量]。

因此,当 Spark Streaming 请求偏移量时,topicPartition [test, 0] 的领导者不可用。因此 Spark 抛出异常消息。您使用的是什么版本的 Spark 和 Kafka?

【讨论】:

  • 我使用的是 spark 2.2 和 kafka 0.9。另外,正如我在工作完成时从 kafka 日志中看到的那样,选举了一个新的领导者。 Spark Streaming 找不到新的领导者并失败了。
  • 我能否将此 spark.streaming.kafka.maxRetries 属性设置为某个值,以使驱动程序重试连接到 Kafka,这样应用程序就不会失败。这是生产中的问题。
  • 是的,增加 spark.streaming.kafka.maxRetries 肯定会有所帮助。同时增加“refreshLeaderBackoffMs”,这将增加每次尝试之间的时间。默认情况下,“refreshLeaderBackoffMs”为 200 毫秒,重试次数为 1。
猜你喜欢
  • 2016-03-21
  • 2017-01-09
  • 2017-03-16
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多