【问题标题】:org.apache.spark.SparkException: Couldn't find leader offsets for Set([test-topic,0])org.apache.spark.SparkException:找不到 Set([test-topic,0])的领导者偏移量
【发布时间】:2016-09-06 10:38:23
【问题描述】:

我尝试使用 Confluent 平台,并以 this code 为例,向 REST 端点发出高级 Kafka 请求。

我使用以下 Kafka 参数:

val kafkaParams = Map(
  "bootstrap.servers" -> "localhost:9092",
  "schema.registry.url" -> "http://localhost:8081",
  "group.id" -> "EventConsumer",
  "auto.offset.reset" -> "smallest"
)

这是我尝试运行代码时遇到的错误。错误发生在以下行:

@transient val kafkaStream: DStream[(String, Object)] =
  KafkaUtils.createDirectStream[String, Object, StringDecoder, KafkaAvroDecoder](
    ssc, kafkaParams, Set(topic)
  )

线程“主”org.apache.spark.SparkException 中的异常: java.nio.channels.ClosedChannelException org.apache.spark.SparkException:找不到领导者偏移量 设置([测试主题,0])在 org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366) 在 org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366) 在 scala.util.Either.fold(Either.scala:98) 在 org.apache.spark.streaming.kafka.KafkaCluster$.checkErrors(KafkaCluster.scala:365) 在 org.apache.spark.streaming.kafka.KafkaUtils$.getFromOffsets(KafkaUtils.scala:222) 在 org.apache.spark.streaming.kafka.KafkaUtils$.createDirectStream(KafkaUtils.scala:484) 在 kafka.EventsConsumer$.delayedEndpoint$kafka$EventsConsumer$1(EventsConsumer.scala:53) 在 kafka.EventsConsumer$delayedInit$body.apply(EventsConsumer.scala:22) 在 scala.Function0$class.apply$mcV$sp(Function0.scala:34) 在 scala.runtime.AbstractFunction0.apply$mcV$sp(AbstractFunction0.scala:12) 在 scala.App$$anonfun$main$1.apply(App.scala:76) 在 scala.App$$anonfun$main$1.apply(App.scala:76) 在 scala.collection.immutable.List.foreach(List.scala:381) 在 scala.collection.generic.TraversableForwarder$class.foreach(TraversableForwarder.scala:35) 在 scala.App$class.main(App.scala:76) 在 kafka.EventsConsumer$.main(EventsConsumer.scala:22) 在 kafka.EventsConsumer.main(EventsConsumer.scala) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:498) 在 com.intellij.rt.execution.application.AppMain.main(AppMain.java:147)

更新:

我尝试将localhost更改为IP,但仍然遇到同样的问题。

【问题讨论】:

  • 有什么解决办法吗?

标签: scala apache-spark apache-kafka confluent-platform


【解决方案1】:

看起来领导者不适用于主题分区。尝试描述主题并检查是否有任何领导者可用于 test-topic 的分区 0。如果分区的所有副本都关闭,则会发生这种情况。如果您的复制因子为 1,那么这是最可能的原因。

【讨论】:

  • 我尝试使用curl -i -X GET "Accept: application/vnd.kafka.avro.v1+json" http://localhost:8082/topics/test-topic获取Kafka主题的内容,并收到了{"name":"test-topic","configs":{},"partitions":[{"partition":0,"leader":0,"replicas":[{"broker":0,"leader":true,"in_sync":true}]}]}的响应。所以,我假设分区没有问题。如何查看更多详情?
猜你喜欢
  • 2016-03-21
  • 2017-01-09
  • 1970-01-01
  • 2015-10-25
  • 2016-02-22
  • 2017-07-12
  • 1970-01-01
  • 1970-01-01
  • 2019-09-26
相关资源
最近更新 更多