【问题标题】:Spark streaming kafka Couldn't find leader offsets for SetSpark Streaming kafka找不到Set的领导者偏移量
【发布时间】:2017-01-09 11:26:26
【问题描述】:

我使用 spark streaming 'org.apache.spark:spark-streaming_2.10:1.6.1' 和 'org.apache.spark:spark-streaming-kafka_2.10:1.6.1' 连接到 kafka 代理版本 0.10.0.1。当我尝试这段代码时:

def messages = KafkaUtils.createDirectStream(jssc,
            String.class,
            String.class,
            StringDecoder.class,
            StringDecoder.class,
            kafkaParams,
            topicsSet)

我收到了这个异常:

    INFO consumer.SimpleConsumer: Reconnect due to socket error: java.nio.channels.ClosedChannelException
Exception in thread "main" org.apache.spark.SparkException: java.nio.channels.ClosedChannelException
org.apache.spark.SparkException: Couldn't find leader offsets for Set([stream,0])
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366)
    at org.apache.spark.streaming.kafka.KafkaCluster$$anonfun$checkErrors$1.apply(KafkaCluster.scala:366)
    at scala.util.Either.fold(Either.scala:97)
    at org.apache.spark.streaming.kafka.KafkaCluster$.checkErrors(KafkaCluster.scala:365)
    at org.apache.spark.streaming.kafka.KafkaUtils$.getFromOffsets(KafkaUtils.scala:222)
    at org.apache.spark.streaming.kafka.KafkaUtils$.createDirectStream(KafkaUtils.scala:484)
    at org.apache.spark.streaming.kafka.KafkaUtils$.createDirectStream(KafkaUtils.scala:607)
    at org.apache.spark.streaming.kafka.KafkaUtils.createDirectStream(KafkaUtils.scala)
    at org.apache.spark.streaming.kafka.KafkaUtils$createDirectStream.call(Unknown Source)
    at org.codehaus.groovy.runtime.callsite.CallSiteArray.defaultCall(CallSiteArray.java:45)
    at org.codehaus.groovy.runtime.callsite.AbstractCallSite.call(AbstractCallSite.java:108)
    at com.privowny.classification.jobs.StreamingClassification.main(StreamingClassification.groovy:48)
    at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.lang.reflect.Method.invoke(Method.java:483)
    at org.apache.spark.deploy.SparkSubmit$.org$apache$spark$deploy$SparkSubmit$$runMain(SparkSubmit.scala:731)
    at org.apache.spark.deploy.SparkSubmit$.doRunMain$1(SparkSubmit.scala:181)
    at org.apache.spark.deploy.SparkSubmit$.submit(SparkSubmit.scala:206)
    at org.apache.spark.deploy.SparkSubmit$.main(SparkSubmit.scala:121)
    at org.apache.spark.deploy.SparkSubmit.main(SparkSubmit.scala)

我试图在这个网站上搜索一些答案,但似乎没有答案,你能给我一些建议吗?话题stream不为空。

【问题讨论】:

  • 这通常是 ZooKeeper 问题的信号。重置 ZooKeeper 并重试。
  • 可能是什么问题?我刚刚在快速入门文档中启动了服务器!
  • 我遇到了 Kafka 和 ZooKeeper 之间存在同步问题的问题。重置它们都解决了。

标签: apache-spark spark-streaming


【解决方案1】:

我也遇到过这个问题。因此,您必须更改 Kafka 上的一些配置。

转到你的Kafka配置并配置listeners;

在套接字服务器设置部分的格式:

listeners=PLAINTEXT://[hostname or IP]:[port]

例如:

listeners=PLAINTEXT://192.168.1.24:9092

【讨论】:

    【解决方案2】:

    我是从 HDP 运行 kafka,所以默认端口是 6667 而不是 9092,当我将 bootstrap.servers 的端口切换到 <hostname>:6667 时,问题得到了解决。

    【讨论】:

      【解决方案3】:

      从经验中我知道,可能导致此错误消息的一件事是,如果 Spark 驱动程序无法使用代理的广告主机名(server.properties 中的advertised.host.name)访问 kafka 代理。即使 spark 配置使用有效的不同地址识别 kafka 代理,情况也是如此。所有代理的广告主机名都必须可以从 Spark 驱动程序访问。

      这发生在我身上,因为集群在一个单独的 AWS 账户中运行,代理使用内部 DNS 记录来标识自己,这些必须复制到另一个 AWS 账户。在此之前,我收到此错误消息,因为 Spark 驱动程序无法联系到代理以询问他们的最新偏移量,即使我们在 spark 配置中使用代理的私有 IP 地址。

      希望对某人有所帮助。

      【讨论】:

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