【问题标题】:Failed to find leader for Set([topic,0]) with Kafka-Spark Integration使用 Kafka-Spark 集成无法找到 Set([topic,0]) 的领导者
【发布时间】:2018-10-27 20:01:12
【问题描述】:

我正在尝试将 SSL 用于 Kafka-Spark 集成。我已经在启用 SSL 的情况下测试了 Kafka,它与示例消费者和生产者完全兼容。

另外,我尝试了 Spark - Kafka 的集成,当在 spark-job 中 没有 SSL 时,它也可以正常工作.

现在,当我在 spark-job 中启用 SSL 时,出现异常并且集成不起作用。

我为在 spark-job 中启用 SSL 所做的ONLY 更改是在我的工作中包含以下代码行:

    sparkConf.set("security.protocol", "SSL");
    sparkConf.set("ssl.truststore.location", "PATH/truststore.jks");
    sparkConf.set("ssl.truststore.password", "passwrd");
    sparkConf.set("ssl.keystore.location", "PATH/keystore.jks");
    sparkConf.set("ssl.keystore.password", "kstore");
    sparkConf.set("ssl.key.password", "keypass");

并且这个 sparkConf 在创建流上下文时被传递。

JavaStreamingContext jssc = new JavaStreamingContext(sparkConf, new Duration(10000));

当我运行 job 时,我得到的错误如下:

17/05/24 18:16:39 WARN ConsumerFetcherManager$LeaderFinderThread: [test-consumer-group_bmj-cluster-1495664195784-5f49cbd0-leader-finder-thread], Failed to find leader for Set([bell,0])
java.lang.NullPointerException
    at org.apache.kafka.common.utils.Utils.formatAddress(Utils.java:312)
    at kafka.cluster.Broker.connectionString(Broker.scala:62)
    at kafka.client.ClientUtils$$anonfun$fetchTopicMetadata$5.apply(ClientUtils.scala:89)
    at kafka.client.ClientUtils$$anonfun$fetchTopicMetadata$5.apply(ClientUtils.scala:89)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.TraversableLike$$anonfun$map$1.apply(TraversableLike.scala:234)
    at scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
    at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
    at scala.collection.TraversableLike$class.map(TraversableLike.scala:234)
    at scala.collection.AbstractTraversable.map(Traversable.scala:104)
    at kafka.client.ClientUtils$.fetchTopicMetadata(ClientUtils.scala:89)
    at kafka.consumer.ConsumerFetcherManager$LeaderFinderThread.doWork(ConsumerFetcherManager.scala:66)
    at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:60)

Kafka 版本 - 2.11-0.10.2.0
Spark 版本 - 2.1.0
Scala 版本 - 2.11.8

流媒体库

  <!-- https://mvnrepository.com/artifact/org.apache.spark/spark-streaming_2.10 -->
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-streaming_2.11</artifactId>
        <version>2.1.0</version>
    </dependency>

    <!-- https://mvnrepository.com/artifact/org.apache.spark/spark-streaming-kafka_2.10 -->
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-streaming-kafka-0-8_2.11</artifactId>
        <version>2.1.0</version>
    </dependency>

对克服这个问题有什么帮助吗?

【问题讨论】:

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


    【解决方案1】:

    通过深入研究,我能够找出我遇到的问题。

    首先,为了启用SSL相关的SSL,kafka-params需要传入KafkaUtils.createDirectStream() JavaStreamingContextsparkConf 方法和 NOT

    然后,给定SSL参数

    "security.protocol", "SSL"
    "ssl.truststore.location", "PATH/truststore.jks"
    "ssl.truststore.password", "passwrd"
    "ssl.keystore.location", "PATH/keystore.jks"
    "ssl.keystore.password", "kstore"
    "ssl.key.password", "keypass"
    

    我正在使用的 spark-kafka-streaming 版本“0-8_2.11”不支持,因此我必须将其更改为版本“0-10_2.11”。

    这反过来对方法进行了完整的API更改:KafkaUtils.createDirectStream(),用于连接到Kafka。

    文档中给出了关于如何使用它的说明here

    所以我连接到 Kafka 的最终代码 sn-p 如下所示:

    final JavaInputDStream<ConsumerRecord<String, String>> stream =
                KafkaUtils.createDirectStream(
                        javaStreamingContext,
                        LocationStrategies.PreferConsistent(),
                        ConsumerStrategies.<String, String>Subscribe(topicsCollection, kafkaParams)
                );
    

    kafka-params 是一个包含所有 SSL 参数的映射。

    谢谢
    沙比尔

    【讨论】:

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